代表的な面接トピック

ストリーミングにおける Exactly-Once 処理と外部副作用をどのように説明しますか?

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

質問

面接官が「ストリーミングジョブは Exactly-Once でなければならない」と言いました。これは何を意味し得るのか、なぜ外部副作用は依然として重複して実行される可能性があるのか、そして処理、コミット、重複排除、モニタリングをどのように設計するかを説明してください。

1. 質問

注文イベントストリームがパース、ウィンドウ処理、集約され、リスクスコアリングされます。結果はウェアハウスに書き込まれ、ダウンストリームの通知がトリガーされます。ワーカーのクラッシュ、ネットワークのタイムアウト、イベントの遅延が発生する可能性があります。Exactly-Once の境界を説明し、再試行時に重複した請求が発生しないフローを設計してください。

2. 制約と確認事項

  • 3 つのレイヤーを分離します:メッセージ配信、パイプライン内結果、および外部副作用。
  • イベントは重複、順不同、または遅延する可能性があります。処理ログはコミットされた結果の証明にはなりません。
  • 結果は再生可能である必要があり、通知やその他の外部呼び出しには冪等性または重複排除が必要です。
  • 実装を選択する前に、レイテンシ、遅延許容ウィンドウ、およびドロップポリシーを明確にしてください。

3. コアコンセプト

Exactly-Once 処理とは通常、システムがレコードの損失を防ぎつつ、レコードのパイプライン結果が永続ストレージ出力に最大 1 回反映されることを意味します。すべてのユーザー関数が 1 回だけ実行されることを意味するわけではなく、HTTP、メール、またはデータベースの呼び出しを自動的にカバーするわけでもありません。At-Least-Once 入力にチェックポイント、決定論的再生、および結果の重複排除を組み合わせることで、検証可能な出力保証を実現できます。

4. 参照フロー

text
onEvent(event):
  key = stableEventId(event)
  state = readCheckpointOrState(key)
  result = deterministicTransform(event, state)
  writeTransactionalResult(key, result)   # unique(key)
  commitCheckpointAfterResult(key)

onExternalSideEffect(result):
  idempotencyKey = result.eventId + ":" + result.version
  callOrOutbox(idempotencyKey, result.payload)

まず一意性制約またはトランザクションを使用して結果とイベント ID をストレージに書き込み、その後にチェックポイントを進めます。外部通知は、冪等な API、または Outbox パターンと独立した送信アプローチを経由してルーティングします。送信側は再試行を行う場合があり、受信側は特定の冪等性キーを 1 回のみ受け入れます。

5. 障害ケースとトレードオフ

外部呼び出しが成功した後、そのチェックポイントがコミットされる前にワーカーがクラッシュした場合、再生によって外部サービスが再度呼び出されます。冪等性キーがなければ、ランナー単体でその重複した副作用を取り除くことはできません。ウィンドウ結果も遅延データやウォーターマークに依存するため、修正の境界を明確にする必要があります。より強力なエンドツーエンドの保証は、重複排除の状態、トランザクションの調整、およびストレージコストを増加させます。重複が許容される場合、At-Least-Once の方が低いレイテンシを提供できる可能性があります。

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

  • クラッシュ、タイムアウト、重複メッセージ、順不同イベントを注入し、単一のビジネスキーについて最終結果を検査します。
  • 入力イベント ID、試行回数、コミットバージョン、重複排除のヒット、および外部呼び出しの結果を記録します。
  • 受信数、処理数、コミット数、通知数の 4 つのカウントを照合します。ワーカーのログだけでは不十分です。
  • 重複率、遅延度、チェックポイントの経過時間、重複排除状態のサイズ、および再生バックログを監視します。

7. よくある間違い

  • Exactly-Once 配信、Exactly-Once 処理、および Exactly-Once 副作用を単一の約束として扱うこと。
  • フレームワークのスイッチを切り替えるだけで、任意のカスタムコードや外部 API に 1 回限りの効果が付与されると仮定すること。
  • 安定したイベント ID ではなくタイムスタンプを使用して重複排除を行い、再試行や再生時に異なるキーを生成してしまうこと。
  • 遅延イベント、バージョンの競合、および重複排除レコードの保持期間を無視すること。

8. 面接の採点ポイント

保証の境界を描けているか

候補者は配信、パイプライン結果、外部副作用を分離し、フレームワークが実際にカバーしているレイヤーを特定できているか。

再生可能なフローを設計しているか

候補者は安定したイベント ID、決定論的変換、トランザクションによる結果のコミット、チェックポイントの順序付けを使用し、クラッシュリカバリを説明できているか。

外部副作用を処理しているか

候補者は冪等性キー、一意性制約、または Outbox パターンを提案し、送信側と受信側が連携して重複を防ぐ方法を説明できているか。

フォールトインジェクションで検証しているか

候補者は重複、順序の乱れ、遅延、タイムアウト、ワーカーのクラッシュを網羅し、コミットされたデータとビジネス照合を用いて主張を検証できているか。

公開情報ソース

関連する質問