代表的な面接トピック

データエンジニアリング面接:Change Data Capture(CDC)パイプラインの設計

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

質問

ある企業には、20個のPostgreSQL OLTPデータベース、200個のテーブル、および8 TBの既存データがあります。平均で毎秒20,000行、ピーク時には毎秒60,000行のコミット済み変更が発生し、エンコードされた変更の平均サイズは約1 KBです。ウェアハウスまたはレイクハウスへのCDCパイプラインを設計してください。要件は、ソースのトランザクションコミットからクエリ可能なキュレーション済みデータまでのp99レイテンシが2分未満、アプリケーションの書き込みを停止しない初期同期が72時間以内、ソース書き込みp99の性能低下(リグレッション)が5%未満であることです。INSERT、UPDATE、DELETE、行ごとの順序保証、at-least-once配信、スキーマ進化、リプレイ、およびリコンシリエーションをサポートしてください。

問題とスコープ

ある企業には、20個のPostgreSQL OLTPデータベース、同期対象の200個のテーブル、および8 TBの既存データがあります。ソースでは平均で毎秒20,000行、ピーク時には毎秒60,000行のコミット済み行変更が発生します。エンコードされた変更は平均約1 KBです。データは分析用ウェアハウスまたはレイクハウスに到達する必要があり、ソーストランザクションのコミットからクエリ可能なキュレーション済みデータまでのp99レイテンシは2分未満でなければなりません。初期同期はアプリケーションの書き込みを停止することなく72時間以内に完了する必要があり、スナップショットとCDCワークロードによるソース書き込みp99の性能低下は5%未満に抑える必要があります。

パイプラインは、INSERT、UPDATE、DELETE、および主キー値の変更を処理できなければなりません。単一ソース内での行ごとの変更順序を保持し、at-least-once(少なくとも1回)配信を使用し、テーブルまたは主キー範囲ごとのスキーマ進化とリプレイをサポートし、ソースおよびシンクの障害から回復し、正確性を証明できるリコンシリエーション(照合機能)を提供する必要があります。高耐久な変更ストリームはRAWイベントを7日間保持し、より低コストなオブジェクトストアがそれ以上の履歴を長期間保持します。データベース数、スループット、サイズ、レイテンシ、およびリソース制限は面接用の前提条件であり、製品の性能保証ではありません。

設計にはブローカーや複数のコンポーネントが含まれますが、核心となるスキルはデータエンジニアリングです。すなわち、保存されたスナップショットを無限の変更ストリームにギャップなく結合し、回復可能な位置とべき等性の境界を定義し、宛先にデータ欠落や気付かれないデータ不正(サイレントエラー)がないことを証明することです。

面接官が評価するポイント

第一の評価ポイントは、コンポーネントを描く前に正確性の定義を行えるかどうかです。質の低い回答は「DebeziumとKafkaを使う」から始まります。質の高い回答は不変条件を述べます。すなわち、コミットされた変更は消失してはならない、1つの行はソースの順序で収束しなければならない、リプレイされたイベントが最終状態を狂わせてはならない、DELETEはターゲットに伝播しなければならない、スナップショットからストリームへの引き継ぎ境界にギャップがあってはならない、回復位置が不確実な場合はサイレントに継続するのではなくフェイルクローズ(安全に停止)しなければならない、ということです。

第二の評価ポイントは、初期スナップショットが「すべてのテーブルをエクスポートしてからCDCを有効化する」という単純なものではないと理解していることです。8 TBのエクスポート中も書き込みは継続します。エクスポート完了後にのみWALキャプチャを開始した場合、その間の変更はすでに失われている可能性があります。信頼性の高い設計では、一貫性のあるスナップショットを取得する前または取得と同時に、回復可能なログ位置と保持メカニズムを確立し、一致する位置からコミット済み変更をストリーミングします。コネクタがこの手順をカプセル化している場合でも、候補者はハンドオフにギャップが生じない理由を説明できなければなりません。

第三の評価ポイントは、シンクに至るまでat-least-once配信を考慮し抜くことです。コネクタは、イベントを高耐久に発行した後、ソースオフセットを記録する前にクラッシュする可能性があります。ウェアハウスのバッチはコミットされたものの、その確認応答(ACK)が失われる可能性があります。「ブローカー内でexactly-once」であっても、外部のウェアハウスまで自動的に網羅されるわけではありません。質の高い回答では、各イベントにソース識別子、ソース位置、トランザクション順序を付与し、すでにその行に記録されているバージョンよりも新しいバージョンのイベントのみを適用するようにします。

