Dataflow を使用して Delta Lake の変更データ キャプチャ(CDC)データを Lakehouse に移行する

頻繁に更新される Delta Lake ソーステーブルとターゲットの Lakehouse Apache Iceberg テーブルを同期させるのは難しい場合があります。データセットの変更をキャプチャするためにテーブル全体を再読み込みすると、コンピューティング費用と処理レイテンシが増加します。また、トランザクション ログを処理するためにカスタム パイプラインを構築すると、複雑な運用オーバーヘッドが発生します。

Dataflow のスケジュール設定されたバッチ変更データ キャプチャ(CDC)パイプラインを使用すると、Delta Lake の移行元テーブルを Lakehouse の Apache Iceberg テーブルと継続的に同期できます。テーブル全体を再読み込みする代わりに、スケジュールされたバッチジョブは、Delta Lake トランザクション ログから増分更新を直接抽出し、構成可能なスケジュールでターゲット テーブルに適用します。

スケジュールされたバッチ CDC パイプラインには、次の機能があります。

  • 費用対効果の高い増分同期: 変更されたデータ(INSERT、UPDATE、DELETE オペレーション)のみを定期的に処理し、コンピューティングのオーバーヘッドとデータ レイテンシを削減します。
  • 変更の自動マッピング: Delta Lake の変更データ フィードログを解析し、変更オペレーションをターゲットの Lakehouse Iceberg テーブルにマッピングします。
  • 柔軟なジョブ スケジューリング: ビジネスニーズに合ったスケジュール(15 分間隔など)でバッチモードの Dataflow ジョブを実行します。

継続的な移行ワークフロー

Delta Lake から Lakehouse へのエンドツーエンドの継続的移行を設定するには、次の操作を順に行います。

  1. 入力 Delta Lake テーブルで delta.enableChangeDataFeed プロパティ(true)を有効にします。
  2. delta.enableChangeDataFeed プロパティを有効にした後、commit ID またはタイムスタンプへの最初のテーブル移行を実行します。手順については、Dataflow を使用して Delta Lake テーブルを Lakehouse にインポートするをご覧ください。
  3. ワークロードの要件に合わせて選択した頻度(コミット ID またはタイムスタンプ範囲)で、後続のデータ更新用にスケジュールされたバッチ CDC ジョブを実行します。

始める前に

スケジュール設定されたバッチ CDC 移行を設定するには、次のものが必要です。

  1. Dataflow API、BigQuery API、Lakehouse API が有効になっていない場合は、有効にします。

    API を有効にするために必要なロール

    API を有効にするには、serviceusage.services.enable 権限が必要です。プロジェクトを作成した場合は、オーナーロール(roles/owner)を介してこの権限がすでに付与されている可能性があります。それ以外の場合は、Service Usage 管理者ロール(roles/serviceusage.serviceUsageAdmin)を介してこの権限を取得できます。ロールを付与する方法を確認する。

    API を有効にする

  2. リソースの作成に必要な権限を取得するには、プロジェクトに必要な Identity and Access Management(IAM)ロールを付与するよう管理者に依頼してください。

  3. Cloud Storage バケットに保存されている有効な Delta Lake テーブル。テーブル ディレクトリには、データファイルと _delta_log/ トランザクション ログ ディレクトリが含まれている必要があります。また、テーブルは次の要件を満たしている必要があります。

    • Delta Lake テーブルで delta.enableChangeDataFeed プロパティを有効にする(true)必要があります。
    • CDC パイプラインに構成された開始コミット ID またはタイムスタンプは、delta.enableChangeDataFeed プロパティが有効になった後の値である必要があります。
  4. 同期されたデータを受け取る既存の Lakehouse Iceberg カタログ。ターゲット Namespace が存在しない場合、CDC パイプラインは自動的に作成します。ターゲット Lakehouse テーブルが存在しない場合、等価列の有効なリストを指定すると、CDC パイプラインによって自動的に作成されます。

サポートと制限事項

