使用 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. 根據工作負載需求選擇頻率 (提交 ID 或時間戳記範圍),執行排定的批次 CDC 工作,以便後續更新資料。

事前準備

如要設定排定的批次 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 bucket 中的有效 Delta Lake 資料表。資料表目錄必須包含資料檔案和 _delta_log/ 交易記錄目錄。此外,資料表必須符合下列規定:

    • Delta Lake 資料表必須啟用 delta.enableChangeDataFeed 屬性 (true)。
    • 為 CDC 管道設定的開始提交 ID 或時間戳記,必須晚於 delta.enableChangeDataFeed 屬性啟用時間。
  4. 現有的 Lakehouse Iceberg 目錄,用於接收同步處理的資料。如果目標命名空間不存在,CDC 管道會自動建立。如果目標 Lakehouse 資料表不存在,只要您提供有效的等值資料欄清單,CDC 管道就會自動建立該資料表。

支援與限制

從 Delta Lake 遷移至 Lakehouse 的排定批次變更資料擷取作業時,請注意下列事項:

  • 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)。
  • 自動變更資料動態饋給:系統不支援自動變更資料動態饋給,因為這項功能不會寫入 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 bucket 名稱。
  • 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 bucket 名稱。
  • 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:Lakehouse 資料表的命名空間。
    • TABLE_NAME:目標 Lakehouse 資料表的名稱。
  3. 按一下「執行」,然後驗證結果。

後續步驟