代表的な面接トピック

ストリーム処理における遅延イベントおよび順序不同イベントの処理

データ難しい
Offer.cc 編集チーム公開日 更新日

質問

モバイルの購入イベントが順序不同、遅延、および重複して到着します。最大で毎秒20,000イベントのピーク負荷において時間別売上を計算し、1分以内に結果を表示し、24時間補正を受け付けるイベント時間ストリームを設計してください。ウォーターマーク、重複排除状態、ウィンドウの更新、カットオフを超えたデータ、復旧、および検証について説明してください。

問題とスコープ

モバイルクライアントは、少なくとも event_idevent_timeuser_id、および単一通貨に変換済みの amount を含む購入イベントをレポートします。ネットワークのリトライによって重複が発生する可能性があり、オフラインのデバイスは数時間後にアップロードすることがあり、異なるパーティション間ではグローバルな順序が保持されません。毎秒最大20,000イベントの入力で、event_time ごとにUTC時計の各時間の売上を計算します。

ビジネス要件として、現在の時間の推定値を1分以内に提供し、ウォーターマークがウィンドウの終了を通過した後に適時(on-time)結果を提供し、そのウィンドウ終了後24時間は自動補正を受け付けることが求められます。24時間を超えて到着したデータは警告なしに消失してはならず、監査およびオフライン照合に送られます。スループット、適時性、および補正期間は問題の前提条件(入力)であり、ストリーム処理エンジンの性能主張ではありません。

プロデューサーは安定した event_id を割り当て、そのIDを持つ有効なリトライは同じビジネスコンテンツを持つと仮定します。未加工イベントは再送可能なストレージに残ります。スコープは、時間セマンティクス、ウォーターマーク、ウィンドウトリガー、重複排除、結果の更新、状態容量、復旧、および照合をカバーします。ブローカーの選定と通貨換算はスコープ外です。これはデータエンジニアリングの設問です。目的は、完全性、レイテンシ、状態コストのトレードオフを明示する、検証可能なデータコントラクトを作成することです。

面接官が評価するポイント

最初のシグナルは、候補者が3つの時計を区別しているかどうかです。event_time はイベントが含まれる営業時間を決定します。processing_time はエンジンがそれを認識する時刻を示し、1分ごとの早期リフレッシュを駆動できます。ウォーターマークは、イベント時間の進捗に関するエンジンの推定量です。処理時間がウィンドウを決定する場合、同じ履歴を別の時間に再生すると異なる時間別結果が生成される可能性があります。

2番目のシグナルは、ウォーターマークを絶対的な保証ではなく、進捗の推定量として扱っているかどうかです。ウォーターマークを素早く進めると適時な出力が得られますが、より多くのレコードが遅延として分類されます。ウォーターマークを遅らせると完全性は向上しますが、ウィンドウのクローズが遅れ、より多くの状態が保持されます。優れた回答は、固定の5分や1時間の遅延を暗記するのではなく、実測された到着遅延データと補正SLOからポリシーを導き出します。

3番目のシグナルは、ウォーターマーク、許容される遅延、重複排除の保持期間を分離していることです。ウォーターマークは適時(on-time)発火を制御します。許容される遅延は、ウィンドウ状態が補正可能である期間を制御します。重複排除状態は、イベントの別のコピーが到着する可能性のある期間をカバーする必要があります。これらの値は関連している場合もありますが、単一の設定ではありません。

最後に、出力が収束する必要があります。遅延イベントにより、1つのウィンドウに対して複数の結果が生じます。すべてのスナップショットを追記すると、下流の合計でウィンドウが重複してカウントされます。結果には、リビジョンまたは確定マーカーを備えた、ウィンドウをキーとする冪等な upsert が必要です。チェックポイントはオペレーターの状態を保護します。シンクがチェックポイントとともにコミットできない場合でも、設計にはトランザクション書き込みまたは冪等なバージョン置換が必要です。