Delta Lake から Lakehouse へのバッチ CDC 移行のスケジュール設定には、次の考慮事項があります。

  • SDK の要件: Apache Beam SDK バージョン 2.77.0 以降が必要です。
  • 宛先テーブルの作成: 宛先 Lakehouse テーブルの作成がサポートされています。ターゲット テーブルが存在しない場合、等価列(equality_columns)の有効なリストを指定すると、CDC パイプラインによって自動的に作成されます。既存の Lakehouse カタログのみが必要です。CDC パイプラインは、名前空間が存在しない場合は自動的に作成します。
  • スキーマ進化: スキーマの自動変更はサポートされていません。ソースの Delta Lake テーブルのスキーマが変更された場合は、CDC 同期ジョブを実行する前に、ターゲットの Lakehouse テーブルのスキーマを手動で更新して一致させる必要があります。
  • バッチ実行モード: Delta Lake CDC はバッチモードで動作します。パイプラインは、指定された開始 commit 境界と終了 commit 境界(commit ID またはタイムスタンプ)の間で、ジョブ実行ごとにトランザクション ログを 1 回取得します。終了 commit の境界(end_version または end_timestamp)が指定されていない場合、パイプラインは最新の commit バージョンまで読み取ります。
  • サポートされているソース ストレージ: ソースデータは、Cloud Storage(gs://)に保存されている有効な Delta Lake テーブルである必要があります。Amazon S3 は、Delta Lake テーブル ソースではサポートされていません。
  • カタログベースのテーブル: カタログベースの Delta Lake テーブル(Unity Catalog など)はサポートされていません。
  • サポートされていない Delta Lake バージョン: Delta Lake バージョン 1.2.1 以前は、これらのバージョンで変更データフィード(CDF)がサポートされていないため、サポートされていません。
  • 自動変更データフィード: CDF 変更ログファイルが書き込まれないため、自動変更データフィードはサポートされていません。

プログラムで CDC パイプラインを作成する

Delta Lake CDC バッチ移行パイプラインをプログラムで構築して実行するには、Apache Beam パイプラインで Managed.read 変換を使用します。

Java

次の依存関係を pom.xml ファイルに追加します。

<dependency>
  <groupId>org.apache.beam</groupId>
  <artifactId>beam-sdks-java-managed</artifactId>
  <version>${beam.version}</version>
</dependency>
<dependency>
  <groupId>org.apache.beam</groupId>
  <artifactId>beam-sdks-java-io-iceberg</artifactId>
  <version>${beam.version}</version>
</dependency>

次の Java スニペットは、Delta Lake から Lakehouse へのバッチ CDC パイプラインを構成して実行する方法を示しています。

import java.util.Arrays;
import java.util.HashMap;
import java.util.Map;
import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.managed.Managed;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.values.PCollection;
import org.apache.beam.sdk.values.Row;

public class DeltaLakeCdcToLakehouse {
  public static void main(String[] args) {
    Pipeline pipeline = Pipeline.create(PipelineOptionsFactory.fromArgs(args).create());

    // 1. Configure Delta Lake CDC reader
    Map<String, Object> readConfig = new HashMap<>();
    readConfig.put("table", "gs://BUCKET_NAME/TABLE_NAME");
    // Starting commit bound (either start_version or start_timestamp must be provided):
    readConfig.put("start_version", START_COMMIT_ID);
    // Optional ending commit bound (end_version or end_timestamp, matching the start bound type):
    // readConfig.put("end_version", END_COMMIT_ID);
    readConfig.put("include_metadata_columns", Arrays.asList("_commit_version", "_change_type"));

    // Read CDC events from Delta Lake
    PCollection<Row> cdcRows =
        pipeline.apply(
            "ReadFromDeltaCDC",
            Managed.read(Managed.DELTA_LAKE_CDC).withConfig(readConfig)).getSinglePCollection();

    // 2. Configure BigLake REST Catalog properties for Lakehouse Iceberg
    Map<String, String> catalogProps = new HashMap<>();
    catalogProps.put("type", "rest");
    catalogProps.put("uri", "https://biglake.googleapis.com/iceberg/v1/restcatalog");
    catalogProps.put("warehouse", "gs://WAREHOUSE_BUCKET");
    catalogProps.put("header.x-goog-user-project", "PROJECT_ID");
    catalogProps.put("io-impl", "org.apache.iceberg.gcp.gcs.GCSFileIO");
    catalogProps.put("rest.auth.type", "org.apache.iceberg.gcp.auth.GoogleAuthManager");

    // 3. Map Delta Lake change types to Iceberg CDC change types
    Map<String, String> changeTypeMap = new HashMap<>();
    changeTypeMap.put("insert", "INSERT");
    changeTypeMap.put("delete", "DELETE");
    changeTypeMap.put("update_preimage", "UPDATE_BEFORE");
    changeTypeMap.put("update_postimage", "UPDATE_AFTER");

    // 4. Configure Lakehouse Iceberg CDC writer
    Map<String, Object> writeConfig = new HashMap<>();
    writeConfig.put("table", "TARGET_NAMESPACE.TARGET_TABLE");
    writeConfig.put("catalog_name", "lakehouse");
    writeConfig.put("catalog_properties", catalogProps);
    writeConfig.put("mode", "merge-on-read");
    writeConfig.put("sequence_number_column", "_commit_version");
    writeConfig.put("equality_columns", Arrays.asList("PRIMARY_KEY_COLUMN"));
    // change_type_column and change_type_map are only required for Portable Runner
    writeConfig.put("change_type_column", "_change_type");
    writeConfig.put("change_type_map", changeTypeMap);

    // Write CDC rows to target Lakehouse Iceberg table
    cdcRows.apply(
        "WriteToLakehouse",
        Managed.write(Managed.ICEBERG).withConfig(writeConfig));

    pipeline.run();
  }
}

次のように置き換えます。

  • BUCKET_NAME: Delta Lake テーブルを含む Cloud Storage バケットの名前。
  • TABLE_NAME: ソースの Delta Lake テーブル ディレクトリの名前。
  • START_COMMIT_ID: Delta Lake CDC イベントの読み取りを開始するコミット バージョン番号。start_version または start_timestamp のいずれかを指定する必要があります。開始コミット ID またはタイムスタンプは、コピー元テーブルで delta.enableChangeDataFeed プロパティが有効になった後の値である必要があります。
  • END_COMMIT_ID: Delta Lake CDC イベントの読み取りの終了コミット バージョン番号。ベスト プラクティスとして、開始(start_version または start_timestamp)に使用される境界型と一致する終了境界(end_version または end_timestamp)を指定します。終了境界が省略されている場合、パイプラインは最新のコミット バージョンまで読み取ります。
  • WAREHOUSE_BUCKET: Lakehouse カタログ ウェアハウスとして使用される Cloud Storage バケットの名前。
  • PROJECT_ID: 実際の Google Cloud プロジェクト ID。
  • TARGET_NAMESPACE: ターゲット Lakehouse テーブルの Namespace。
  • TARGET_TABLE: ターゲット Lakehouse テーブルの名前。
  • PRIMARY_KEY_COLUMN: CDC の更新と削除の行の識別に使用される主キー列名。

ジョブの出力を調べる

CDC データが Lakehouse テーブルに正常に統合されたことを確認します。

  1. Google Cloud コンソールで、BigQuery Studio ページに移動します。

    BigQuery に移動

  2. クエリエディタで SQL クエリを実行して、同期されたデータを確認します。

    SELECT * FROM `PROJECT_ID`.`CATALOG`.`NAMESPACE`.`TABLE_NAME` LIMIT 10;
    

    次のように置き換えます。

    • PROJECT_ID: 実際の Google Cloud プロジェクト ID。
    • CATALOG: Lakehouse カタログの名前。
    • NAMESPACE: Lakehouse テーブルの Namespace。
    • TABLE_NAME: ターゲット Lakehouse テーブルの名前。
  3. [実行] をクリックして結果を確認します。

次のステップ