要讓目標 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 資料,請確認您具備下列條件:
如果尚未啟用 Dataflow、BigQuery 和 Lakehouse API,請啟用這些 API。
啟用 API 時所需的角色
如要啟用 API,您必須具備
serviceusage.services.enable權限。如果您建立了專案,可能已透過「擁有者」角色 (roles/owner) 取得這項權限。否則,您可以透過「服務使用管理員」角色 (roles/serviceusage.serviceUsageAdmin) 取得這項權限。瞭解如何授予角色。如要取得建立資源所需的權限,請要求管理員授予您專案的必要 Identity and Access Management (IAM) 角色。
儲存在 Cloud Storage 中,或在 Iceberg REST 目錄或 metastore 中註冊的有效來源 Apache Iceberg 資料表。
現有的 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 資料表:
前往 Google Cloud 控制台的「BigQuery Studio」頁面。
在查詢編輯器中執行 SQL 查詢,驗證同步處理的資料:
SELECT * FROM `PROJECT_ID`.`CATALOG`.`NAMESPACE`.`TABLE_NAME` LIMIT 10;更改下列內容:
PROJECT_ID:您的 Google Cloud 專案 ID。CATALOG:Lakehouse 目錄的名稱。NAMESPACE:Lakehouse 資料表的命名空間。TABLE_NAME:目標 Lakehouse 資料表的名稱。
按一下「執行」,然後驗證結果。
後續步驟
- 進一步瞭解如何將 Iceberg 資料表匯入 Lakehouse。
- 瞭解 Lakehouse 中的變更資料擷取擷取作業。
- 進一步瞭解 Dataflow 中的受管理 I/O。