回答前に明確にすべき質問

  • ビジネスメトリクスは発生ベースですか、それとも到着ベースですか? この問題では購入 event_time を使用します。処理時間が適切なのは、メトリクスが明確に「システムによって現在処理されたリクエスト」である場合のみです。
  • 1分の結果は推定値ですか、それとも最終値ですか? ここでは推定値であるため、処理時間による早期トリガーが有効です。表示されるすべての結果が完全でなければならない場合、システムはより長く待機するかバッチ計算を使用する必要があります。
  • 24時間後にレポートを変更できますか? 24時間はストリーミングジョブの自動補正期間です。それ以降のデータは照合に送られます。すべての遅延レコードを補正するという規制上または決済上の要件がある場合、状態のクリーンアップがデータの最終的な真実を定義することはできません。
  • 重複はIDで識別できますか? ここには安定した event_id が存在します。ユーザー、金額、時間から推測すると、正当な購入が削除されたり、重複が見逃されたりします。まずプロデューサーのコントラクトを修正してください。
  • 同じIDが異なるコンテンツを持つ場合はどうしますか? first-write-wins や last-write-wins を選択しないでください。プロデューサーが冪等性コントラクトに違反したため、ペイロードフィンガープリントの不一致を記録し、隔離(quarantine)します。
  • シンクはスナップショット、差分(delta)、リトラクション(retraction)のどれを受け入れますか? この設計では完全な自治スナップショットを出力し、(window_start, dimensions) で upsert します。追記専用(append-only)シンクには、読み取りレイヤーが最新リビジョンを選択するバージョニングされたチェンジログが必要です。
  • イベントのタイムスタンプは信頼できますか? はるか未来のタイムスタンプ、RAWデータの保持期間より古いもの、または合意されたタイムゾーンに対して無効なものは隔離します。最大イベント時間に含まれる不正な未来のタイムスタンプは、ウォーターマークを過度に進めてしまう可能性があります。

30秒の回答

event_time ごとにUTC1時間ウィンドウを割り当て、アクティブなソースパーティションごとにウォーターマークを生成し、安全な最小進捗で進めます。処理時間トリガーが毎分推定値を出力し、ウォーターマークが適時(on-time)結果を出力し、24時間の補正のために状態を保持します。event_id とペイロードフィンガープリントでイベントを重複排除し、シンクはウィンドウとリビジョンで upsert します。遅すぎるデータはサイド出力と日次照合に送られます。復旧は、再生可能なソース、チェックポイント、およびトランザクション対応または冪等なシンクを組み合わせて行います。」

ステップごとの詳細解説

まず結果コントラクトから始めます。[10:00, 11:00) のような半開区間ウィンドウを使用し、キーに少なくとも window_start とレポートディメンションを含めます。すべての出力は、その現在バージョンにおけるそのウィンドウの完全な売上スナップショットです。

text
HourlyRevenue {
  window_start
  window_end
  revenue
  revision
  result_state  // EARLY | ON_TIME | FINAL
}

revision はウィンドウに対して単調増加し、シンクはより大きなリビジョンのみを受け入れます。EARLY は1分ごとのリフレッシュであり、完全性を主張しません。ON_TIME は、ウォーターマークがウィンドウの終了を通過したことを意味します。FINAL は、24時間の自動補正期間が終了したことを意味します。FINAL は運用上のコントラクトであり、これ以上のデータが存在しないという主張ではありません。カットオフを超えたレコードは照合によって処理されます。

制約から時間ポリシーを導き出します。アクティブなソースパーティションごとに、検証済みの event_time を抽出します。一般的な有界無順序(bounded-out-of-orderness)戦略は次のとおりです。

text
partition_watermark = max_valid_event_time_seen - out_of_order_bound
operator_watermark = min(active_partition_watermarks)

観測された到着遅延分布、許容される遅延データ率、および ON_TIME レイテンシ目標から out_of_order_bound を選択します。パーティションごとのウォーターマークにより、高速なパーティションが低速なパーティションを完了として宣言するのを防ぎます。複数入力オペレーターは、より古いレコードを生成する可能性のある入力を追い越さないように、最小の進捗を採用します。合意された期間データが生成されなかった後にのみ、パーティションをアイドルとしてマークします。そうしないと、グローバルウォーターマークが永久に停止する可能性があります。アイドル閾値が短すぎると、再開されたパーティションからの古いレコードが遅延扱いになるため、サイド出力は引き続き必要です。

