使用 Dataflow 将 Delta Lake 变更数据捕获 (CDC) 数据迁移到 Lakehouse

使目标 Lakehouse Apache Iceberg 表与经常更新的 Delta Lake 源表保持同步可能很困难。执行全表重新加载来捕获数据集更改会增加计算费用和处理延迟时间,而构建自定义流水线来处理事务日志会带来复杂的运营开销。

通过使用 Dataflow 预定批处理变更数据捕获 (CDC) 流水线,您可以持续将 Delta Lake 源表与 Lakehouse Apache Iceberg 表同步。与重新加载整个表不同,预定的批量作业会直接从 Delta Lake 事务日志中提取增量更新,并按照可配置的时间表将其应用到目标表。

预定批处理 CDC 流水线提供以下功能:

  • 经济高效的增量同步:定期仅处理更改的数据(INSERT、UPDATE 和 DELETE 操作),从而减少计算开销和数据延迟时间。
  • 自动更改映射:解析 Delta Lake Change Data Feed 日志,并将更改操作映射到目标 Lakehouse Iceberg 表。
  • 灵活的作业调度:按符合您业务需求的计划(例如,每隔 15 分钟)执行批处理模式的 Dataflow 作业。

持续迁移工作流

如需设置从 Delta Lake 到 Lakehouse 的端到端持续迁移,请按顺序完成以下操作:

  1. 在输入 Delta Lake 表上启用 delta.enableChangeDataFeed 属性 (true)。
  2. 在启用 delta.enableChangeDataFeed 属性后,执行初始表迁移到提交 ID 或时间戳。如需相关说明,请参阅使用 Dataflow 将 Delta Lake 表导入 Lakehouse。
  3. 运行预定的批量 CDC 作业,以根据工作负载要求选择的频率(提交 ID 或时间戳范围)进行后续数据更新。

准备工作

如需设置定期批处理 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 存储桶中的有效 Delta Lake 表。表目录必须包含数据文件和 _delta_log/ 事务日志目录。此外,该表还必须满足以下要求:

    • 必须在 Delta Lake 表上启用 delta.enableChangeDataFeed 属性 (true)。
    • 为 CDC 流水线配置的开始提交 ID 或时间戳必须在启用 delta.enableChangeDataFeed 属性之后。
  4. 用于接收同步数据的现有 Lakehouse Iceberg 目录。如果目标命名空间不存在,CDC 流水线会自动创建它。如果目标 Lakehouse 表不存在,并且您提供了有效的等值列列表,CDC 流水线会自动创建该表。

支持和限制

从 Delta Lake 到 Lakehouse 的预定批量 CDC 迁移需注意以下事项:

  • SDK 要求:需要使用 Apache Beam SDK 2.77.0 版或更高版本。
  • 目标表创建:支持创建目标 Lakehouse 表。如果目标表不存在,并且您提供了有效的等值列列表 (equality_columns),CDC 流水线会自动创建该表。您只需提供现有的 Lakehouse 目录;如果命名空间不存在,CDC 流水线会自动创建该命名空间。
  • 架构演变:不支持自动架构更改。如果源 Delta Lake 表的架构发生变化,您必须手动更新目标 Lakehouse 表架构以匹配,然后才能运行 CDC 同步作业。
  • 批处理执行模式:Delta Lake CDC 以批处理模式运行。流水线在指定的开始和结束提交边界(提交 ID 或时间戳)之间,每次作业执行时拉取一次事务日志。如果未指定结束提交边界(end_version 或 end_timestamp),流水线会读取到最新的提交版本。
  • 支持的源存储空间:源数据必须是存储在 Cloud Storage (gs://) 中的有效 Delta Lake 表。Delta Lake 表源不支持 Amazon S3。
  • 基于目录的表:不支持基于目录的 Delta Lake 表(例如 Unity Catalog)。
  • 不支持的 Delta Lake 版本:不支持 Delta Lake 版本 1.2.1 或更低版本,因为这些版本不支持变更数据源 (CDF)。
  • 自动变更数据 Feed:不支持自动变更数据 Feed,因为它不会写入 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 表的命名空间。
  • 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:湖仓一体表的命名空间。
    • TABLE_NAME:目标 Lakehouse 表的名称。
  3. 点击 Run(运行),然后验证结果。

后续步骤