代表的な面接トピック

ストリーミング集計において遅延および順序不同のイベントをどのように処理しますか?

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

質問

5分間のイベント時間集計ジョブを設計してください。イベントは順序不同、遅延、または一時的にアイドル状態のパーティションから届く可能性があります。ウォーターマーク、遅延データポリシー、結果の修正、およびオブザーバビリティについて説明してください。

1. 質問

注文イベントストリームにおいて、customerId ごとの金額と注文数を5分間のウィンドウで集計します。デバイスクロックのドリフト、ネットワークの再試行による順序不同の配信が発生する可能性があり、一部の Kafka パーティションには一時的に新しいメッセージが存在しない場合があります。結果を迅速に表示しつつ、一定期間は遅延データによってウィンドウを修正できるようにする必要があります。イベント時間の処理を設計してください。

2. 制約と確認事項

  • ウィンドウがイベント時間、書き込み時間、処理時間のいずれを使用するかを確認します。本問題ではイベント時間を使用します。
  • 30秒などの最大乱れ許容境界を設定し、それを超えた場合の動作を定義します。
  • 結果が追記専用(append-only)イベントか更新可能なスナップショットかを決定します。ダウンストリームのコンシューマはリビジョンを識別できる必要があります。
  • アイドル状態のパーティション、無効なタイムスタンプ、重複イベント、リプレイジョブについて検討し、メッセージがない入力がストリーム終了と誤認されないようにします。

3. コアコンセプト

イベント時間はレコード自体に由来します。ウォーターマークは、システムがイベント時間が特定の進捗位置まで進んだと判断していることを示します。通常、ウォーターマークがウィンドウの終了を通過したときにウィンドウがトリガーされます。そのウォーターマーク以降に到着したレコードで、タイムスタンプがまだそのウィンドウに属しているものは遅延データとなります。並行入力を結合する場合、オペレータは通常、最小の入力ウォーターマークを待機するため、非アクティブなパーティションがグローバルな進捗を停滞させる可能性があります。そのため、アイドル検出またはパーティションごとのタイムアウトポリシーが必要です。

4. 参照フロー

text
onRecord(event):
  ts = extractEventTimestamp(event)
  key = canonicalKey(event.customerId)
  updateWatermarkGenerator(ts)

  window = floorToFiveMinutes(ts)
  if ts <= currentWatermark + allowedLateness:
    state[window, key] = aggregate(state[window, key], event)
    emitUpsert(window, key, state[window, key], revision + 1)
  else:
    routeToLateData(window, key, event)

onWatermark(wm):
  finalizeWindowsBefore(wm)
  expireStateAfterRetention()

ソースはイベントタイムスタンプを抽出し、有界順序不同(bounded-out-of-orderness)ウォーターマークを生成します。データがない状態が長期間続いたパーティションをアイドルとしてマークし、マージされたウォーターマークの進行を妨げないようにします。許容遅延期間中はウィンドウ状態を保持し、upsert または修正イベントによって結果を更新します。期限を超えたデータは、確認やオフラインバックフィルのためにサイド出力(side output)へルーティングします。

5. 精度とレイテンシのトレードオフ

許容遅延を長くするとイベント時間のビューはより完全になりますが、状態サイズと修正の量が増加します。最終結果を1回だけ出力するとダウンストリームのレイテンシは最小限に抑えられますが、遅延による修正を表現できません。ウィンドウキー、リビジョン、および理由を含む更新イベントは、ダウンストリームでのべき等なマージを必要とします。金額のような損失を許容できない指標の場合、生イベントを保持してオフラインの再計算をスケジュールします。リアルタイムのリーダーボードでは、期限後の次のバッチでのみ更新することで十分な場合があります。

6. 検証とオブザーバビリティ

  • 順序通り、順序不同、30秒未満の遅延、期限超過のイベントを生成し、各ウィンドウをオフラインのベースラインと比較します。
  • アイドルパーティション、クロックジャンプ、重複イベント、タスク再起動を注入し、ウォーターマークが停滞したり逆行したりしないことを検証します。
  • 現在のウォーターマーク、処理時間とイベント時間のラグ、ウィンドウ状態のサイズ、遅延イベント率、サイド出力のボリューム、修正回数を記録します。
  • すべての更新にウィンドウキー、リビジョン、入力イベント ID を含め、同じログをリプレイして最終スナップショットを比較します。

7. よくある間違い

  • ネットワークの乱れに強いと主張しながら、イベント時間ウィンドウを処理時間ウィンドウに置き換えてしまうこと。
  • ウォーターマークを絶対的な完了マーカーとして扱うこと。実際には遅延の仮定に基づくヒューリスティックにすぎません。
  • アイドルパーティションを無視し、無音のパーティションによってすべてのウィンドウのクローズがブロックされること。
  • サイド出力、バージョン管理、またはオフラインバックフィルパスを設けずに遅延データを破棄すること。

8. 面接の評価ポイント

3つの時間の概念を区別できているか

候補者はイベント時間、取り込み時間、処理時間を定義し、ウィンドウの選択が精度とレイテンシをどのように変化させるかを説明する必要があります。

ウォーターマークの生成とマージができるか

回答には、有界順序不同、並行入力にわたる最小ウォーターマーク、およびアイドル処理が含まれ、ウォーターマークが絶対的な保証ではなく進捗の推計であることを指摘する必要があります。

遅延データと修正パスを設計できているか

候補者は、期限内および期限後のデータの両方について、許容遅延期限、サイド出力、リビジョン、およびダウンストリームでのべき等マージを定義する必要があります。

リプレイによる結果検証ができるか

候補者は、ジョブが稼働し続けているかどうかを確認するだけでなく、オフラインのベースラインとライブメトリクスに対して、順序不同、アイドル状態、再起動、重複をテストする必要があります。

公開情報ソース

関連する質問