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

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

Dataflow バッチまたはストリーミング CDC パイプラインを使用すると、Apache Iceberg ソーステーブルを Lakehouse Apache Iceberg テーブルと継続的または定期的に同期できます。CDC パイプラインは、テーブル全体を再読み込みするのではなく、移行元の Apache Iceberg 変更ログから増分更新を直接抽出し、移行先のテーブルに適用します。

CDC パイプラインには次の機能があります。

  • 費用対効果の高い増分同期: 変更されたデータ(INSERT、UPDATE、DELETE オペレーション)のみを処理し、コンピューティングのオーバーヘッドとデータ レイテンシを削減します。
  • 柔軟な実行モード: ビジネス要件に合わせてバッチモードとストリーミング モードの両方をサポートします。

CDC パイプラインは、次の 2 つの実行モードをサポートしています。

  • バッチ CDC モード: 指定された開始スナップショットと終了スナップショットの境界の間で、移行元 Apache Iceberg テーブルから増分アップデートを定期的に読み取ります。
  • ストリーミング CDC モード: 構成可能な頻度(デフォルトは 1 分、最小 1 秒)で、ソースの Apache Iceberg テーブルを新しい変更 commit について継続的にポーリングし、準リアルタイムで Lakehouse テーブルに適用します。

始める前に

Apache Iceberg から Lakehouse への 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 に保存されているか、Iceberg REST カタログまたは metastore に登録されている有効なソース Apache Iceberg テーブル。

  4. 同期されたデータを受け取る既存の Lakehouse Iceberg カタログ。

    ターゲット Namespace が存在しない場合、CDC パイプラインは自動的に作成します。ターゲット Lakehouse テーブルが存在しない場合、等価列の有効なリストを指定すると、CDC パイプラインによって自動的に作成されます。

サポートと制限事項

Apache Iceberg から Lakehouse への CDC 移行には、次の考慮事項があります。

  • SDK の要件: Apache Beam SDK バージョン 2.77.0 以降が必要です。
  • 宛先テーブルの作成: 宛先 Lakehouse テーブルの作成がサポートされています。ターゲット テーブルが存在しない場合、等価列(equality_columns)の有効なリストを指定すると、CDC パイプラインによって自動的に作成されます。既存の Lakehouse カタログのみが必要です。名前空間が存在しない場合は、CDC パイプラインによって自動的に作成されます。
  • スキーマ進化: スキーマの自動変更はサポートされていません。ソースの Apache Iceberg テーブルのスキーマが変更された場合は、CDC 同期ジョブを実行する前に、ターゲットの Lakehouse テーブルのスキーマを手動で更新して一致させる必要があります。
  • 実行モード:
    • バッチ CDC: 指定された開始スナップショットと終了スナップショットの境界(スナップショット ID またはタイムスタンプ)の間で、ジョブ実行ごとにトランザクション ログの変更を 1 回取得します。
    • ストリーミング CDC: 新しい commit ログについて、移行元の Apache Iceberg テーブルを継続的にポーリングします。ポーリング頻度は構成可能です(デフォルトは 1 分、最小は 1 秒)。

Java を使用して CDC パイプラインを作成する

Apache Beam Java SDK を使用して Apache Iceberg CDC バッチまたはストリーミング移行パイプラインをプログラムで構築して実行するには、Managed.ICEBERG_CDC で Managed.read 変換を使用します。

次の Java の例は、バッチ ワークロードとストリーミング ワークロード用に Apache Iceberg CDC パイプラインを構成して実行する方法を示しています。

バッチ CDC パイプライン

