問題とスコープ
あるEコマースプラットフォームは、注文イベントを永続的でリプレイ可能なログに書き込み、イミュータブルな生データを保持します。定常レートは指定されていませんが、ピークは毎秒20,000イベントであり、プロモーションによる短時間のスパイクではその最大10倍に達します。在庫の異常検知はイベント発生から10秒以内のアラートが必要です。運用ダッシュボードは5分以内に更新される必要があります。財務部門は翌日に帳簿を締め、同一の営業日を再実行可能かつ照合可能にする必要があります。イベントの約2%は最大6時間遅れて到着します。
これらのレート、スパイク倍率、遅延割合、および期限は問題の前提条件であり、特定のエンジンのパフォーマンスを主張するものではありません。各イベントには安定した event_id、order_id、event_time、およびスキーマバージョンが存在すると仮定します。財務部門は明確なタイムゾーンの営業日を使用します。4人のデータチームはすでにウェアハウスとスケジューラを安定運用していますが、ステートフルなストリーミングジョブに対する成熟したオンコール運用体制はまだありません。
この問題は、3つのコンシューマーすべてに対して単一の処理モードを適用することを求めてはいません。目標は、アクションの期限、結果の完全性、リカバリ、および保守コストから、必要十分で最小限の設計を導き出すことです。この問題の本質はドメイン横断的なプラットフォームアーキテクチャではなく、処理セマンティクス、鮮度、およびデータ品質のトレードオフにあるため、data に該当します。
面接官が評価するポイント
最初のシグナルは、候補者がメッセージキューを見ただけで安易にストリーミングを選択するのではなく、最終的に役立つビジネスアクションから逆算して考えているかどうかです。アンバウンデッド(無界)なデータソースは、データが到着し続けることだけを示しています。すべてのダウンストリームの結果が1レコードずつ処理されることを要求しているわけではありません。翌日の財務処理は、同じログのバウンデッド(有界)なスナップショットを読み取ることで、バッチ処理による高い再現性の恩恵を受けることができます。
2つ目のシグナルは、処理レイテンシとアクションレイテンシを分離できているかです。エンジンが500ミリ秒で計算できたとしても、シンクの更新が5分間隔であれば、秒単位のダッシュボードは実現できません。優れた回答では、取り込み、キューイング、計算、書き込み、キャッシュ更新、アラート配信のバジェットを個別に配分し、各セグメントを測定します。
3つ目のシグナルは、正確性の境界を明確に定義しているかです。ストリーミングエンジンにおける Exactly-Once 処理は、遅延データが完全であることを証明するものではなく、外部の副作用を自動的に保護するわけでもありません。再実行可能なバッチジョブも、自動的に安全であるとは限りません。定義された入力範囲、ビジネスキー、スナップショットバージョン、およびアトミックな公開がなければ、再実行によって定義が重複したり混在したりする可能性があります。
最後に、面接官は運用上の判断力をテストしています。継続的ストリーミングには、持続的なキャパシティ、バックログリカバリ、状態管理、チェックポイント、安全なデプロイ、およびオンコール体制が必要です。マイクロバッチで5分の期限を確実に満たせる場合、レコードごとの状態管理は判断を変えることなく障害モードを増やすだけです。逆に、1時間ごとのバッチ処理では、実際の補充や販売停止のトリガーとなる10秒のアラートの価値が失われます。
回答前に明確にすべき質問
- 期限の開始と終了はどこか? ここでは、ソースシステムがイベントをコミットした時点で開始し、アラートが到達するかクエリ結果が表示された時点で終了します。ウェアハウスへの取り込みで完了とする約束は、シンクとキャッシュの遅延を無視しています。
- 10秒のアラートは自動アクションをトリガーするか? 在庫の自動変更、販売停止、または通知には、重複排除、冪等性、および監査が必要です。監視目的のアラートであれば、より多くの誤検知や重複を許容できます。
- 5分のダッシュボードは推計値か、完全な確定値か? 「基準日時(as of)」を持つ改訂可能な推計値であれば、マイクロバッチで十分です。すべての表示に6時間の遅延テールを含める必要がある場合、5分の鮮度と完全性の要件が競合するため、製品仕様の契約を変更する必要があります。
- 財務部門はいつ締めを確定(フリーズ)し、後から調整できるか? 翌日にドラフトを作成して後から調整仕訳を入れるのか、フリーズする前に6時間待つのかでは異なります。この回答によってバッチのカットオフと改訂プロトコルが決まります。
- 3つの用途で未加工レイヤーと変換定義を共有できるか? 正規化されたイベント、ビジネスキー、およびテストフィクスチャは共有可能です。フィルタ、タイムゾーン、金額計算のルールを個別にコピーすると、最終的にバッチとストリーム間で乖離(ドリフト)が発生します。
- チームはステートフルなストリームを24時間体制で運用できるか? バックログアラート、チェックポイントリカバリの訓練、および安全なデプロイ体制がない場合、継続的ストリーミングは秒単位のアクションが真に必要な最小限のパスに限定すべきです。
30秒の回答
「ソース全体に単一のモードを適用するのではなく、ビジネスアクションの期限ごとにコンシューマーを分割します。10秒の在庫異常検知パスには継続的ストリーミングを使用します。5分の運用ダッシュボードには1〜2分のマイクロバッチを使用し、書き込みとキャッシュ更新のバジェットを確保します。翌日の財務処理にはバッチで営業日のスナップショットを使用し、照合の基準とし続けます。3つすべてがイミュータブルな生イベント、正規化ルール、およびビジネスキーを共有します。
ストリーミングの結果は改訂可能なビューであり、遅延データが完全であると主張するものではありません。外部アクションは、イベントおよびルールバージョンに対して冪等にします。本番切り替え前に、バッチ、マイクロバッチ、ストリームの各パスを通じて同じ履歴をリプレイし、件数、金額、テールレイテンシをシャドー比較した上で、秒単位のパスのみアクションを有効にします。アクションのSLOを満たし、リカバリが収束する最もシンプルなモードを選択するのが原則です。」
ステップごとの詳細解説
ステップ1: 3つの結果契約を定義する
Spark、Flink、その他の製品を選択する前に、各出力のキー、期限、完全性、および改訂動作を定義します。
| 用途 | 結果キー | 表示期限 | 完全性と改訂 | 推奨モード |
|---|---|---|---|---|
| 在庫異常検知 | item, rule version, time window | 10秒 | 高速に改訂可能、アクションは冪等 | continuous stream |
| 運用ダッシュボード | metric, dimensions, window | 5分 | 基準日時を表示、遅延データはバージョンを置換 | 1–2 minute micro-batch |
| 財務締め処理 | business day, account, currency | 翌日 | フリーズされたスナップショットを再計算、不整合を調整 | batch |
1つのイベントが3つの異なる契約に利用されます。メッセージキューは入力のトランスポートに過ぎず、3つの行すべてに同じ実行モードを強制することはできません。
ステップ2: エンドツーエンドのアクションバジェットを割り当てる
10秒のアラートに対するテスト可能な初期バジェットとして、ソースのコミットとトランスポート、キューイング、計算、シンクとルールアクション、アラート配信にそれぞれ2秒を割り当てます。負荷テストによってこれらの暫定的な配分を更新する必要がありますが、合計がビジネス上の期限を超えてはなりません。各セグメントの p95、p99、および最大バックログ経過時間を記録します。オペレーターの実行時間だけを見ていても、遅いシンクを見逃してしまいます。
5分のダッシュボードは、1分または2分ごとにマイクロバッチを開始できます。2分のスライスの場合、スケジュールの待機時間は最悪で約2分、計算と書き込みに各1分、キャッシュの更新に最後の1分が割り当てられます。10倍のスパイクによって計算がバジェットを超過した場合は、まず並列度を上げるか、スライスを短くするか、優先度の低いディメンションの計算を後回しにします。平均実行時間はテールSLOの証明にはなりません。
財務処理は、最小レイテンシではなく、再現性とフリーズされた定義によって制約されます。バッチは明示的なスナップショットまたはオフセット範囲を読み取り、run_id でラベル付けされたステージング領域に結果を書き込み、検証してアトミックに公開します。同じ入力とルールバージョンを再実行すれば、同じ結果が得られる必要があります。
ステップ3: 二重実装の負債を作らずにファクトを共有する
生イベントは、元のペイロード、スキーマバージョン、取り込み時刻、ソース位置とともに、リプレイ可能なログまたはオブジェクトストアに最初に入力されます。正規化レイヤーが、スキーマの進化、タイムゾーン、金額の単位、キャンセルステータス、およびビジネスキーを一貫して処理します。各実行モードは、その境界の後から読み取ります。
ファクトを共有することは、3つのエンジン間でコードを1行ずつ共有することを意味しません。より有用な境界は、共有データ契約、ゴールデンフィクスチャ、および決定論的なルールです。バッチとストリームでウィンドウロジックを個別に実装する必要がある場合は、同じフィクスチャに対して等価性テストを実行します。財務部門は独立したより厳格なチェックを保持できます。コードの再利用によって統制(コントロール)を排除すべきではありません。
ステップ4: 重複、遅延データ、副作用を個別に設計する
在庫ストリームは event_id で重複を排除し、短いイベント時間ウィンドウを計算して、ルールと結果のバージョンを出力します。販売停止、補充、または通知では、安定したアクション冪等性キーを使用します。チェックポイントの成功は外部HTTP呼び出しが1回だけ実行されたことを証明するものではないため、再試行時にはすでに実行されたアクションを認識できる必要があります。
運用マイクロバッチは、固定された半開区間のソースオフセット範囲を読み取り、別の合計値を追加(アペンド)するのではなく、バージョン管理されたターゲットウィンドウを置換します。遅延イベントは後続のバッチに入り、結果のバージョンを上げます。ダッシュボードには「基準日時」が表示され、5分の鮮度が6時間のテールの完全性を意味しないことを明確にします。
財務バッチは、合意されたカットオフ後にフリーズされた入力を読み取ります。それ以降のイベントは、元の実行、入力範囲、ルールバージョン、および不整合の承認を保持する調整実行に入ります。Exactly-Once ストリーム処理は、特定の境界内で永続化される処理結果の重複を防ぐことができますが、遅延データが存在する場合の完全性を証明するものではなく、外部の副作用に自動的に拡張されることもありません。
ステップ5: キャパシティとコストの再現性を確保する
規定のピークは毎秒20,000イベントであるため、10倍のスパイクは毎秒200,000イベントになります。正規化されたイベントが仮に 1 KiB である場合、ピーク時の取り込みは約 195 MiB/秒 になります。これはキャパシティの概算に過ぎず、本番のサイジングは圧縮および非圧縮のサイズ分布から再計算する必要があります。
継続的ストリーミングには、ピークトラフィックとバックログリカバリの両方に対する持続的なキャパシティが必要です。毎秒20,000イベントの状態で10分間障害が発生すると、1,200万件のイベントがキューに滞留します。復旧後の処理能力が現在の入力と同等である場合、バックログは決して解消されません。マイクロバッチは、次のスライスが到着する前に現在のスライスを完了する必要があります。バッチはより安価な時間帯にリソースを集中させることができますが、極端に小さく頻繁なバッチは起動、コミット、小さなファイルのオーバーヘッドが支配的になります。
比較対象には、常時稼働の計算リソース、状態およびチェックポイントのストレージ、シンクへの書き込み、スキャン、オンコール対応、デプロイの複雑さが含まれます。クラウドの利用料金だけでは、4人のチームが3つの類似した実装を保守するコストが見落とされます。推奨される設計では、継続的ストリーミングを在庫アラートのみに限定し、財務やダッシュボードのために24時間体制の状態管理を行うことを回避します。
ステップ6: 一括切り替えではなくシャドーリプレイで移行する
通常トラフィック、10倍のスパイク、重複、順序の乱れ、および6時間の遅延テールを含む履歴セグメントをフリーズします。従来のバッチ結果をベースラインとして使用し、新しいマイクロバッチおよびストリームジョブを在庫アクションをトリガーしないシャドーモードで実行します。ビジネスキーごとに件数、金額、バージョン、遅延改訂履歴を比較し、すべての不整合の理由を明らかにします。
段階的に公開範囲を拡大します。シャドーテーブルへの書き込みのみを行い、社内ダッシュボードを公開し、その後ストリームアラートによる可逆的なアクションのトリガーを許可します。フォールトインジェクションでは、チェックポイント前後でのクラッシュ、シンクのタイムアウト、パーティションの停止、スキーマ変更、バックログの追いつきを検証します。受け入れメトリクスには、エンドツーエンドの p99、最大バックログ経過時間、マイクロバッチ完了時間、遅延改訂率、重複アクション数、バッチとストリームの金額差、およびリカバリ時間が含まれます。
終了基準も定義します。マイクロバッチが5分の期限を一貫して逃す場合は、継続的ストリーミングを検討する前に、スケジューリング、スキュー、シンク、およびスパイク処理能力を調査します。10秒のアラートがアクションをトリガーしなくなった場合は、オンコールコストを削減するためにマイクロバッチにダウングレードします。処理モードは検証可能なビジネス上の選択肢であり、恒久的なアイデンティティではありません。
優れた回答例
「私はアンバウンデッドな入力をその処理モードから切り離して考えます。これらのコンシューマーはアクションの期限が異なるため、技術的な統一性だけを目的にすべてをストリーミングに強制することはありません。在庫異常検知は実際に10秒以内にアクションを起こす必要があるため、継続的ストリーミングを使用し、ソース、キュー、計算、書き込み、配信にバジェットを割り当て、event_id とルールバージョンによって各アクションを冪等にします。運用ダッシュボードは5分で十分であるため、ウィンドウとバージョンによって結果を置換し、基準日時を表示する1〜2分のマイクロバッチから始めます。財務の締め処理は、バッチでフリーズされた営業日スナップショットを読み取り、run_id ごとにステージングし、検証してアトミックに公開します。遅延データは調整実行に入ります。
3つのパスはリプレイ可能な生イベント、正規化された契約、およびゴールデンフィクスチャを共有し、財務部門は独立した統制を維持します。毎秒20,000イベントで短時間の10倍のスパイクが発生する場合、定常状態だけでなく、毎秒200,000イベントの取り込みと障害後の追いつき率をテストします。ストリームのチェックポイントは外部アクションの冪等性を代替するものではなく、Exactly-Once 処理は6時間の遅延テールが完全であることを証明しません。
移行については、同じ履歴をリプレイし、シャドーテーブルでバッチ、マイクロバッチ、ストリームの各パスのキーレベルの結果と改訂を比較します。継続的ストリーミングは監視アラートから開始します。p99が10秒未満になり、重複アクションがゼロになり、バックログが目標時間内に解消されることを確認した後にのみ、自動アクションをトリガーします。私の判断基準は、アクション期限、完全性、およびリカバリ契約を満たす最もシンプルなモードを選択することです。低レイテンシ化による状態管理とオンコールのコストは、それがビジネス上の意思決定を変える場合にのみ正当化されます。」
よくある間違い
- Kafka があるからといってすべてをストリーミングに移行する → 継続的な入力であってもすべての結果を1レコードずつ処理する必要はないため、財務やダッシュボードにメリットのない状態管理と運用コストが発生する → コンシューマーのアクション期限ごとに選択する。
- ストリーミングは低レイテンシであるとだけ述べる → シンク、キャッシュ、通知配信がバジェット全体を消費する可能性がある → エンドツーエンドのパスと各テールレイテンシを測定する。
- 5分のビューを確定値と呼ぶ → イベントは最大6時間遅れて到着する可能性がある → 推計、改訂、フリーズ、調整のセマンティクスを定義する。
- チェックポイントを Exactly-Once の外部アクションとして扱う → 再試行によって在庫変更や通知呼び出しが重複して実行される可能性がある → アクション冪等性キー、監査ログ、リプレイテストを使用する。
- 各マイクロバッチから新しい合計値を追加する → 同じウィンドウが二重にカウントされる → ウィンドウとバージョンによってアトミックに置換するか、取り消し可能な差分(retractable-delta)プロトコルを定義する。
- バッチとストリームにビジネスルールを個別にコピーする → タイムゾーン、キャンセル、金額の定義が乖離する → 契約とフィクスチャを共有し、キーレベルで比較する。
- 定常入力のみを考慮してサイジングする → 10倍のスパイクや障害によるバックログを期限内に解消できなくなる → ピーク、正味のリカバリレート、シンク容量を併せてテストする。
- 初回のデプロイで実際のアクションを有効にする → セマンティクスの不整合が直接在庫を変更してしまう → 可逆的なアクションを段階的に有効にする前に、シャドー書き込みとアラート監視を行う。
フォローアップ質問と回答
フォローアップ1: なぜ5分のダッシュボードに継続的ストリーミングを使用しないのですか?
ピーク時およびリカバリ時においても1〜2分のマイクロバッチが安定して完了し、最終的な書き込みとキャッシュ更新が5分以内に収まるのであれば、継続的ストリーミングを導入しても運用上の判断は変わりません。むしろ、長期的な状態管理、チェックポイント、バックログリカバリ、およびデプロイコストが増加します。将来的に期限が30秒に短縮されたり、スケジューリングと起動がマイクロバッチのバジェットの大半を消費するようになったりした場合は、同じ履歴をリプレイして継続的ストリーミングと比較検討します。「リアルタイム」という言葉ではなく、SLOに対する根拠に基づいてアップグレードを判断します。
フォローアップ2: 単一のストリームでアラート、ダッシュボード、財務結果をすべて生成できますか?
生成自体は可能ですが、1つのチェックポイント、スキーマ変更、または不正なウィンドウ処理によって3つの用途すべてがブロックされてはなりません。アラートは正規化されたストリームを読み取って短い状態を保持し、ダッシュボードはバージョン管理された集計を読み取り、財務はフリーズされたスナップショットでイミュータブルな履歴から再計算できます。デプロイと障害ドメインを分離しながら、入力と定義を共有します。障害の影響、バックフィル、および監査が許容範囲内であることが証明された後にのみ、ランタイムユニットを統合します。
フォローアップ3: Exactly-Once 処理によって、ストリーム出力を財務の締め処理にそのまま適したものにできますか?
そのラベルだけでは不十分です。処理セマンティクスによってコミットされた出力が再試行時に重複するのを防ぐことはできますが、遅延レコードがすべて到着したことを保証するものではなく、外部の副作用を自動的にカバーするわけでもありません。財務部門には依然として、フリーズされた入力範囲、ルールバージョン、再現可能な実行、元帳の制約、および不整合の承認が必要です。ストリーム出力は初期の推計値として役立ちます。それらの監査統制が満たされ、長期的な照合によって等価性が証明された場合にのみ、締め処理の入力となります。
フォローアップ4: 10倍のスパイクからのリカバリをどのように証明しますか?
スパイクの継続時間と障害の継続時間を記録し、キューに滞留したイベント数を計算して、正味の消化レートを測定します。現在の入力が毎秒20,000件でコンシューマーの処理能力が毎秒30,000件の場合、正味の消化レートは毎秒10,000件に過ぎず、滞留した1,200万件のイベントの処理には約20分かかります。受け入れ検証では、エンドツーエンドのアラートレイテンシ、シンクの制限、状態の増大、およびオートスケーリングの遅延も監視します。正味のリカバリ計算を伴わないピークスループットの数値は、ロングテールのSLO違反を見逃す原因になります。
フォローアップ5: 継続的ストリーミングをマイクロバッチにダウングレードすべきなのはどのような場合ですか?
ビジネスが秒単位の出力に基づいてアクションを行わなくなった場合、マイクロバッチが新しい期限を満たせる場合、またはストリーミングのオンコールおよび状態管理コストがそれによって防止できる損失を持続的に上回る場合にダウングレードします。まずマイクロバッチをシャドー実行し、結果とレイテンシを比較します。次に、リプレイ可能な入力とロールバックウィンドウを保持したまま、実際のアクションを無効化します。コスト削減によってデータの陳腐化が見落とされないよう、変更後も遅延改訂とピーク完了時間を監視し続けます。