第四の評価ポイントは、PostgreSQLのレプリケーションスロットが持つ両刃のリスクを認識していることです。スロットを使用するとコネクタはLSNから再開できますが、コネクタの処理が遅延している間はWALが保持され続けます。保持バイト数やディスク空き容量のアラートがないと、シンクの障害が最終的にソースストレージを枯渇させる可能性があります。スロットが削除されて再作成された場合、新しいスロットは古い位置を再現できません。パイプラインは、「現在から」継続して何も失われていないと主張するのではなく、停止、リコンシリエーション、または再スナップショットを行う必要があります。

最後に、実際のデータセマンティクスを網羅している必要があります。業務上の削除(DELETE)イベントは、ログ圧縮(log compaction)に使用されるトゥームストーンとは異なります。主キー値の更新は通常、古いキーの削除と新しいキーの作成として現れます。別々のデータベースからのLSNを比較することはできません。コンシューマが複数行トランザクション内のすべての行についてアトミックな可視性を必要とする場合、設計はトランザクション境界メタデータを使用し、追加のバッファリング、レイテンシ、およびリカバリの複雑さを受け入れる必要があります。

回答前に確認すべき質問

  • 宛先は最新の状態、完全な変更履歴、またはその両方を必要としていますか? 最新の状態であれば主キーのMERGE操作に適しています。監査やリプレイには、イミュータブルなRAW変更レイヤーも必要です。最終状態だけでは履歴を再構築できません。
  • 2分のレイテンシクロックはどこから始まってどこで終わりますか? この問題では、ソースのコミットからクエリ可能なキュレーション済みデータまでを測定します。RAWイベントがブローカーに入る時点で終了する要件であれば、はるかに容易です。
  • マルチテーブルトランザクションはアトミックに可視化される必要がありますか? ここでのデフォルトは、ウェアハウスでの一時的な部分的一貫性を許容する行単位の順序付けです。トランザクション全体のアトミック性を必要とする金融系コンシューマの場合は、BEGIN/ENDのアセンブリが必要です。
  • すべてのテーブルに安定した主キーがありますか? 主キーがない場合、UPDATEおよびDELETEイベントの特定と重複排除が難しくなります。ビジネスキーを追加するか、行全体の完全一致、サロゲートキー、およびストレージコストの増加を明示的に許容する必要があります。
  • ソースのフェイルオーバー時にロジカルレプリケーションスロットは保持されますか? 保持されない場合、リカバリには制御された書き込み一時停止、最終LSNの検証、スロットの再作成、およびリコンシリエーションが必要です。プライマリのフェイルオーバーを通常の再接続として扱うことはできません。
  • どのスキーマ変更を自動的に通過させることができますか? この設計では、NULL許容カラムの追加など互換性のある変更を自動的に受け入れます。削除、リネーム、型の縮小変更、および主キーの変更は、制御されたマイグレーションを使用します。
  • 72時間のスナップショット期限と5%のソース負荷バジェットのどちらが厳しいですか? 72時間で8 TBを読み取るには、約30.9 MB/sの合計実効スループットが必要です。負荷テストでこれがソースバジェットを破綻させることが判明した場合は、期限を延長するか、セマンティクス的に有効なレプリカを使用するか、段階的にスコープを縮小します。
  • オンラインストリームは宛先の停止をどのくらいの期間吸収する必要がありますか? この設計では、ピーク時の30分間の停止をオンラインでカバーします。7日間の高耐久ストリームとオブジェクトアーカイブが、より長期のリプレイを処理します。

30秒の回答

「私はまず正確性を定義します。スナップショットとストリームにギャップがなく、行はソース位置によって収束し、リプレイは安全であり、削除は伝播し、レプリケーションスロットの履歴喪失時はフェイルクローズすることです。各PostgreSQLソースはWALに対してロジカルデコーディングを使用し、一貫性のあるスナップショットを取得し、対応するLSNから再開します。イベントはソース、テーブル、および主キーごとにパーティション化されて高耐久ストリームに入ります。RAWレイヤーはイミュータブルのまま維持され、キュレーションされたシンクは主キーとソースバージョンによってべき等なMERGE操作を実行します。スキーマの互換性は適用前に検証され、削除とキー変更は明示的に処理されます。セグメント化された鮮度、LSNラグ、保持WAL、スナップショットの進捗、およびリコンシリエーションの差分を監視し、コネクタのクラッシュ、スロット喪失、スキーマ変更、シンク停止の障害注入テストを行います。」

