プロンプトとコンテキスト
これはデータプラットフォームのシステム設計に関する質問です。リバースETLは、信頼できるDWHモデルをCRM、マーケティング、またはプロダクトツールに配信します。Hightouchのドキュメントではこのフローを「ソース → モデル → 同期 → 送信先」と説明しており、CensusのガイドではDWHからビジネスプラットフォームへの配信をオペレーショナルアナリティクスと位置付けています。この面接では、バッチ処理と増分処理、送信先APIの制限、データガバナンスが総合的に試されます。
面接官が評価するポイント
- モデルのスナップショット、変更検出、スケジューリング、キュー、送信先アダプターを適切に分離できているか。
- ソースの配信が少なくとも1回である場合に、冪等性キーとバージョンを使用して送信先を保護できるか。
- 削除、同意の撤回、スキーマドリフト、テナントの分離を明示的なコントラクトに落とし込めるか。
- 単にデータフロー図を描くだけでなく、鮮度、成功率、バックログ、照合メトリクスによって信頼性を証明できるか。
明確化のための質問
- 10分のSLOは高優先度コホートのみに適用されますか、それともすべてのレコードに適用されますか?
- 送信先はバッチアップサート、削除、冪等性キー、サーバーサイドカーソルをサポートしていますか?
- モデルは安定したキー、更新時刻、削除墓石(tombstone)を公開していますか?スナップショットはどのくらいの期間保持されますか?
- テナントのクォータは独立していますか?また、1つの大規模テナントがグローバルのスループット全体を消費する可能性はありますか?
- 同意の撤回はどれくらい迅速に有効になる必要があり、削除処理が失敗している間、新規同期はブロックされる必要がありますか?
「私ならシステムを、バージョン管理されたモデル、変更検出器、テナント別キュー、送信先アダプター、照合ジョブに分割します。各レコードには tenantid、安定したビジネスキー、モデルバージョン、rowversion、削除状態を持たせます。高水位線(ハイウォーターマーク)またはCDCにより、少なくとも1回実行されるタスクを作成します。アダプターは送信先の制限内でアップサートをバッチ化し、テナント、送信先、recordid、rowversion を冪等性キーとして使用します。リトライによってバージョンが下がることはありません。削除と同意の撤回は回避不可能なフェンスを書き込みます。鮮度ラグ、バックログ、スロットリング、障害クラス、照合の差異、送信先での削除レイテンシを監視します。」
ステップごとの詳細解説
ステップ 1: ソースモデルとバージョンの定義
DWHモデルを同期の入力として扱い、ワーカーがアドホックに複数の運用データベースを結合しないようにします。安定した record_id、tenant_id、ビジネスフィールド、row_version、updated_at、consent_state、deleted_at を出力します。モデルの実行ごとに model_run_id が付与されます。レコードが削除されたときは、単に省略するのではなく墓石(tombstone)を出力し、ワーカーが『まだスキャンされていない』と『ダウンストリームで明示的に削除する』を区別できるようにします。
ステップ 2: 変更の検出と作業のスケジューリング
モデルの更新カラムまたはCDCの高水位線を優先して使用します。チェックポイントを永続化し、同一のタイムスタンプによって行の取りこぼしが発生しないよう重複ウィンドウを使用します。イミュータブルな変更バッチを書き込み、スケジューラーにテナントの優先度ごとに分割させます。10分のSLOに対して許容される遅延をキューの経過時間から計算します。低優先度の作業は必要に応じてキャパシティを譲りますが、同意撤回キューを迂回することはできません。
ステップ 3: アップサートの冪等化
少なくとも1回の配信では、『送信成功後にワーカーがクラッシュする』ケースを安全に再実行できる必要があります。送信先が冪等性をサポートしている場合は tenant_id + destination + record_id + row_version を使用し、現在のバージョン以上のバージョンのみを受け入れます。送信先の冪等性がない場合は、リクエストのフィンガープリントとレスポンスを保持し、同時実行数を制限し、送信先を読み取ることで照合します。システム間トランザクションが正確に1回(exactly-once)を提供すると主張してはなりません。リトライは、リトライ可能なエラー、指数バックオフ、最大試行回数によって分類します。
ステップ 4: スロットリングと過負荷の分離
テナントごとにトークンバケットまたは送信先から報告されるクォータを維持し、グローバルな同時実行上限も設けます。テナントキュー、フェアスケジューリング、デッドレターキューにより、1つの大規模テナントが他のテナントを枯渇させるのを防ぎます。429、5xx、ネットワークタイムアウトは遅延を伴ってリトライし、スキーマや認可に関する4xxエラーは手動処理にルーティングします。バックログがSLOに近づいたときにアラートを発報し、低優先度のバックフィルを一時停止できるようにします。
ステップ 5: 削除と同意撤回の処理
各撤回をテナント、record_id、イベントバージョンを含む独立した削除フェンスに書き込みます。ワーカーはアップサートの直前にフェンスをチェックします。撤回されたレコードは、ガバナンスによって明示的にクリアされるまで削除のみを送信できます。送信先の削除レシートとタイムスタンプを保持します。照合処理では、ダウンストリームに依然として存在する禁止レコードを検索する必要があります。リクエストの成功だけでは不十分です。
ステップ 6: スキーマドリフトとロールバックの管理
モデルスキーマをバージョン管理し、デプロイ前に入出力フィールドのマッピングを検証します。加算的なオプショナルフィールドはグレーリリースし、互換性のない型の変更や削除はすべてのテナントを破壊するのではなくレポートを出してブロックします。各アダプタータスクで mapping_version を保持し、失敗したバッチは古いマッピングでリトライし、最新の成功で部分的な移行を上書きするのではなく、検証済みのバージョンにロールバックします。
ステップ 7: 監視と照合
source_run、タスク状態、試行回数、最終成功バージョン、APIレイテンシ、スロットル、キューの経過時間、削除レイテンシをテナントおよび送信先ごとに記録します。コアメトリクスは、高優先度の鮮度ラグp95、成功率、デッドレター、スキーマエラー率、ソースと送信先のカウント差分、サンプリングされたフィールドハッシュ差分です。毎日またはリリース後に完全な照合を実行し、安全な再実行を自動修復し、照合不能なギャップは運用チームにルーティングします。
トレードオフと境界条件
スナップショット、増分、またはCDC
スナップショットはシンプルですがデータを再スキャンします。タイムスタンプの増分は低コストですが、安定したクロックと更新カラムに依存します。CDCは削除を表現できますが、ソースまたはモデリング層で変更ファクトを保持する必要があります。選択はモデルの更新頻度、削除セマンティクス、送信先のキャパシティに依存することを説明し、行の取りこぼしに対する保護策として定期的な完全照合を維持します。
キューの配置と整合性
テナントパーティションは分離性と順序性を向上させます。グローバルキューはキャパシティを効率的に使用しますが、フェアスケジューリングが必要です。単一レコードのバージョンに対する単調増加の可視性は保証できますが、DWHと送信先をまたぐアトミックコミットは保証できません。バージョン条件付き書き込み、リプレイ、照合によって、説明可能な結果整合性(Eventual Consistency)を提供します。
バックフィルとライブ更新
バックフィルには独立した低優先度のバジェット、一時停止可能なカーソル、スロットリングへの配慮を与え、ライブ更新は高優先度キューに送信します。双方が1つのレコードで競合する場合、より高い row_version が勝ち、送信先の条件付き書き込みは古いバージョンを拒否する必要があります。
模範解答
「私ならDWHモデルをバージョン管理されたソースとして扱い、高水位線またはCDCを使用してイミュータブルな変更バッチを作成し、テナントおよび送信先ごとにキューイングします。レコードには安定したキー、row_version、モデルバージョン、マッピングバージョンを持たせます。アップサートには送信先の冪等性またはリクエストのフィンガープリントを使用し、指数バックオフによるリトライで新しいバージョンを上書きできないようにします。テナントごとのトークンバケットとグローバルな同時実行上限でスロットリングを処理します。削除と同意の撤回は送信直前にチェックされるフェンスを書き込み、削除レシートを保持します。鮮度ラグ、キューの経過時間、デッドレター、スキーマエラー、フィールドハッシュ照合、禁止レコードの残存を検証します。コントラクトは結果整合性を伴う少なくとも1回の配信であり、システム間をまたぐ正確に1回ではありません。」
よくある間違い
- リバースETLをモデルの更新やフィールドマッピングを無視したリアルタイムデータベースレプリケーションと呼ぶこと。
- 送信先への重複リクエストを処理せずに『キューが正確に1回を保証する』と主張すること。
- テナント分離を行わず単一のグローバルレート制限を使用し、大規模テナントがすべてのキャパシティを消費できるようにしてしまうこと。
- 墓石、同意フェンス、ダウンストリームの照合なしに、現在のモデルからの消失を単なる削除として扱うこと。
- 各バッチがどのマッピングバージョンを使用したかを記録せずにスキーマ変更を配信すること。
フォローアップの質問
10分の鮮度SLOをどのように証明しますか?
モデルのコミット時刻または変更イベントの時刻から送信先が読み取り可能であることを確認するまでの時間を測定し、高優先度レコードのp95とタイムアウト率をレポートします。ワーカーの開始時刻や平均値だけでは不十分です。
送信先がコレクション全体の置換しかサポートしていない場合はどうしますか?
model_run_id を使用してテナントごとにバージョン管理されたスナップショットを構築し、一時コレクションをアップロードし、カウントとハッシュを検証して、アトミックにバージョンを切り替えます。削除と撤回には依然として個別のフェンスが必要です。次回のフルロードによって禁止データが安全に削除されると見なすことはできません。
送信先では成功したがレシートが失われた場合はどうしますか?
同じ冪等リクエストを再実行するか、リクエストのフィンガープリントと送信先の読み取りを使用して照合します。冪等性キーも読み取りも存在しない場合は、証拠なしに成功とマークするのではなく、不確定な結果として手動確認に回します。
これをいつ専用の同期プラットフォームに移行すべきですか?
送信先数、テナントクォータ、マッピングバージョン、ガバナンスフェンス、照合が1つのDAGの保守性を超えたときに、永続的なタスクサービスとアダプター層に分割します。小規模な単一送信先の構成はオーケストレーターと冪等スクリプトから開始できますが、それでも削除とリトライのコントラクトは必要です。