ウィンドウには2つのトリガーファミリーがあります。処理時間の早期トリガーは、現在の完全なスナップショットを1分に1回出力します。ウォーターマークが window_end を通過すると、ON_TIME リビジョンを出力します。ウォーターマークが window_end + 24h を通過するまでウィンドウ状態を保持します。有効な遅延イベントごとに集計が更新され、より高いリビジョンが生成されます。クリーンアップ時に FINAL を出力し、状態を解放します。Beam のトリガーと蓄積モデルは、「いつ出力するか」と「ペインが差分か累積スナップショットか」が別個の選択である理由を示しています。この設計では累積スナップショットと upsert を使用するため、下流システムがペインを組み立て直す必要はありません。

集計の前に event_id で重複排除します。event_id → (event_time, payload_fingerprint) を保存します。最初のイベントを通過させ、同じフィンガープリントを持つ別のコピーを破棄し、異なるフィンガープリントを持つ同じIDをコンフリクトストリームに送信します。保持期間は、単なるウィンドウサイズではなく、上流が別のコピーを再送する可能性のある最長期間に依存します。Spark のウォーターマークベースの重複排除セマンティクスでも同様に、最も早い重複と最も遅い重複の間のタイムスタンプギャップよりも長い遅延閾値が必要です。状態クリーンアップが早すぎると、遅れてきた重複が再びカウントされてしまいます。

状態を明示的に見積もります。毎秒20,000イベントのピークが24時間続き、すべてのイベントが一意である場合、ジョブは 20,000 × 86,400 = 1,728,000,000 個のIDを記憶します。エントリあたりわずか40バイトという説明用の論理ペイロードであっても、この上限は約 69.12 GB になります。エンジンオブジェクト、インデックス、状態バックエンド、チェックポイント、およびレプリカによって物理フットプリントが増加します。したがって、状態にはキーベースのパーティショニングと増分チェックポイントが必要であり、保持は再送コントラクトに従います。重複の99%がより短い期間内に到着する場合、より短いストリーミング期間をシンクのビジネスキーおよびオフライン照合と組み合わせることができますが、残余の重複リスクはTTLの背後に隠すのではなく定量化する必要があります。

出力パスは再送に耐える必要があります。理想的には、シンクがチェックポイントトランザクションに参加し、ソース位置、オペレーター状態、およびウィンドウ出力が一緒にコミットされるようにします。そうでない場合は、(window_key, revision) を冪等な upsert にします。クラッシュ後に同じリビジョンまたはより古いリビジョンを再生しても、新しい値を上書きすることはできません。外部データベースの書き込みがチェックポイント確認の外部にある場合、ブローカーの「exactly-once」というラベルだけでは不十分です。エンドツーエンドのセマンティクスは、ソース、状態、およびシンクの共有コミット境界に依存します。

24時間を超えるイベントは遅すぎるデータのサイド出力に送られ、RAWログに残ります。日次バッチジョブは、同じIDと金額のルールを使用して影響を受けるウィンドウを event_time ごとに再計算し、ストリーミングの FINAL スナップショットと比較します。重大な差異がある場合は、より高いリビジョンを作成するか、財務承認に送られます。バッチジョブはウィンドウをアトミックに置換するか、バージョニングされた upsert を実行します。バッチ結果を既存の合計に加算すると、バッチが再実行されたときに同じデータが2回カウントされてしまいます。

最小限のシーケンスを使用してセマンティクスを検証します。ウィンドウ [10:00, 11:00) は、A=100、B=50、そしてAの同一のコピーを受信します。重複排除された ON_TIME の結果は 150 です。ウォーターマークが11:00を通過した後、許容される遅延期間中に C=20 が到着し、より高いリビジョンによってウィンドウが 170 に変更されます。Dはウォーターマークがクリーンアップポイントを通過した後に到着し、クリアされた状態を直接変更するのではなく照合に送られます。Aの2番目のペイロードの金額が 120 である場合、コンフリクトストリームに送られます。結果は 170190 になってはなりません。