ステップごとの詳細解説

ステップ1: バッファ、スナップショット、およびリカバリバジェットのサイジング

平均負荷は次のとおりです。

text
20,000 events/s × 1 KB ≈ 20 MB/s
20,000 × 86,400 = 1,728,000,000 events/day
20 MB/s × 86,400 ≈ 1.728 TB/day

ピーク時のイングレスは約60 MB/sです。ピーク時に宛先が30分間停止した場合、RAWバックログは約次のようになります。

text
60 MB/s × 1,800 seconds = 108 GB

平均レートでの7日間は、レプリケーション、インデックス、およびエンコーディングのオーバーヘッドを除いて約12.096 TBのRAWデータになります。ブローカー、オブジェクトストア、およびネットワークは、20 MB/sの平均だけでなく、ピークトラフィックおよび追いつき(catch-up)処理に合わせてサイジングします。復旧したコンシューマがライブ入力レートと同等の処理しかできない場合、108 GBのバックログを解消することはできません。設計には余剰の消費キャパシティ、またはキュレーションレイヤーのレイテンシの一時的な緩和が必要です。

8 TBのスナップショットを72時間で完了するには、少なくとも以下が必要です。

text
8 TB ÷ 72 hours ≈ 30.9 MB/s

これは、スキャン増幅、シリアライゼーション、ネットワーク再試行、およびターゲットへの書き込みを除いた下限値です。本番に近い形状のデータでチャンク化スナップショットのベンチマークを実施し、ソースI/O、キャッシュの挙動、レプリケーションラグ、およびp99に基づいてレート制限を設けます。期限が5%のソースバジェットと競合する場合は、OLTPを過負荷にするのではなく、期限またはソースパスを変更します。

ステップ2: すべてのコンポーネントが正確性またはキャパシティの要件を満たすようにする

メインのデータフローは次のとおりです。

text
PostgreSQL WAL / logical slot
  → source connector
  → durable change stream keyed by source + table + primary key
  → immutable raw archive
  → schema validation and light normalization
  → sink staging tables
  → idempotent MERGE / DELETE into curated tables
  → warehouse and lakehouse consumers

すべてのソースデータベースに独自のソース識別子、コネクタ、およびレプリケーションスロットを割り当て、1つの障害ドメインがすべてのソースを停止させないようにします。高耐久ストリームはバーストを吸収し、コンシューマを疎結合にし、短期的なリプレイをサポートします。オブジェクトストレージは長期履歴を保持します。CDC処理レイヤーはスキーマ検証、正規化、およびルーティングのみに限定します。メインパスにリプレイ不可能な外部参照を設けるとリカバリが困難になります。2分のSLAと効率的なカラムナー書き込みのバランスを取るため、ステージングテーブルと小さなアトミックMERGEバッチを介してウェアハウスに書き込みます。

1つの行に対する変更が1つの順序付けられたパーティションに到達するように、パーティションキーとして(source_id, table_id, primary_key)の安定したエンコーディングを使用します。これによりグローバルな順序は提供されませんが、グローバルな順序は不要です。20個のソースデータベースは独立したLSNシーケンスを持っており、その数値にはソース間で横断的な意味はありません。

ステップ3: 書き込みオンラインでの初期スナップショットの実行

論理シーケンスは次のとおりです。

  1. コネクタが消費する前に必要なWALが回収されないよう、ロジカルレプリケーションスロットを確立する。
  2. 一貫性のあるスナップショットと、それに対応するソースログ位置を取得する。
  3. テーブルおよび主キーのチャンク単位でスナップショットを読み取り、現在の行をREADイベントとして発行する。
  4. スナップショットに関連付けられた位置から、コミット済みのINSERT、UPDATE、およびDELETEレコードをストリーミングする。
  5. シンクをソースバージョンによって収束させ、リカバリ時の境界リプレイを許容しつつ、欠落区間は絶対に発生させないようにする。

