代表的な面接トピック

データエンジニアリング面接:大規模な履歴データの安全なバックフィルをどのように行いますか?

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

質問

注文データパイプラインは、オブジェクトストレージから毎日 2 TB の不変な生データを処理しています。変換の不具合が過去 90 の business_date パーティションに影響を与えたため、5 日以内に 180 TB を再計算する必要があります。日次のパイプラインは P95 の鮮度を 45 分以下に維持し続ける必要があり、遅延した修正がバックフィルと同じ注文を更新する可能性があります。再現可能、一時停止可能、再開可能、かつロールバック(可逆)可能な履歴バックフィルを設計してください。キャパシティバジェット、本番処理と履歴処理の分離、データ検証、および安全な公開について説明してください。

プロンプトと適用されるコンテキスト

注文データパイプラインは、オブジェクトストレージから毎日 2 TB の不変な生データを読み取り、business_date でパーティショニングされた orders_daily テーブルを生成します。チームは 90 個のパーティションに影響を与える税金変換の不具合を発見しました。そのため、5 日以内に 180 TB を再計算する必要があります。日次のインクリメンタルパイプラインを停止することはできず、その P95 鮮度は 45 分以下に維持されなければなりません。また、遅延した返金や注文の修正がバックフィルと同じビジネスキーを更新する可能性があります。

再現可能で、一時停止後に再開可能、監査可能、かつロールバック(可逆)可能なバックフィルを設計してください。入力およびコードのバージョン固定、パーティションの分割とスケジューリング、履歴処理による本番キャパシティの枯渇防止、バックフィルと本番インクリメントの重複解決、公開をブロックする検証の定義、および障害発生後に最後に信頼されたバージョンを復元する方法を網羅してください。

近年の公開データエンジニアリング面接の教材では、履歴バックフィル、冪等な再実行、およびリアルタイム処理の保護がパイプラインの信頼性に関する質問として明確に扱われています。公開されているシナリオでも、テラバイト規模の入力処理、パーティションの再処理、検証、ロールバックへの対応が候補者に求められます。公式のパイプラインドキュメントでは、履歴実行、再処理ポリシー、および個別の同時実行制限が公開されています。検索意図は具体的です。候補者には、「Airflowで過去90日分を再実行する」というスケジューリングのエントリポイントにとどまらない、実行可能な本番計画が求められています。

面接官が評価しているポイント

最初のシグナルは、候補者が再現可能なデータバージョンを定義しているかどうかです。優れた回答では、バックフィル対象期間、論理日付、ソースのスナップショットまたはソースバージョン、変換コード、依存するディメンションのバージョン、およびターゲットスキーマが固定されます。2 回の試行で異なる入力が読み取られたり、コードが now()、ランダム値、または変更可能な外部ルックアップに依存している場合、「安全に再実行可能」という言葉は検証可能な意味を持ちません。

2 番目のシグナルは、オーケストレーションの成功とデータの正当性を区別できているかです。オーケストレータは過去の論理日付の実行を作成し、同時実行数を制限できます。しかし、ビジネス書き込みを冪等にするわけではなく、本番インクリメントと履歴バックフィルの間の書き込み競合を解決するわけでもありません。候補者は、パーティションの置換、安定したキーに基づく MERGE、またはバージョニングされたターゲットを選択し、その前提条件を説明する必要があります。

3 番目のシグナルは、キャパシティと分離です。180 TB を 5 日で完了するための生データの最小平均読み取り速度は次のとおりです。

text
180 TB / (5 × 24 h) = 1.5 TB/h ≈ 417 MB/s

これは、リトライを含まない連続実行の下限値です。スキャンの増幅、シャッフル、ターゲットへの書き込み、検証、失敗した再計算は除外されています。優れた回答では、まず 1 つのパーティションをベンチマークし、ステージごとのスループットとピークリソースを測定した上で、個別のキューまたはコンピュートプール、同時実行制限、および本番優先度を設定します。本番パイプラインの鮮度は、バックフィル速度を下げるためのフィードバックシグナルとなります。

