問題と背景
アップストリームシステムから1日あたり2億件の注文イベントが送信され、そのうち約0.2%に必須フィールドの欠落、パースの失敗、またはビジネスルール違反が含まれる可能性があります。有効なイベントはファクトテーブルやダウンストリームのメトリクスへと流れ続けなければなりません。無効なイベントを警告なしに消失させてはならず、修復後は同じイベントを2回カウントすることなく、元の入力から再送できる必要があります。
入力ごとに安定したevent_id、不変の未加工(raw)コピー、スキーマのバージョン管理されたルール、および失敗理由、ルールバージョン、リトライ状態を保持する隔離レコードを備えた、バッチおよびストリーム設計を前提とします。この面接の課題は、単一のSparkオプションではなく、データエンジニアリングにおけるルーティング、オブザーバビリティ、およびリカバリに関するものです。
面接官が見ているポイント
- パース失敗、フィールド欠落、ビジネス検証の失敗を実行可能なクラスに分離できているか。
- 監査および再送のために、未加工のエビデンス、ルールバージョン、およびリネージを保持できているか。
- 影響範囲に基づいて、
fail-fast、drop、redirect/quarantineの中から適切に選択できているか。 - 二重カウントを防ぐために、冪等性キー、重複排除状態、および出力バージョンを活用できているか。
- 品質メトリクス、アラートしきい値、修復ワークフロー、および再送ゲートを設計できているか。
- 隔離領域における有害データ、PII、保持期間、およびアクセス制御に対処できているか。
最初に確認すべき明確化の質問
- 0.2%は許容される不具合バジェットですか、それともいかなる違反もリリースをブロックする必要がありますか?これにより、しきい値ベースかゼロトレランス(許容度ゼロ)のゲーティングかが決まります。
- イベントのリトライ、順不同の到着、重複は発生しますか?発生する場合、
event_idで冪等性の境界を定義する必要があります。 - ルール違反は自動修復可能ですか、それとも人間の承認が必要ですか?これにより、再送キューと制御方法が変わります。
- ダウンストリームのメトリクスは遅延や修正を許容できますか?許容できない場合、修復されたデータには補償パーティションとバージョン管理されたレポートが必要です。
- 未加工のペイロードに個人データ(PII)は含まれていますか?それにより暗号化、マスキング、アクセス、削除の要件が変わります。
30秒の回答フレームワーク
パイプラインを、不変の未加工層、パース、ルール検証、有効データ、隔離の各ステージに分割します。パースやルールの失敗は、event_id、スキーマとルールのバージョン、理由コード、未加工データへの参照を含む隔離レコードを生成します。有効なイベントは冪等な書き込みによってファクトテーブルに入力されます。隔離領域は、修復、承認、再送のキューを提供します。再送は同じevent_idを使用し、ターゲットの境界で重複排除を行います。有効率、理由コードの比率、隔離期間、再送の成功率を監視し、ビジネスSLOと重大度に基づいてアラートまたはブロッキングゲートを選択します。
ステップごとの詳細解説
1. まず未加工のエビデンスを保持する
不変のオブジェクトまたはログをバッチ、ソース、受信時刻ごとにパーティション分割し、チェックサムを保存します。処理ジョブは、未加工のペイロードを上書きするのではなく、ステータスを追記します。これにより、ルールのアップグレード、パーサーの修正、サプライヤーとの不一致の検証を、同一の入力から再現できるようになります。未加工層と隔離層の権限を分離することで、サポートユーザーがファクトを直接編集することを防ぎます。
2. 失敗を分類し、理由を保持する
最初にバイト/フォーマットのパースを実行し、次にスキーマ型と必須フィールドのチェック、3番目にクロスフィールドのビジネスルールを実行します。MALFORMED_JSON、MISSING_ORDER_ID、INVALID_CURRENCYなどの構造化されたreason_codeの値を保存します。1つのレコードに複数の理由が存在する場合がありますが、後からの修復で原因を説明できるように、最初に失敗したステージとルールのバージョンを保持します。
quarantine_record = {
event_id, source_batch, raw_uri, payload_hash,
schema_version, rule_version, failed_stage,
reason_codes, first_seen_at, status
}Sparkのファイルオプションは不良ファイルを記録したり破損ファイルを無視したりできますが、ジョブを継続させることは、ビジネスレコードを安全に保持することと同じではありません。単に無視スイッチを有効にするのではなく、回復可能なレコードを隔離領域に明示的にルーティングする設計にする必要があります。
3. 失敗(Fail)、破棄(Drop)、隔離(Quarantine)の選択
読み取り不能なインフラストラクチャ入力、信頼できない署名、またはバッチ全体を汚染する可能性のある破損が発生した場合は、バッチを失敗(Fail)させ、アラートを保持する必要があります。他のイベントに影響を与えることなく隔離できる単一レコードのエラーは、有効なデータの処理を継続できるように隔離(Quarantine)する必要があります。破棄(Drop)が許容されるのは、レコードが回復不能であり、ビジネスオーナーが損失を許容し、監査証跡が求められる場合のみです。すべての破棄はカウント可能かつ追跡可能でなければなりません。
4. 冪等な有効データ書き込みの定義
event_idとソースバージョンを一意のキーとして使用し、冪等なアップサート(upsert)またはコミットログを介してファクトテーブルに書き込みます。隔離データを再送する前に、ターゲットがすでにそのイベントを受け入れているかどうかを確認し、スキップ、更新、または補償バージョンのいずれかを選択します。注文金額などの修正可能なファクトについては、履歴を警告なしに上書きしないでください。修正イベントを発行し、ダウンストリームのレポートがバージョンまたは有効時刻によって再計算できるようにします。
5. 修復、承認、および再送
修復ツールは新しいペイロードまたはパッチを作成し、未加工層を変更することはありません。操作者、理由、入力ハッシュ、およびルールバージョンを記録し、承認フローに沿って変更をルーティングします。再送ワーカーは隔離状態を読み取り、完全な検証チェーンを再実行します。成功すると、アトミックにQUARANTINEDをREPLAYEDに移行します。失敗した場合は試行回数をインクリメントし、次回をスケジュールします。リースまたはデータベースロックにより、2つのワーカーが同じイベントを同時に適用することを防止します。
6. 品質メトリクスとリリースゲート
総取り込み量、有効率、各reason_codeの比率、隔離期間のパーセンタイル、再送成功率、重複イベント、およびダウンストリームの修正を監視します。重大度に基づいてゲートを設定します。署名の失敗はゼロトレランスにする一方、必須ではないフィールドの欠落はアラートのみにとどめることができます。0.2%であっても、それを正常と決めつけるのではなく、過去のベースライン、ソースの構成比、ビジネス上の損失と比較する必要があります。ゲートが作動した場合は、ダウンストリームへの公開を凍結するか、以前のルールバージョンにロールバックし、手動リリースの記録を残します。
7. 保持期間、プライバシー、およびリカバリ
修復に必要な最小限の未加工フィールドのみを保持し、機密ペイロードを暗号化して、アクセスを制限します。保持および削除のリクエストは、event_idから未加工オブジェクト、隔離行、および派生インデックスにマッピングされる必要があります。キュー、メタデータテーブル、およびオブジェクトストレージは個別にバックアップします。ターゲットへの書き込みが成功したもののステータス更新に失敗した場合は、一意キーによってリトライします。再送済みとマークされているが書き込みが不確実な場合は、ワーカーのレスポンスから成功を推測するのではなく、コミットログまたはターゲットテーブルのチェックからリカバリします。
質の高い模範解答
まず、許容される不具合バジェット、重複または遅延配信の有無、注文の修正がレポート処理を待てるかどうかを確認します。パイプラインは不変の未加工層を保持し、パース、スキーマ、ビジネスルールの順に段階的に実行します。レコードレベルの回復可能なエラーは隔離され、インフラやセキュリティの障害はバッチをブロックします。隔離行には、冪等性キー、未加工データへの参照、ルールバージョン、理由コード、および状態が保存されます。有効なイベントは冪等にファクトテーブルに書き込まれます。修復では新しいペイロードが作成されて承認が必要となり、再送では完全な検証チェーンが再実行され、成功がアトミックにマークされます。メトリクスは有効率、理由、隔離期間、再送率、重複の影響をカバーし、重大度ベースのゲートを備えます。これにより、異常を正常なデータとして隠蔽することなく、小規模な不正データによってバッチ全体がブロックされるのを防ぎます。
よくある間違い
ignoreCorruptFilesを有効にして成功とみなす → ジョブは継続しますがレコードが消失する可能性があります → 回復可能なレコードは理由コード付きの隔離領域にルーティングしてください。- 編集可能なテーブルに失敗を保存する → 未加工のエビデンスが改ざんされる可能性があります → 未加工データは不変に保ち、新しい修復バージョンを作成してください。
- 再送イベントを直接再挿入する → ダウンストリームのメトリクスが二重カウントされます →
event_idの一意性境界とコミットログを使用してください。 - すべてのエラーでバッチ全体をブロックする → わずかな不良データが鮮度を損ないます → ステージと重大度に応じて失敗(Fail)または隔離(Quarantine)を選択してください。
- 失敗の総数のみをカウントする → ソースやルールのデグレード(リグレッション)が隠れてしまいます → 理由、スキーマバージョン、ソース、および時間ごとにメトリクスを分解してください。
- 修復後に検証をスキップする → パッチによって2つ目の不具合が混入する可能性があります → 再送時は完全な検証チェーンを実行する必要があります。
フォローアップ質問と回答
隔離率が0.2%から8%に急上昇しました。公開を継続しますか?
まず理由とソースごとに分割して分析します。1つのサプライヤーで回復可能なフィールド欠落が発生している場合は、そのソースのみを一時停止し、他のソースは継続します。パースや署名検証が失敗している場合は、ダウンストリームへの公開を凍結し、ルールバージョンをロールバックします。しきい値は絶対的なパーセンテージだけでなく、ビジネス上の損失や過去のベースラインを反映する必要があります。
再送によって決済済みの注文が変更されるのをどのように防ぎますか?
元のイベントと修正イベントを分離し、ファクトテーブルにバージョンまたは有効時刻のカラムを保持します。決済スナップショットは、その入力の世代を固定します。修復されたイベントは補償ワークフローを通過し、決済済み金額を警告なしに上書きするのではなく、調整を作成するかどうかを財務ルールに基づいて決定します。
隔離領域にPIIが含まれています。サポート担当者はどのように調査できますか?
デフォルトではマスキングされたフィールドと理由コードを表示します。未加工のペイロードに対しては、有効期限の短い認可アクセスを付与し、すべての監査イベントを記録します。削除リクエストはevent_idに従って未加工オブジェクト、隔離行、および派生インデックス全体に適用され、処理後は不可逆的な監査ダイジェストのみを保持します。
ターゲットへの書き込み後、隔離状態を変更する前にワーカーがクラッシュしました。何が起こりますか?
リトライ時に、冪等性キーを使用してターゲットテーブルとコミットログを確認します。存在する場合はステータスのみを修復し、副作用(再書き込み)を繰り返さないようにします。存在しない場合は再送信します。不明な結果を成功または失敗と推測することがないよう、条件付きでリトライ可能な状態遷移を使用します。
レコードを保持せずに破棄(Drop)すべきなのはどのような場合ですか?
レコードが回復不能であり、コンプライアンスルールで保持が義務付けられておらず、ビジネスオーナーがその損失を明示的に受け入れている場合に限られます。その場合でも、監査可能なカウント、理由、バッチ、ポリシーバージョンを保持し、品質レポートで破棄されたことを開示する必要があります。