コネクタは、完全な一貫性スナップショットまたはインクリメンタルスナップショットウィンドウを実装する場合があります。「LSNを記録してから通常のSELECT文を実行する」という独自プロトコルを考案してはいけません。分離レベル、長時間トランザクション、および同時書き込みにより、欺瞞的な複雑さが生じるためです。検証済みのコネクタセマンティクスに依存し、スナップショットチャンクの実行中に対象の主キーが繰り返し更新、削除、再作成されるテストを実施してください。

主キー範囲でチャンク化し、進捗を永続化します。チャンクを小さくすると、長時間トランザクション、キャッシュの乱れ、障害時の手戻りが軽減されますが、極端に小さくするとクエリとスケジューリングのオーバーヘッドが増加します。各ソースのレート制限を個別に設定し、まず小さなテーブルでエンドツーエンドの信頼性を確立してから大きなテーブルを処理します。コネクタはスナップショット中もWALのドレイン(排出)を継続しなければならず、そうでなければ保持WALが増大します。選択したコネクタがインクリメンタルスナップショット中のスキーマ変更をサポートしていない場合は、そのテーブルのDDLを凍結するか、スナップショットを一時停止します。

ステップ4: イベント、順序付け、およびべき等な適用の定義

正規化されたイベントには、少なくとも以下が含まれます。

text
ChangeEvent {
  event_id
  source_id
  table_id
  primary_key
  operation          // READ | CREATE | UPDATE | DELETE
  before
  after
  source_lsn
  transaction_id
  transaction_order
  source_commit_time
  schema_version
  captured_at
}

1つのソース内では、source_lsnにトランザクション順序を加えることでイベント順序が確立されます。異なるデータベースのLSNは座標系を共有しないため、バージョンキーにはsource_idを含める必要があります。各(source_id, table_id, primary_key)に対して、ターゲットは最後に適用されたソースバージョンを保存します。シンクは、行を変更することなく、同一またはより古いバージョンに対して確認応答を返します。より新しいバージョンが来た場合、単一のターゲットトランザクションでビジネス行と適用済みバージョンの両方を更新します。

これにより、主要なat-least-once障害ウィンドウが処理可能になります。

  • イベントが高耐久ストリームに到達したが、そのソースオフセットが記録されていない場合:リカバリによって再発行され、シンクは同一バージョンを無視します。
  • ターゲットのMERGEはコミットされたが確認応答が失われた場合:バッチがリプレイされ、同じ状態に収束します。
  • バッチの途中でコンシューマがクラッシュした場合:コミット済みの行は安全にリプレイされ、未コミットの行は再処理されます。

宛先が複数行トランザクションのアトミック性を必要とする場合は、トランザクション境界メタデータを有効にし、ENDまでtransaction_idごとにバッファリングしてから、完全なトランザクションをまとめてコミットします。大規模なトランザクションはより多くのメモリを消費し、テールレイテンシを増加させます。また、タイムアウトやリカバリロジックによってトランザクションが完全かどうかを判断する必要があります。行レベルの結果整合性で十分な通常の分析用ウェアハウスでは、真の要件がない限りこの複雑さを受け入れるべきではありません。

ステップ5: 削除、キー変更、およびスキーマ進化の正確な処理

DELETEイベントは、キュレーションレイヤーが行を物理削除(ハードデリート)するか、is_deletedを設定するか、履歴を保持するために十分なキー情報を保持している必要があります。削除に続くトゥームストーンは、主にログ圧縮をサポートするためのものです。業務上の削除イベントの代わりにはならず、ミドルウェアはターゲットに到達する前に削除イベントを破棄してはなりません。

主キー値が変更された場合、一般的なCDCセマンティクスでは古いキーの削除と新しいキーの作成が発行されます。シンクはその両方を適用しなければならず、そうでなければ古いキーがゴースト行として残ってしまいます。主キーの定義を変更することは、値を変更することよりもリスクが高くなります。移行中はイベントキーの形状に不整合が生じる可能性があるため、読み取り専用または書き込み一時停止のウィンドウを設け、コネクタをドレインし、スキーマを更新してから再開します。

明示的な互換性ポリシーを使用します。

  • NULL許容カラムの追加:新しいバージョンを登録し、古いコンシューマには未知のフィールドを無視させ、書き込みを有効にする前にキュレーションテーブルを拡張します。
  • カラムの削除またはリネーム:新しいフィールドを導入して値を設定し、すべてのコンシューマを移行した後に古いフィールドを削除します。
  • 型の拡張:ターゲット互換性チェックの後にアップグレードします。型の縮小やセマンティクスの変更は、互換性のないイベントを隔離(quarantine)します。
  • 主キーの変更:サイレントな自動進化ではなく、個別のマイグレーションとして処理します。