4 番目のシグナルは、公開とリカバリの境界です。90 個のタスクが成功したからといって、コンシューマーにとって新しいバージョンが安全であるとは限りません。公開前に、システムはパーティションの完全性、キーの一意性、ビジネス不変条件、ソースとの照合、および不具合と整合する差分を証明する必要があります。公開は 1 回の制御されたバージョン切り替え(cutover)であるべきです。観察期間中は古いバージョンを保持し、ロールバックが 180 TB の再計算ではなくポインタの変更だけで済むようにします。

回答前に明確にすべき質問

  • 生データの入力は本当に不変か? オブジェクトのバージョン、スナップショット ID、または再生可能なログ位置を取得します。ソースがデータをインプレースで変更する場合は、まず参照可能な入力バージョンを作成します。
  • どのクロックが 90 日間の期間を定義しているか? business_date、イベント時刻、取り込み時刻、およびタイムゾーンを一致させます。また、遅延した返金をどのパーティションが保持するかを定義します。
  • ターゲットの粒度と安定したキーは何か? 行が注文、注文項目、または日次集計のいずれであるか、また order_idsource_version、決定論的な競合順序が存在するかどうかを確定します。
  • 変換はどの変更可能なデータに依存しているか? 為替レート、税規則、SCD ディメンション、および削除レコードは、本日の値に暗黙的に置き換えるのではなく、過去のイベント発生時点のデータとして読み取る必要があります。
  • 本番のインクリメントはどのパーティションを更新できるか? 最新の 7 日間のみに書き込む場合、バックフィルはそれより古いパーティションを専有できます。過去の任意の注文が変更される可能性がある場合は、バージョンの順序付けまたは差分のキャッチアップが必要です。
  • どのような分離機能およびアトミック公開機能が存在するか? 独立したウェアハウス、リソースプール、優先度キュー、パーティショントランザクション、テーブルクローン、ビューの切り替え、またはカタログポインタによって設計が変わります。
  • 5 日間は厳密な期限か、それとも目標か? 本番の P95 鮮度目標、ソースの読み取りクォータ、コスト上限、および許可される短い切り替えウィンドウを取得します。
  • 誰が承認するか? 技術的なアサーション、財務照合、ダウンストリームのサンプリング、および観察期間のそれぞれに、担当者とブロックの閾値が必要です。

30秒の回答フレームワーク

「90 のビジネス日付、入力スナップショット、およびコードバージョンを固定し、各論理日付パーティションをマニフェスト付きの独立したステージングに書き込みます。180 TB を 5 日で処理するには最低でも約 417 MB/s が必要なため、まずベンチマークを実施し、本番の P95 が 45 分に近づいた場合はスロットリングを行います。W0 をバックフィルした後、W1 までの差分をキャッチアップし、90/90 のパーティション、照合、およびビジネスチェックに合格した後にのみ切り替えます。ロールバックに備えて古いバージョンはそのまま利用可能にしておきます。」

ステップバイステップの詳細解説

ステップ 1:バックフィルを不変の実行仕様として定義する

backfill_id を作成し、以下を記録します。

フィールド目的
範囲[2026-04-01, 2026-06-29]、90 パーティション実行中に境界がずれるのを防止する
入力raw_snapshot=s_1042、ウォーターマーク W0すべての試行で同じファクトを読み取る
ロジックcode_sha=abc123tax_rules=v17変換と依存関係のバージョンを固定する
出力orders_daily__bf_20260718候補結果を信頼されたバージョンから分離する
リソースバックフィルプール、最大同時実行数、読み取り/書き込みクォータ本番 SLO を保護する
ゲートユニークキー、金額差分、完全性、承認者'完了' を判定可能な状態にする

business_date を各パーティションの実行に渡します。変換内で論理時刻を実世界時刻(ウォールクロック)に置き換えてはいけません。過去の為替レートや SCD ディメンションが依存関係にある場合は、イベント時刻での as-of 結合を実行します。入力、コード、またはいずれかの依存関係のバージョンが変更されると新しい backfill_id が作成されます。1 回の実行内で 2 つの結果バージョンを混在させないでください。

