頻繁に更新される Delta Lake ソーステーブルとターゲットの Lakehouse Apache Iceberg テーブルを同期させるのは難しい場合があります。データセットの変更をキャプチャするためにテーブル全体を再読み込みすると、コンピューティング費用と処理レイテンシが増加します。また、トランザクション ログを処理するためにカスタム パイプラインを構築すると、複雑な運用オーバーヘッドが発生します。
Dataflow のスケジュール設定されたバッチ変更データ キャプチャ(CDC)パイプラインを使用すると、Delta Lake の移行元テーブルを Lakehouse の Apache Iceberg テーブルと継続的に同期できます。テーブル全体を再読み込みする代わりに、スケジュールされたバッチジョブは、Delta Lake トランザクション ログから増分更新を直接抽出し、構成可能なスケジュールでターゲット テーブルに適用します。
スケジュールされたバッチ CDC パイプラインには、次の機能があります。
- 費用対効果の高い増分同期: 変更されたデータ(
INSERT、UPDATE、DELETEオペレーション)のみを定期的に処理し、コンピューティングのオーバーヘッドとデータ レイテンシを削減します。 - 変更の自動マッピング: Delta Lake の変更データ フィードログを解析し、変更オペレーションをターゲットの Lakehouse Iceberg テーブルにマッピングします。
- 柔軟なジョブ スケジューリング: ビジネスニーズに合ったスケジュール(15 分間隔など)でバッチモードの Dataflow ジョブを実行します。
継続的な移行ワークフロー
Delta Lake から Lakehouse へのエンドツーエンドの継続的移行を設定するには、次の操作を順に行います。
- 入力 Delta Lake テーブルで
delta.enableChangeDataFeedプロパティ(true)を有効にします。 delta.enableChangeDataFeedプロパティを有効にした後、commit ID またはタイムスタンプへの最初のテーブル移行を実行します。手順については、Dataflow を使用して Delta Lake テーブルを Lakehouse にインポートするをご覧ください。- ワークロードの要件に合わせて選択した頻度(コミット ID またはタイムスタンプ範囲)で、後続のデータ更新用にスケジュールされたバッチ CDC ジョブを実行します。
始める前に
スケジュール設定されたバッチ CDC 移行を設定するには、次のものが必要です。
Dataflow API、BigQuery API、Lakehouse API が有効になっていない場合は、有効にします。
API を有効にするために必要なロール
API を有効にするには、
serviceusage.services.enable権限が必要です。プロジェクトを作成した場合は、オーナーロール(roles/owner)を介してこの権限がすでに付与されている可能性があります。それ以外の場合は、Service Usage 管理者ロール(roles/serviceusage.serviceUsageAdmin)を介してこの権限を取得できます。ロールを付与する方法を確認する。リソースの作成に必要な権限を取得するには、プロジェクトに必要な Identity and Access Management(IAM)ロールを付与するよう管理者に依頼してください。
Cloud Storage バケットに保存されている有効な Delta Lake テーブル。テーブル ディレクトリには、データファイルと
_delta_log/トランザクション ログ ディレクトリが含まれている必要があります。また、テーブルは次の要件を満たしている必要があります。- Delta Lake テーブルで
delta.enableChangeDataFeedプロパティを有効にする(true)必要があります。 - CDC パイプラインに構成された開始コミット ID またはタイムスタンプは、
delta.enableChangeDataFeedプロパティが有効になった後の値である必要があります。
- Delta Lake テーブルで
同期されたデータを受け取る既存の Lakehouse Iceberg カタログ。ターゲット Namespace が存在しない場合、CDC パイプラインは自動的に作成します。ターゲット Lakehouse テーブルが存在しない場合、等価列の有効なリストを指定すると、CDC パイプラインによって自動的に作成されます。
サポートと制限事項
Delta Lake から Lakehouse へのバッチ CDC 移行のスケジュール設定には、次の考慮事項があります。
- SDK の要件: Apache Beam SDK バージョン 2.77.0 以降が必要です。
- 宛先テーブルの作成: 宛先 Lakehouse テーブルの作成がサポートされています。ターゲット テーブルが存在しない場合、等価列(
equality_columns)の有効なリストを指定すると、CDC パイプラインによって自動的に作成されます。既存の Lakehouse カタログのみが必要です。CDC パイプラインは、名前空間が存在しない場合は自動的に作成します。 - スキーマ進化: スキーマの自動変更はサポートされていません。ソースの Delta Lake テーブルのスキーマが変更された場合は、CDC 同期ジョブを実行する前に、ターゲットの Lakehouse テーブルのスキーマを手動で更新して一致させる必要があります。
- バッチ実行モード: Delta Lake CDC はバッチモードで動作します。パイプラインは、指定された開始 commit 境界と終了 commit 境界(commit ID またはタイムスタンプ)の間で、ジョブ実行ごとにトランザクション ログを 1 回取得します。終了 commit の境界(
end_versionまたはend_timestamp)が指定されていない場合、パイプラインは最新の commit バージョンまで読み取ります。 - サポートされているソース ストレージ: ソースデータは、Cloud Storage(
gs://)に保存されている有効な Delta Lake テーブルである必要があります。Amazon S3 は、Delta Lake テーブル ソースではサポートされていません。 - カタログベースのテーブル: カタログベースの 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 バケットの名前。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 テーブルの Namespace。TARGET_TABLE: ターゲット Lakehouse テーブルの名前。PRIMARY_KEY_COLUMN: CDC の更新と削除の行の識別に使用される主キー列名。
ジョブの出力を調べる
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 テーブルの Namespace。TABLE_NAME: ターゲット Lakehouse テーブルの名前。
[実行] をクリックして結果を確認します。
次のステップ
- Delta Lake テーブルを Lakehouse にインポートするの詳細を確認する。
- Lakehouse での変更データ キャプチャの取り込みについて学習する。
- Dataflow のマネージド I/O の詳細を確認する。