互換性のないイベントは隔離キューに入り、アラートを発報します。メインオフセットは、イベントが高耐久に保存され、リプレイ可能であり、明確な修正担当者が割り当てられた後にのみ進めることができます。不正なスキーマをサイレントにスキップすると、目に見えないデータの欠落が発生します。

ステップ6: リカバリ、リプレイ、およびリコンシリエーションを通常パスとして扱う

コネクタのリカバリには、永続化されたオフセットと、対応する履歴をまだ保持しているレプリケーションスロットの両方が必要です。restart_lsnconfirmed_flush_lsn、現在のWAL位置、保持バイト数、増加率、およびディスク枯渇までの時間を監視します。宛先の停止による影響は、ソースで無制限にWALが保持されるほどコネクタの確認応答を止めるのではなく、高耐久ストリーム内に蓄積させる必要があります。

スロットが失われた場合、保存されたオフセットがスロットの利用可能位置より遅れている場合、または必要なWALが消失している場合は、フェイルクローズします。該当ソースのキュレーション公開を停止し、最後に信頼できた位置を記録し、影響を受けるテーブルを再スナップショットし、範囲ごとにリコンシリエーションを行います。新しく作成されたスロットは作成後の変更のみをキャプチャするため、それ以前の区間が完全であることを証明できません。

リプレイには、独自のreplay_job_id、テーブルおよび主キーのスコープ、ソース時間の境界があります。イミュータブルなRAWレイヤーからステージングへ読み込みます。古いバックフィルが新しい状態を上書きしないように、ライブデータとリプレイデータは同じソースバージョンMERGEルールを使用します。変換ロジックが変更された場合は、まず新しいキュレーションバージョンまたはシャドウテーブルを構築します。本番を即座に上書きするよりも、比較してロールバックできるようにする方が安全です。

リコンシリエーションは、少なくとも次の4つのレイヤーをカバーします。

  1. ソースコミット位置からコネクタ、高耐久ストリーム、ターゲット適用位置までの連続性。
  2. テーブル、日付、および主キーバケットごとの行数、削除数、およびチェックサム。
  3. ソースの現在の状態と宛先の最新状態との間での主キーサンプリング比較。
  4. SLA内でのINSERT、UPDATE、およびDELETEの可視性を検証する、定期的で識別可能なカナリアトランザクション。

エンドツーエンドのレイテンシを、ソースコミットからキャプチャ、キャプチャからストリーム、ストリームからステージング、ステージングからキュレーションに分解します。また、ソースごとのイベントレート、LSNラグ、保持スロットWAL、スナップショットチャンクの進捗、重複または古いバージョン、スキーマ隔離、ターゲットMERGE障害、リコンシリエーション差分、および推定追いつき時間も監視します。

ステップ7: 代替案の境界を説明する

updated_atによるポーリングは、書き込みが少なく、分単位の鮮度で十分であり、再スキャンが許容されるテーブルにはシンプルです。ただし、クエリ負荷が増加し、ハードデリートを見落としやすく、タイムスタンプの同一値や時計の精度の問題に対処する必要があります。トリガーは削除を監査テーブルに書き込むことができますが、ソースの書き込みパスに余分な作業と運用の結合をもたらします。本問のような高負荷OLTPワークロードには、ログベースのCDCがより適切なデフォルトです。

Outboxパターンは別の問題を解決します。サービスがビジネス状態とドメインイベントを単一のデータベーストランザクション内で書き込み、データベースコミット成功後のメッセージ送信失敗を防ぎます。これは特定のビジネスイベントには適しています。しかし、200テーブルにわたる汎用的な行レベルの同期を自動的に置き換えるものではありません。両者は共存可能です。サービス統合はOutboxを消費し、分析や監査はCDCを消費します。

質の高い模範解答

「私は保証すべき事項から始めます。WALまたは高耐久ストリームにまだ存在するすべてのコミット済み変更は回復可能でなければなりません。行はソース位置によって収束し、リプレイによって最終状態が狂うことはなく、削除はターゲットに到達します。レプリケーションスロットと履歴位置が損なわれた場合、公開は停止され、影響を受けるデータは再スナップショットされます。コネクタがサイレントに最新位置から再起動することはありません。