テストマトリクスには、パーティション間の順序の乱れ、ウォーターマークを停止させるアイドルパーティション、未来のタイムスタンプの隔離、ウォーターマーク直前および直後の境界、シンク書き込み前後のクラッシュ、チェックポイントの復旧と再生、シンクに到達する順序不同のリビジョン、および繰り返しの照合実行も含まれます。本番メトリクスには、p50/p95/p99 およびテールの到着遅延、オペレーターごとの現在のウォーターマークと遅延(lag)、早期/適時/遅延/遅すぎるカウントと金額、重複排除ヒットとフィンガープリントのコンフリクト、状態バイト数、チェックポイント時間、拒絶されたシンクリビジョン、バッチとストリームの差分が含まれます。

優れた回答例

「1分の数値を推定値とし、24時間を自動補正期間として定義します。event_time は購入をUTC時間に割り当て、処理時間は1分の早期トリガーのみを駆動します。アクティブな入力パーティションごとに、検証された最大イベント時間から無順序許容量を引いてウォーターマークを生成し、下流では最も遅いアクティブな入力を使用します。アイドル状態のパーティションは明示的にマークされ、無効な未来のタイムスタンプは隔離されます。

集計の前に、イベント時間とペイロードフィンガープリントを保存して event_id で重複排除します。同一のリトライは1回としてカウントされ、異なるコンテンツを持つ同じIDはコンフリクトストリームに入ります。自然な時間ウィンドウは、毎分完全な EARLY スナップショットを出力し、その終了がウォーターマークの後ろになった後に ON_TIME スナップショットを出力し、24時間状態を保持します。その期間内の遅延イベントは合計を更新しリビジョンを増やします。クリーンアップ時に FINAL を出力します。

シンクは、各スナップショットを差分として加算するのではなく、ウィンドウキーとリビジョンによって upsert します。ソースは再生可能であり、オペレーター状態はチェックポイントから復元されます。シンクがチェックポイントトランザクションに参加できない場合、バージョニングされた書き込みにより、復旧時の再送が新しい結果を上書きするのを防ぎます。24時間より遅いイベントはサイド出力に送られ、日次バッチジョブがRAWログを再計算して FINAL と比較します。

毎秒20,000イベントのピークが24時間続き、すべてのIDが一意である場合、上限は17億2800万エントリの重複排除状態になります。エントリあたり論理サイズ40バイトでも69.12GBとなるため、測定された遅延分布によって保持を正当化する必要があり、状態とチェックポイントのコストを監視する必要があります。最後に、重複、順序不同、遅延、未来のタイムスタンプ、アイドルパーティション、復旧に対するフォールトインジェクションを行い、最終的なウィンドウがオフラインの再計算と一致することを検証します。」

よくある間違い

  • 処理時間でビジネスウィンドウを割り当てる → 履歴の再送中に同じイベントが異なる時間に入る可能性があります → 所属にはイベント時間を使用し、早期出力にのみ処理時間を使用します。
  • ウォーターマークを「これ以上古いデータは到着しない」と説明する → 通常はヒューリスティックな進捗推定量であるため、より古いイベントが依然として現れる可能性があります → 許容される遅延、遅すぎるデータのサイド出力、および照合を定義します。
  • ウォーターマーク遅延、許容される遅延、重複排除TTLに説明のない単一の値を使用する → これらはそれぞれ発火、ウィンドウ状態、重複認識を制御します → レイテンシSLO、補正期間、上流の再送コントラクトから導き出します。
  • 発火のたびに合計を追記する → EARLY、ON_TIME、および遅延発火が繰り返し合算されます → 完全なスナップショットを出力してウィンドウとリビジョンで upsert するか、リトラクション可能な差分プロトコルを定義します。
  • グローバルな最大イベント時間から進める → 高速なパーティションや不正な未来のタイムスタンプにより、低速パーティションのデータが早すぎるタイミングで遅延扱いになります → パーティションごとの進捗を生成し、アクティブ入力の最小値を採用し、タイムスタンプを検証します。
  • 永続的に静止しているパーティションを最小値計算に残す → ウォーターマークが停止し、ウィンドウや重複排除状態がクリーンアップされなくなります → 観測可能なアイドル処理を使用し、再開された遅延データを正しくルーティングします。
  • イベントIDのみを保存する → プロデューサーによるIDの再利用が無言で飲み込まれます → ペイロードフィンガープリントも保存し、コンフリクトを隔離します。
  • 「exactly-onceが有効になっている」とだけ言う → 外部シンクがチェックポイント境界を共有していない可能性があります → エンドツーエンドのコミットパスを追跡し、トランザクション対応シンクまたはバージョニングされた冪等書き込みを使用します。
  • 状態クリーンアップ後にデータを破棄する → ストリームは安定して見えても、財務結果が説明不能になります → RAWイベントを保持し、サイド出力、バッチ再計算、および不一致レビューを実装します。

