使用 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 API、BigQuery API 和 Lakehouse API(如果尚未启用)。

    启用 API 所需的角色

    如需启用 API,您需要拥有 serviceusage.services.enable 权限。如果您是项目创建者,则可能已通过 Owner 角色 (roles/owner) 获得此权限。否则,您可以通过 Service Usage Admin 角色 (roles/serviceusage.serviceUsageAdmin) 获得此权限。了解如何授予角色。

    启用 API

  2. 如需获得创建资源所需的权限,请让管理员向您授予项目的必要 Identity and Access Management (IAM) 角色。

  3. 存储在 Cloud Storage 中或在 Iceberg REST Catalog 或 metastore 中注册的有效源 Apache Iceberg 表。

  4. 用于接收同步数据的现有 Lakehouse Iceberg 目录。

    如果目标命名空间不存在,CDC 流水线会自动创建它。如果目标 Lakehouse 表不存在,并且您提供了有效的等值列列表,CDC 流水线会自动创建该表。

支持和限制

从 Apache Iceberg 迁移到 Lakehouse 的 CDC 需注意以下事项:

  • SDK 要求:需要使用 Apache Beam SDK 2.77.0 版或更高版本。
  • 目标表创建:支持创建目标 Lakehouse 表。如果目标表不存在,CDC 流水线会自动创建该表,前提是您提供有效的等值列列表 (equality_columns)。您只需提供现有的 Lakehouse 目录;如果命名空间不存在,CDC 流水线会自动创建该命名空间。
  • 架构演变:不支持自动架构更改。如果源 Apache Iceberg 表的架构发生变化,您必须先手动更新目标湖仓一体表架构以匹配,然后才能运行 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 存储桶的名称。
  • 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 存储桶的名称。
  • 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:湖仓一体表的命名空间。
    • TABLE_NAME:目标 Lakehouse 表的名称。
  3. 点击 Run(运行),然后验证结果。

后续步骤