使用 Dataflow 將 Apache Iceberg 變更資料擷取 (CDC) 資料遷移至 Lakehouse

要讓目標 Lakehouse Apache Iceberg 資料表與經常更新的 Apache Iceberg 來源資料表保持同步,可能會有困難。執行完整資料表重新載入作業來擷取資料集變更,會增加運算成本和處理延遲時間,而建構自訂管道來處理變更資料擷取 (CDC) 記錄,則會造成複雜的作業負擔。

使用 Dataflow 批次或串流 CDC 管道,即可持續或定期將 Apache Iceberg 來源資料表與 Lakehouse Apache Iceberg 資料表同步。CDC 管道會直接從來源 Apache Iceberg 變更記錄中擷取增量更新,並套用至目標資料表,不必重新載入完整資料表。

CDC 管道提供下列功能:

  • 經濟實惠的增量同步:只處理變更的資料 (INSERT、UPDATE 和 DELETE 作業),減少運算負擔和資料延遲。
  • 彈性執行模式:支援批次和串流模式,可滿足您的業務需求。

CDC 管道支援兩種執行模式:

  • 批次 CDC 模式:在指定的開始和結束快照界限之間,定期從來源 Apache Iceberg 資料表讀取增量更新。
  • 串流 CDC 模式:持續以可設定的頻率 (預設為 1 分鐘,最短為 1 秒) 輪詢來源 Apache Iceberg 資料表,找出新的變更提交內容,並近乎即時地將這些內容套用至 Lakehouse 資料表。

事前準備

如要設定從 Apache Iceberg 遷移至 Lakehouse 的 CDC 資料,請確認您具備下列條件:

  1. 如果尚未啟用 Dataflow、BigQuery 和 Lakehouse API,請啟用這些 API。

    啟用 API 時所需的角色

    如要啟用 API,您必須具備 serviceusage.services.enable 權限。如果您建立了專案,可能已透過「擁有者」角色 (roles/owner) 取得這項權限。否則,您可以透過「服務使用管理員」角色 (roles/serviceusage.serviceUsageAdmin) 取得這項權限。瞭解如何授予角色。

    啟用 API

  2. 如要取得建立資源所需的權限,請要求管理員授予您專案的必要 Identity and Access Management (IAM) 角色。

  3. 儲存在 Cloud Storage 中,或在 Iceberg REST 目錄或 metastore 中註冊的有效來源 Apache Iceberg 資料表。

  4. 現有的 Lakehouse Iceberg 目錄,用於接收同步處理的資料。

    如果目標命名空間不存在,CDC 管道會自動建立。如果目標 Lakehouse 資料表不存在,只要您提供有效的等值資料欄清單,CDC 管道就會自動建立該資料表。

支援與限制

從 Apache Iceberg 遷移至 Lakehouse 的 CDC 資料時,請注意下列事項:

  • SDK 規定:必須使用 Apache Beam SDK 2.77.0 以上版本。
  • 建立目的地資料表:系統支援建立目的地 Lakehouse 資料表。如果目標資料表不存在,只要您提供有效的等值資料欄清單 (equality_columns),CDC 管道就會自動建立該資料表。您只需要現有的 Lakehouse 目錄,CDC 管道就會在命名空間不存在時自動建立。
  • 漸進式結構定義更新:系統不支援自動變更結構定義。如果來源 Apache Iceberg 資料表的結構定義有所變更,您必須先手動更新目標 Lakehouse 資料表結構定義,再執行 CDC 同步工作。
  • 執行模式:
    • 批次 CDC:在指定開始和結束快照界限 (快照 ID 或時間戳記) 之間,每次執行工作時擷取一次交易記錄變更。
    • 串流 CDC:持續輪詢來源 Apache Iceberg 資料表,取得新的提交記錄。輪詢頻率可設定 (預設為 1 分鐘,最少 1 秒)。

使用 Java 建立 CDC 管道

如要使用 Apache Beam Java SDK,以程式輔助方式建構及執行 Apache Iceberg CDC 批次或串流遷移管道,請搭配使用 Managed.read 轉換和 Managed.ICEBERG_CDC。

下列 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 bucket 名稱。
  • PROJECT_ID:您的 Google Cloud 專案 ID。
  • TARGET_NAMESPACE:目標 Lakehouse 資料表的命名空間。
  • TARGET_TABLE:目標 Lakehouse 資料表的名稱。

串流 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 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 bucket 名稱。
  • PROJECT_ID:您的 Google Cloud 專案 ID。
  • TARGET_NAMESPACE:目標 Lakehouse 資料表的命名空間。
  • 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 資料表的命名空間。
    • TABLE_NAME:目標 Lakehouse 資料表的名稱。
  3. 按一下「執行」,然後驗證結果。

後續步驟