本番ターゲットに書き込まずにプランを生成します。全 90 パーティション、依存順序、推定入力バイト数、提案する同時実行数、および出力先パスを一覧表示します。欠落、重複した日付、ソースの保持期間外のパーティション、およびダウンストリームへの副作用を検出します。電子メール、請求処理、外部 API 呼び出しなどのデータ以外の副作用を無効化するか、監査モードにリダイレクトして、履歴再生によって実際のビジネスアクションが再度トリガーされないようにします。

ステップ 2:一時停止、再開、監査のためにパーティションマニフェストを使用する

business_date を有界な作業単位として扱います。マニフェストには少なくとも以下を記録する必要があります。

text
backfill_id, business_date, input_snapshot, code_sha,
state, attempt, input_rows, output_rows, output_checksum,
staging_location, published_version, started_at, completed_at

有用な状態モデルは PENDING → RUNNING → VALIDATED → PUBLISHED であり、失敗した場合は FAILED に遷移します。条件付き更新またはリースを使用して作業単位を要求し、パーティションのアクティブな所有者が 1 つだけになるようにします。期限切れのリースは再要求可能です。失敗したパーティションのみを再試行し、そのパーティション用に別の分離されたステージング試行を書き込みます。最終テーブルにやみくもにアペンド(追記)してはいけません。

日付パーティションが完全にクローズされている場合、最も単純な冪等書き込みは、完全なパーティションを構築してからそのパーティションをトランザクション的に置換することです。注文が日付をまたいで修正される可能性がある場合は、安定したビジネスキーとソースバージョンを MERGE で使用します。タイブレークの明示的な勝者(例えば source_updated_at、同値の場合は単調増加する source_sequence など)を定義します。主キー、バージョン、および削除セマンティクスのすべてが信頼できる場合にのみ、MERGE は古い履歴レコードが新しい修正を上書きするのを防ぎます。

ステップ 3:スレッド数を推測するのではなく、測定から同時実行数を導出する

180 TB / 5 日の要件から、生データ読み取りの下限は約 417 MB/s となります。代表的な 1 パーティションでカナリア実行を行い、読み取り、解凍、シャッフル、変換、書き込み、検証に関するバイト数、所要時間、CPU、メモリ、ウェアハウスのキュー待機時間、および一時領域を測定します。1 つのパーティションが r MB/s の実効スループットを提供する場合、理論上の同時実行数の下限はおよそ ceil(417/r) となります。ソースのクォータ、シャッフルのピーク、ターゲットのコミットキャパシティ、およびコスト制限によって、実際の値はさらに制約されます。

バックフィルには、日次インクリメントよりも低い優先度で個別のコンピュートプールまたはキューを割り当てます。スケジューラの同時実行数、ソースの読み取り、ターゲットの書き込み、および総コストを組み合わせて制限します。どれか 1 つだけを制限しても通常は不十分です。コントローラは、本番パイプラインの P95 鮮度、ウェアハウスのキューイング、およびソースのスロットリングを監視します。鮮度が 45 分に近づくと、新しいパーティションの要求を停止するか、同時実行数を引き下げます。本番パイプラインが安全な範囲に戻った後にのみ徐々に引き上げます。コミット中のパーティションを強制終了(ハードキル)して不完全な出力を残してはいけません。キャンセルシグナルは新規の要求を防止し、アクティブな作業を完了させるか、ステージングを安全にロールバックさせる必要があります。

カナリア実行が成功した後は、段階的にスケールさせます。例えば、最初は 1 パーティション、次に 3、そして測定された安全な同時実行数へとスケールさせます。各レベルで本番の完全なインクリメンタルサイクルを 1 回観察します。キャパシティが期限に間に合わない場合は、早期に期限、一時的なキャパシティ、またはスコープを変更します。本番の鮮度を犠牲にして見積もりエラーを隠蔽してはいけません。

ステップ 4:履歴パイプラインと本番パイプラインの間で書き込みの所有権を定義する