平均負荷では、イングレスは約20 MB/s、1日あたり17.28億イベントおよび1.728 TBです。ピーク時は60 MB/sであるため、宛先が30分停止すると約108 GBのバックログが発生します。7日間のRAWデータ保持は約12.096 TBです。8 TBのスナップショットを72時間以内に完了するには、少なくとも30.9 MB/sの実効読み取り速度が必要であるため、ベンチマークを実施し、p99、I/O、および保持WALに基づいて各ソースにレート制限を設けます。

各PostgreSQLソースには、独立したロジカルスロットとコネクタを割り当てます。コネクタは一貫性のあるスナップショットを取得し、対応するLSNからコミット済み変更を再開します。イベントはソース、テーブル、主キーによって高耐久ストリームにパーティション化され、イミュータブルなRAWレイヤーにコピーされます。軽量な処理でスキーマを検証し、エンベロープを正規化します。ウェアハウスはステージングテーブルをロードし、主キーとソースバージョンによってべき等なMERGEまたはDELETE操作を実行します。LSNがデータベース間で比較されることはありません。トランザクション順序はトランザクション内のLSNを補完します。

スナップショットは主キー範囲でチャンク化され、進捗は永続化され、ソース負荷は動的に制限されます。スナップショット実行中もWALはドレインされ、ターゲットのソースバージョンルールにより境界でのリカバリリプレイが吸収されます。イベントには、操作区分、前後の値、ソースLSN、トランザクション識別子と順序、およびスキーマバージョンが含まれます。シンクは新しいバージョンのみを適用するため、コネクタのリプレイ、ターゲット確認応答の喪失、コンシューマのクラッシュがあってもすべて安全に収束します。

削除イベントはキュレーションレイヤーまでそのまま保持されます。トゥームストーンは圧縮専用です。主キー値の変更は、古いキーの削除と新しいキーの作成として適用されます。NULL許容カラムの追加は互換性を保って進化できますが、削除、リネーム、型の縮小、キー変更には制御されたマイグレーションを使用します。互換性のないレコードは、サイレントにスキップされるのではなく、隔離されてリプレイ可能になります。

運用面では、セグメント化されたレイテンシ、ソースごとのLSNラグ、保持スロットWALとディスク空き容量、スナップショットの進捗、スキーマ隔離、ターゲット障害、およびリコンシリエーション差分を監視します。スロット喪失時はフェイルクローズし、影響を受けるスコープの再スナップショットをトリガーします。最後に、スナップショット中の同時書き込み、重複配信、削除、キー変更、コネクタクラッシュ、シンク確認応答の喪失、30分間のシンク停止、スロット喪失、破壊的なスキーマ変更などの障害注入テストを実施します。行数、キーバケットのチェックサム、サンプリングされた行、およびカナリアトランザクションにより、変更の欠落がなく、古いイベントが新しい状態を上書きしないことを証明します。」

よくある間違い

  • 全テーブルをエクスポートしてからCDCを有効化する → エクスポート期間中のWALが消失し、スナップショットとストリームの間にギャップが生じる可能性があります → 一致するポイントからストリーミングする前に、ログ保持と一貫性のあるスナップショット位置を確立してください。
  • 「DebeziumとKafkaを使う」とだけ言う → 製品名はレイテンシ、順序付け、リカバリ、シンクのべき等性を定義しません → 正確性の不変条件とキャパシティを最初に定義し、各コンポーネントをそれらにマッピングしてください。
  • 異なるデータベースのLSNを1つのグローバルバージョンとして比較する → すべてのソースは独立したログ座標を持っています → バージョンにsource_idを含め、単一ソースのシーケンス内でのみ位置を比較してください。
  • ブローカーのexactly-onceによってウェアハウスもexactly-onceになると想定する → 外部のMERGEはブローカートランザクションの外側に位置する可能性があり、確認応答の喪失後にリプレイされます → 主キーとソースバージョンによってべき等に変更を適用してください。
  • 停止したコネクタを無期限に待機させる → そのスロットがソースストレージを枯渇させるまでWALを保持し続ける可能性があります → 保持バイト数と枯渇までの時間を監視し、高耐久ストリームでシンク障害を隔離し、損切り(stop-loss)閾値を設定してください。
  • 失われたスロットを再作成して最新位置から継続する → 新しいスロットには過去の履歴がなく、データ欠落を隠蔽する可能性があります → フェイルクローズし、再スナップショットを行い、影響を受けるスコープをリコンシリエーションしてください。
  • トゥームストーンを唯一の削除シグナルとして扱う → 圧縮マーカーは、行キーを運ぶ業務上のDELETEの代わりにはなりません → 削除イベントを保持し、宛先でハードデリート、ソフトデリート、または履歴保持を選択してください。
  • すべてのスキーマ変更を自動的に受け入れる → 削除、型の縮小、キー変更は、コンシューマやイベントキーを破壊する可能性があります → 互換性マトリクスを定義し、互換性のないレコードを隔離し、破壊的変更のマイグレーションを段階的に行ってください。
  • 72時間のスナップショット期限のみに最適化する → 大規模なスキャンはソースのp99、キャッシュの挙動、レプリケーションを損なう可能性があります → 30.9 MB/sの下限値をベンチマークし、動的にレート制限を行い、ソースのバジェットを保護してください。
  • 合計行数のみを比較する → 誤って適用された更新、見落とされた削除、相殺されたエラーが隠れたままになる可能性があります → キーバケットごとの件数とチェックサムを比較し、行をサンプリングし、カナリアを注入してください。