次の Java の例では、指定された開始スナップショットと終了スナップショットの範囲内で Apache Iceberg CDC イベントを読み取るバッチ パイプラインを構成します。

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 スニペットは、Iceberg テーブルから 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.io.iceberg.cdc.IcebergCdcMetadataColumns;
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 IcebergBatchCdcToLakehouse {
  public static void main(String[] args) {
    Pipeline pipeline = Pipeline.create(PipelineOptionsFactory.fromArgs(args).create());

    // Configure source catalog properties
    Map<String, String> sourceCatalogProps = new HashMap<>();
    sourceCatalogProps.put("type", "SOURCE_CATALOG_TYPE");
    sourceCatalogProps.put("warehouse", "SOURCE_WAREHOUSE_LOCATION");
    // Add additional properties if required by your catalog type (for example, uri)

    // Configure batch Iceberg CDC reader with snapshot bounds
    Map<String, Object> readConfig = new HashMap<>();
    readConfig.put("table", "SOURCE_NAMESPACE.SOURCE_TABLE");
    readConfig.put("catalog_properties", sourceCatalogProps);
    readConfig.put("start_snapshot_id", START_SNAPSHOT_ID);
    readConfig.put("end_snapshot_id", END_SNAPSHOT_ID);
    readConfig.put(
        "include_metadata_columns",
        Arrays.asList(
            IcebergCdcMetadataColumns.CHANGE_TYPE,
            IcebergCdcMetadataColumns.COMMIT_SNAPSHOT_SEQUENCE_NUMBER));

    // Read batch CDC events from source Iceberg table
    PCollection<Row> cdcRows =
        pipeline.apply(
            "ReadFromIcebergCDC",
            Managed.read(Managed.ICEBERG_CDC).withConfig(readConfig)).getSinglePCollection();

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

    // Configure Lakehouse Iceberg writer
    Map<String, Object> writeConfig = new HashMap<>();
    writeConfig.put("table", "TARGET_NAMESPACE.TARGET_TABLE");
    writeConfig.put("catalog_properties", destinationCatalogProps);
    writeConfig.put("mode", "merge-on-read");
    writeConfig.put("change_type_column", IcebergCdcMetadataColumns.CHANGE_TYPE);

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

    pipeline.run();
  }
}

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

  • SOURCE_CATALOG_TYPE: ソース Iceberg カタログのタイプ(hadoop や hive など)。
  • SOURCE_WAREHOUSE_LOCATION: ソース Iceberg カタログのウェアハウスのロケーション(例: s3://source-warehouse、gs://source-warehouse)。
  • SOURCE_NAMESPACE: ソース Iceberg テーブルの名前空間。
  • SOURCE_TABLE: ソース Iceberg テーブルの名前。
  • START_SNAPSHOT_ID: バッチ CDC 読み取りの開始スナップショット ID。
  • END_SNAPSHOT_ID: バッチ CDC 読み取りの終了スナップショット ID。
  • WAREHOUSE_BUCKET: Lakehouse カタログ ウェアハウスとして使用される Cloud Storage バケットの名前。
  • PROJECT_ID: 実際の Google Cloud プロジェクト ID。
  • TARGET_NAMESPACE: ターゲット Lakehouse テーブルの Namespace。
  • TARGET_TABLE: ターゲット Lakehouse テーブルの名前。

ストリーミング CDC パイプライン

次の Java の例では、新しい CDC の変更についてソースの Apache Iceberg テーブルを継続的にポーリングするストリーミング パイプラインを構成します。

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 スニペットは、Iceberg テーブルから 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.io.iceberg.cdc.IcebergCdcMetadataColumns;
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 IcebergStreamingCdcToLakehouse {
  public static void main(String[] args) {
    Pipeline pipeline = Pipeline.create(PipelineOptionsFactory.fromArgs(args).create());

    // Configure source catalog properties
    Map<String, String> sourceCatalogProps = new HashMap<>();
    sourceCatalogProps.put("type", "SOURCE_CATALOG_TYPE");
    sourceCatalogProps.put("warehouse", "SOURCE_WAREHOUSE_LOCATION");
    // Add additional properties if required by your catalog type (for example, uri)

    // Configure streaming Iceberg CDC reader with continuous polling
    Map<String, Object> readConfig = new HashMap<>();
    readConfig.put("table", "SOURCE_NAMESPACE.SOURCE_TABLE");
    readConfig.put("catalog_properties", sourceCatalogProps);
    readConfig.put("streaming", true);
    readConfig.put("poll_interval_seconds", POLL_INTERVAL_SECONDS);
    readConfig.put(
        "include_metadata_columns",
        Arrays.asList(
            IcebergCdcMetadataColumns.CHANGE_TYPE,
            IcebergCdcMetadataColumns.COMMIT_SNAPSHOT_SEQUENCE_NUMBER));

    // Read streaming CDC events from source Iceberg table
    PCollection<Row> cdcRows =
        pipeline.apply(
            "ReadFromIcebergCDC",
            Managed.read(Managed.ICEBERG_CDC).withConfig(readConfig)).getSinglePCollection();

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

    // Configure Lakehouse Iceberg writer
    Map<String, Object> writeConfig = new HashMap<>();
    writeConfig.put("table", "TARGET_NAMESPACE.TARGET_TABLE");
    writeConfig.put("catalog_properties", destinationCatalogProps);
    writeConfig.put("mode", "merge-on-read");
    writeConfig.put("change_type_column", IcebergCdcMetadataColumns.CHANGE_TYPE);

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

    pipeline.run();
  }
}

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

  • SOURCE_CATALOG_TYPE: ソース Iceberg カタログのタイプ(hadoop や hive など)。
  • SOURCE_WAREHOUSE_LOCATION: ソース Iceberg カタログのウェアハウスのロケーション(例: s3://source-warehouse、gs://source-warehouse)。
  • SOURCE_NAMESPACE: ソース Iceberg テーブルの名前空間。
  • SOURCE_TABLE: ソース Iceberg テーブルの名前。
  • POLL_INTERVAL_SECONDS: 継続的なポーリングの頻度(秒単位)(デフォルトは 60、最小値は 1)。
  • WAREHOUSE_BUCKET: Lakehouse カタログ ウェアハウスとして使用される Cloud Storage バケットの名前。
  • PROJECT_ID: 実際の Google Cloud プロジェクト ID。
  • TARGET_NAMESPACE: ターゲット Lakehouse テーブルの Namespace。
  • TARGET_TABLE: ターゲット Lakehouse テーブルの名前。

ジョブの出力を調べる

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. [実行] をクリックして結果を確認します。

次のステップ