まず、ソーススナップショットまたはログウォーターマーク W0 から履歴期間を新しいバージョンに再計算します。本番パイプラインが直近の 7 日間のみを修正する場合、バックフィルにそれ以前の 83 日間を排他的に所有させます。最新の 7 日間は本番の所有権下に維持し、最後にその重複ウィンドウを新しいバージョンに再計算またはマージします。

本番の修正が過去の任意の注文に影響を与える可能性がある場合は、スナップショット+差分キャッチアップ(snapshot-plus-delta-catch-up)フローを使用します。

  1. W0 を記録します。バックフィルは W0 以前の決定論的な入力のみを読み取ります。
  2. 本番パイプラインは古いバージョンへのデータ提供を継続し、W0 以降の変更は再生可能なログに保持されます。
  3. 90 個の履歴パーティションすべての検証が完了した後、W1 を記録し、同じバージョンルールを使用して (W0, W1] の変更を新しいバージョンに適用します。
  4. 遅延が切り替えバジェット内に収まったら、公開ポインタを一時的に凍結するか書き込みフェンスを取得し、最終差分を適用します。
  5. コンシューマーをアトミックに新しいバージョンへ移行し、その後それを本番パイプラインの書き込みターゲットとします。

この設計には、完全な変更ログ、安定した順序付け、および真にアトミックな切り替えプリミティブが必要です。プラットフォームにアトミックなビューまたはカタログ変更機能がない場合は、明示的なメンテナンスウィンドウ、または古いパーティションのバックアップを伴うトランザクションパーティション置換を使用します。目に見える中間状態を文書化してください。シームレスな切り替えを口頭だけで約束してはいけません。

ステップ 5:階層的に検証し、公開前に厳格なゲートを設ける

各ステージングパーティションが完了した後、3 つのレイヤーでチェックを実行します。

  • 構造と完全性: 互換性のあるスキーマ、必須フィールドの存在、一意なビジネスキー、正しい日付の境界、および 90 パーティション内での欠落や重複がないこと。
  • ソースとの照合: 日付、地域、通貨、およびステータス別に入力と出力の行数、ユニークな注文数、税抜金額、税金、および純額を比較します。集計については、調査可能なレコードレベルの差分を保持します。
  • ビジネスルールと差分: 金額の保存則、返金可能額の範囲内の返金、および有効な状態遷移。新旧バージョンの差分は、不具合の影響を受けた注文に集中している必要があります。影響を受けていないスライスは、説明なしに変更されてはなりません。

行数の一致は弱い証拠です。不正な結合によって行の追加と欠落が同時に発生する可能性があるためです。安定したキーによるバケット化されたチェックサムも計算し、既知の不具合、境界日、遅延修正、削除をサンプリングし、財務部門またはデータプロダクトオーナーに修正の方向性を確認してもらいます。検証クエリをバージョニングし、閾値、実際の値、および結果を保持します。

全体的な公開ゲートには、少なくとも以下が必要です。マニフェスト内の 90/90 VALIDATED、アクティブまたは失敗したパーティションが存在しないこと、入力スナップショット・コード・依存関係のバージョンの一貫性、W1 までのキャッチアップ、すべての厳格なアサーションの合格、バックフィル中の日次パイプラインの 45 分 P95 鮮度の維持、およびオーナーの承認。何らかの障害が発生した場合は、古いバージョンが表示されたままになります。

ステップ 6:1 回で切り替え、継続的に観察し、迅速にロールバックする

公開前に古い読み取りポインタと出力バージョンを保存します。切り替えは安定したビューまたはカタログポインタのみを変更し、切り替え中に 180 TB を移動させることはありません。重要なダッシュボードやダウンストリームジョブに対してコンシューマークエリスイートを即座に実行し、クエリレイテンシと最新のインクリメンタル書き込みを確認します。観察期間中は、新旧バージョン間で重要なメトリクスの比較を継続します。