フォローアップ質問

フォローアップ1: スナップショット実行中の更新や削除によって、古いスナップショット値が新しい状態を上書きしないことをどのように証明しますか?

アプリケーションのタイミングを推測するのではなく、コネクタの検証済みの一貫性スナップショットとログハンドオフのセマンティクスを使用します。シンクは行ごとにソースバージョンを保存し、より新しいバージョンのみを受け入れるため、リプレイされたREADがより新しいUPDATEやDELETEを上書きすることはできません。テストでは、スナップショットチャンクの前後で1つのキーを繰り返し更新、削除、再作成し、最終的なソースとターゲットの一致およびソース位置の連続性をアサート(検証)します。

フォローアップ2: ウェアハウスが30分間停止した後のリカバリ時間をどのように見積もりますか?

ピーク時のバックログは約108 GBです。リカバリ時間は、復旧したスループットから継続するライブイングレスを引いた値に依存します。ライブトラフィックが平均20 MB/sに戻り、コンシューマが80 MB/sを維持できる場合、正味の追いつきレートは約60 MB/sとなり、MERGE増幅と安全マージンを考慮する前の理論上のドレイン時間は約30分です。処理能力がライブイングレスと同等にしかならない場合、バックログは決して減少しません。

フォローアップ3: 1つのソーストランザクション内のすべての行をまとめて可視化する必要がある場合、何が変化しますか?

トランザクション境界メタデータを有効にし、トランザクションIDごとにイベントをアセンブルし、完全なEND境界の後にのみステージングと公開へトランザクションをコミットします。大規模トランザクションをディスクに永続化し、タイムアウトと修復パスを定義し、高耐久なアセンブリ状態から再開できるようにします。これによりテールレイテンシと状態保持コストが増加するため、テーブル間の一時的な不整合がコンシューマにとって真に許容できないものであるかをまず確認してください。

フォローアップ4: レプリケーションスロットが失われたものの、高耐久ストリームにまだ7日分のデータがある場合、フルスナップショットを実行する必要がありますか?

まず、発生しうるギャップを特定します。高耐久ストリームが最後に信頼できたLSN以降のすべての変更を保持していることが証明できる場合、フルスナップショットを行わずにリプレイとリコンシリエーションによって連続性を復元できます。イベントがストリームに到達する前にスロットが消失した場合、またはギャップ境界を証明できない場合は、影響を受けるテーブルまたは主キー範囲を再スナップショットします。この判断は、手戻りコストではなく、証明可能な連続性に基づいて行います。

フォローアップ5: 設計全体をOutboxパターンに置き換えないのはなぜですか?

Outboxパターンは、サービスが注文作成や支払い完了などのドメインイベントを明示的に発行し、ビジネス状態と同じトランザクションでイベントを書き込む場合に優れています。これにはアプリケーション側の実装が必要であり、選択された事実のみが含まれます。本問題は分析のために200テーブルにわたるINSERT、UPDATE、DELETEを同期するものであるため、汎用的な行レベルのCDCが依然として必要です。異なるコンシューマ向けに両方のパターンが共存できます。

公開情報ソース

関連する質問