フォローアップの質問と回答

フォローアップ1:ウォーターマークの無順序許容量はどのくらいにすべきですか?

本番環境で processing_time - event_time を測定し、ソース、クライアントバージョン、地域ごとにセグメント化します。まず、許容可能な最新の ON_TIME 結果と補正パスに許容されるイベントの割合を決定します。次に、その両方を満たすパーセンタイルを選択します。p50、p95、p99、およびテールを監視し続けます。単一の外れ値によってポリシーが自動的に動くのではなく、分布の変化は構成のリリースを通じて行うべきです。ウォーターマークは適時発火を制御し、24時間の補正期間はより長いテールをカバーします。

フォローアップ2:なぜ1つの静止したKafkaパーティションがウィンドウを停止させる可能性があるのですか?

複数入力オペレーターの安全な進捗は、最小入力ウォーターマークです。進まないパーティションがあると、その最小値は変更されません。アイドル閾値の間イベントを生成していないことを確認した後、アイドルとしてマークして一時的に最小値から除外します。閾値を短く設定しすぎないでください。再開されたパーティションからのレコードが新しいウォーターマークより古くなる可能性があります。それらは許容される遅延処理またはサイド出力に入る必要があり、アイドル状態の遷移には独自のメトリクスが必要です。

フォローアップ3:なぜ24時間の重複排除状態で重複が永久に不可能であると保証できないのですか?

24時間は、この問題の自動補正期間のみをカバーします。上流が25時間目にIDを再送した場合、ストリーミング状態は消去されているため、イベントが集計に再進入する可能性があります。永続的な重複排除には、より長期間有効なビジネス一意インデックス、クエリ可能なイベントレジストリ、または完全なオフライン重複排除が必要であり、それぞれストレージと書き込みのコストが増加します。「宣言された再送期間内での重複排除」として保証を定義し、それ以降のレコードは照合します。

フォローアップ4:下流のデータベースが INSERT をサポートし、upsert をサポートしていない場合はどうしますか?

すべての結果を、ウィンドウとリビジョンをキーとするイミュータブルなチェンジログに書き込みます。読み取りレイヤーは、各ウィンドウの最大リビジョンから現在のビューを構築します。コンシューマーは、各レコードが差分ではなく完全なスナップショットであることを認識し、すべてのリビジョンを合算してはなりません。クエリコストが高すぎる場合は、監査と復旧のためのバージョンログを保持しながら、アトミックな置換をサポートする配信テーブルに非同期でコンパクションします。

フォローアップ5:復旧時に過大カウントも過小カウントも発生しないことをどのように証明しますか?

固定入力に対して、期待される重複排除セットとウィンドウリビジョンを記録します。3つの境界でジョブを終了させます:ソース読み取り後で状態チェックポイントの前、状態チェックポイント後でシンク確認の前、およびシンク書き込み後でチェックポイント確認の前。復旧後に再生し、各イベントIDが1回だけ寄与していること、シンクが最大リビジョンのみを保持していること、そして最終ウィンドウがオフライン再計算と一致することを検証します。これらのコミット境界でのフォールトインジェクションを行わずに RUNNING ステータスを確認するだけでは、エンドツーエンドの正確性の証明にはなりません。

公開情報ソース

関連する質問