キーの一意性、金額の照合、本番の鮮度、またはコンシューマークエリが失敗した場合は、新しいバージョンへの公開を停止し、ポインタを元に戻します。失敗したバージョン、マニフェスト、および検証証拠は調査のために保持します。インプレースで編集して同じバージョンとして提示しないでください。入力またはコードを修正した後、新しい実行仕様を作成します。信頼性が証明されバージョン互換性のあるパーティションは再利用し、バージョンが異なるパーティションのみを再計算します。

監視ビューには、処理済みおよび残りの TB、状態別のパーティション数、スループット、完了推定時間、リトライおよび検証失敗、ソースのスロットリング、コンピュートプールのキューイング、ターゲットのコミットレイテンシ、および本番パイプラインの P50/P95 鮮度を表示する必要があります。アラートには、オペレーターが同時実行数の削減、パーティションの再試行、または公開全体の中止を判断できるように、backfill_id、パーティション、コードバージョン、失敗したゲート、および担当者を含める必要があります。

高品質な回答例

「まず、このバックフィルをバージョニングされた実行として凍結します。90 個の business_date 値、生のスナップショットとウォーターマーク W0、コード SHA、税ディメンションのバージョン、ターゲットテーブルのバージョン、および承認閾値です。各日は分離されたステージング出力を書き込む 1 つのタスクとなります。マニフェストには PENDING/RUNNING/VALIDATED/PUBLISHED とすべての試行が記録されます。変換には論理時刻のみを使用します。クローズされたパーティションは完全に構築されて置換されます。遅延修正を受け取る可能性のある注文は、古いバージョンが新しいバージョンを上書きできないよう、安定した主キーとソースバージョンを MERGE 内で使用します。」

「180 TB を 5 日で処理するには、リトライや検証の前に、生データの平均読み取りで最低約 417 MB/s が必要です。代表的な 1 パーティションをベンチマークし、エンドツーエンドのスループットと各ステージのピークを測定してから、段階的に同時実行数を増やします。バックフィルには、ソースの読み取り、ターゲットの書き込み、およびスケジューラの同時実行数に制限を設けた、分離された低優先度のリソースプールを割り当てます。コントローラは本番の P95 鮮度を保護し、遅延が 45 分に近づくと新しいパーティションの要求を停止し、回復後にゆっくりと再開します。」

「並行処理の正当性を保つため、本番パイプラインが古いバージョンの提供を継続している間に、スナップショット W0 から新しいバージョンを構築します。履歴パーティションの完了後、W1 を記録し、(W0, W1] の修正を新しいバージョンに再生し、短いフェンスを使用して最終キャッチアップを行った後、ポインタをアトミックに移動します。公開には、90/90 のパーティション、一意なキー、ソース金額の照合、ビジネス不変条件、説明可能な新旧の差分、および健全な本番の鮮度が必要です。観察期間中は古いバージョンを保持します。厳格なゲートまたはコンシューマークエリが失敗した場合はポインタを元に戻し、失敗したバージョンを調査用に保存します。」

よくある間違い

  • スケジューラの同時実行数を最大に設定する → ソース、シャッフル、ターゲットのコミット、または日次ジョブが実際のボトルネックになる → 1 パーティションのベンチマーク、スループットの下限、および本番 SLO から同時実行数を導出し、動的にスロットリングする。
  • タスクのリトライを冪等性と同一視する → オーケストレータはタスクを再実行するだけであり、無秩序なアペンドは依然としてデータを重複させる → 決定論的な入力と論理時刻を使用し、パーティションを置換するか、安定したキーとバージョンでマージする。
  • 履歴ジョブと本番ジョブに同じパーティションを書き込ませる → 完了の遅れた古い計算が新しい修正を上書きする可能性がある → パーティションの所有権を割り当てるか、W0/W1 差分キャッチアップと明示的なバージョンの順序付けを使用する。
  • 90 個のタスクが成功した時点で公開する → 成功ステータスは、完全なレコード、正しい金額、または妥当な差分を証明しない → パーティション、ソース、ビジネス、差分、およびコンシューマーのチェックでゲートを設ける。
  • 検証前にオンラインテーブルを上書きする → 検証が失敗した場合、リカバリ手段として再び大規模な再計算が必要になる → バージョニングされた出力を書き込み、それを検証し、1 回で切り替えを行い、古いポインタを保持する。
  • 総行数のみを比較する → 重複と欠落が互いに相殺される可能性がある → 一意キー、バケット化されたチェックサム、金額、状態の分布、およびレコード差分も比較する。
  • 履歴変換内で now() を呼び出す → 2 回の試行で異なるパーティションセマンティクスが生成される → 論理日付、入力スナップショット、および依存関係のバージョンを実行パラメータとして渡す。
  • 障害発生後に 180 TB すべてを再実行する → コストとリスクが増大し、検証済みの進捗が破棄される → パーティションマニフェストから失敗した作業単位を再開し、バージョンが変更された場合にのみ影響を受けるパーティションを再計算する。

フォローアップの質問と回答

フォローアップ 1:ソースに不変のスナップショットがない場合、バックフィルをどのように再現可能にしますか?

実行前に作成されたバージョニングされた履歴エクスポートを優先するか、データベースのスナップショット位置、CDC ログ位置、およびオブジェクトバージョン ID を記録します。唯一のソースが変更可能なテーブルである場合は、読み取り時刻、ウォーターマーク、およびパーティションごとのチェックサムをマニフェストに記録し、後でキャッチアップできるように実行中の変更を継続的にキャプチャします。元の入力を再構築できない場合は、再現性の制限を開示し、2 回の試行が必ず同一になると主張してはいけません。

フォローアップ 2:修正された SCD ディメンションがファクトテーブルに供給されています。どちらを先にバックフィルすべきですか?

影響を受けるリネージのサブグラフを構築し、トポロジカル順序で処理します。まず新しいディメンションバージョンを作成し、次にイベント発生時点のデータとしてファクトテーブルをそのバージョンに結合し、最後に集計とマートを再構築します。すべてのレイヤーで 1 つのバックフィルバージョン識別子を共有します。古いディメンションを新しいファクトテーブルと一緒に公開してはいけません。ディメンションのビジネスキー、重複のない有効期間、ファクトの外部キー一致率、および重要な集計を検証します。

フォローアップ 3:プラットフォームにアトミックなテーブルスワップ機能がない場合、どのように公開しますか?

独立したバージョンテーブルを書き込み、安定したビューを介してコンシューマーをルーティングすることを推奨します。ビュー定義をアトミックに置換できる場合は、ビューのみを切り替えます。それすら利用できない場合は、古いパーティションのバックアップを伴うトランザクション形式の小バッチ置換、または切り替えと受け入れのために一時的に関連する読み書きを停止する明示的なメンテナンスウィンドウを使用します。目に見える中間状態とロールバック時間を明示してください。複数ステップの上書きをアトミックと表現してはいけません。

フォローアップ 4:新旧のバージョンで行数と合計金額が一致しています。他に何をチェックすべきですか?

安定したビジネスキーによるバケット化されたチェックサムを計算し、レコードの差分を検査します。地域、通貨、注文ステータス、税区分、および境界日ごとの分布を比較します。既知の不良サンプルが変更され、影響を受けていないサンプルが変更されていないことを確認します。主キー、参照整合性、状態遷移、および返金上限を確認します。合計が一致していても、過剰カウントと欠落が相殺し合っている可能性があります。

フォローアップ 5:実行は最も古い日付から開始すべきですか、それとも最新の日付から開始すべきですか?

依存関係とビジネス価値に依存します。以降のパーティションが以前のパーティションの状態に依存している場合は、昇順(古い順)で実行します。パーティションが独立しており、直近のレポートの緊急度が高い場合は、逆順(新しい順)が役立ちます。どちらの場合も、プロダクトオーナーがフェーズごとに個別のバージョン、コンシューマースコープ、およびロールバック境界を設定した段階的公開を明示的に承認しない限り、公開ゲートは期間全体を対象とします。

公開情報ソース

関連する質問