ahmedabu98 commented on code in PR #40387: URL: https://github.com/apache/beam/pull/40387#discussion_r4187274616
########## examples/java/src/main/java/org/apache/beam/examples/DeltaLakeToLakehouseCdcExample.java: ########## @@ -0,0 +1,643 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.beam.examples; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import java.io.BufferedReader; +import java.io.IOException; +import java.io.InputStreamReader; +import java.nio.channels.Channels; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Comparator; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; +import org.apache.beam.sdk.Pipeline; +import org.apache.beam.sdk.PipelineResult; +import org.apache.beam.sdk.extensions.gcp.options.GcpOptions; +import org.apache.beam.sdk.io.FileSystems; +import org.apache.beam.sdk.io.fs.MatchResult; +import org.apache.beam.sdk.managed.Managed; +import org.apache.beam.sdk.options.Default; +import org.apache.beam.sdk.options.Description; +import org.apache.beam.sdk.options.ExperimentalOptions; +import org.apache.beam.sdk.options.PipelineOptionsFactory; +import org.apache.beam.sdk.options.Validation; +import org.apache.beam.sdk.values.PCollection; +import org.apache.beam.sdk.values.Row; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Splitter; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap; +import org.checkerframework.checker.nullness.qual.Nullable; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * An Apache Beam Java example pipeline that reads Change Data Feed (CDC) records from a Delta Lake + * table on Google Cloud Storage (GCS) and applies those row-level changes (`INSERT`, + * `UPDATE_BEFORE`, `UPDATE_AFTER`, `DELETE`) to a Google Cloud Platform (GCP) Lakehouse (BigLake + * Metastore Iceberg REST Catalog) table using Beam's {@link Managed} I/O connectors. + * + * <h2>Overview</h2> + * + * <p>This pipeline performs the following steps: + * + * <ol> + * <li><b>Validates Delta Lake Change Data Feed (CDF):</b> Inspects the input Delta Lake table's + * transaction log ({@code _delta_log/*.json}) to verify that the table property {@code + * delta.enableChangeDataFeed = true} is enabled, failing fast with an error if it is not. + * <li><b>Reads Delta Lake CDC Data:</b> Uses {@code Managed.read(Managed.DELTA_LAKE_CDC)} to read + * change events over either a commit version range ({@code --startVersion} / {@code + * --endVersion}) or an ISO-8601 timestamp range ({@code --startTimestamp} / {@code + * --endTimestamp}). Standard GCS Hadoop filesystem properties and Delta CDC metadata columns + * are configured by default. + * <li><b>Writes CDC Data to GCP Lakehouse:</b> Uses {@code Managed.write(Managed.ICEBERG)} in + * {@code merge-on-read} mode against the BigLake Iceberg REST catalog to apply inserts, + * updates, and deletes ordered by {@code _commit_version}. Standard BigLake REST catalog and + * Iceberg CDC sink properties are configured by default. + * </ol> + * + * <h2>Prerequisites</h2> + * + * <ul> + * <li><b>Java 17+:</b> Both the Delta Lake Kernel API and Iceberg require Java 17 or later. If + * using SDKMAN, switch to Java 17 before running: + * <pre>{@code + * sdk use java 17.0.15-tem + * }</pre> + * <li><b>Source Delta Lake Table on GCS:</b> Must have Change Data Feed enabled: + * <pre>{@code + * ALTER TABLE delta.`gs://my-bucket/path/to/delta_table` + * SET TBLPROPERTIES ('delta.enableChangeDataFeed' = 'true'); + * }</pre> + * <li><b>Target GCP Lakehouse (Iceberg V2) Table:</b> Must be an Iceberg format-version 2 table + * registered in BigLake Metastore REST Catalog. Its schema must match the source Delta Lake + * table's columns. If the Iceberg table does not declare identifier (primary-key) fields in + * its table schema, pass {@code --equalityColumns=<col1,col2>} to specify the primary-key + * columns used for equality deletes. + * <li><b>GCP Authentication:</b> Set {@code GOOGLE_APPLICATION_CREDENTIALS} to a service account + * key JSON file (or configure Application Default Credentials) with permissions to access the + * GCS buckets, BigLake Metastore ({@code roles/biglake.admin}), and Dataflow. + * </ul> + * + * <h2>Running the Example on Dataflow Portable Runner (Example Default)</h2> + * + * <p>This example configures the <b>Dataflow Portable Runner</b> ({@code + * --experiments=use_runner_v2}) by default ({@code --usePortableRunner=true}). Because the Dataflow + * Portable Runner does not yet propagate an element's native Beam {@code ValueKind} metadata across + * the FnAPI boundary, the source includes {@code _change_type} in {@code include_metadata_columns} + * and the Iceberg CDC sink is configured with {@code change_type_column: "_change_type"} and {@code + * change_type_map}. + * + * <h3>1. Reading by Commit Version Range</h3> + * + * <pre>{@code + * export GOOGLE_APPLICATION_CREDENTIALS=/path/to/service-account-key.json + * + * ./gradlew :examples:java:execute \ + * -PmainClass=org.apache.beam.examples.DeltaLakeToLakehouseCdcExample \ + * -Pexec.args="--runner=DataflowRunner \ + * --project=my-gcp-project \ + * --region=us-central1 \ + * --tempLocation=gs://my-temp-bucket/temp \ + * --deltaTable=gs://my-delta-lake-bucket/delta_lake/employee_data/ \ + * --lakehouseTable=my-gcp-project.my-lakehouse-warehouse-bucket.my_namespace.my_iceberg_table \ + * --equalityColumns=employee_id \ + * --startVersion=2 \ + * --endVersion=4" + * }</pre> + * + * <h3>2. Reading by Timestamp Range</h3> + * + * <pre>{@code + * export GOOGLE_APPLICATION_CREDENTIALS=/path/to/service-account-key.json + * + * ./gradlew :examples:java:execute \ + * -PmainClass=org.apache.beam.examples.DeltaLakeToLakehouseCdcExample \ + * -Pexec.args="--runner=DataflowRunner \ + * --project=my-gcp-project \ + * --region=us-central1 \ + * --tempLocation=gs://my-temp-bucket/temp \ + * --deltaTable=gs://my-delta-lake-bucket/delta_lake/employee_data/ \ + * --lakehouseTable=my-gcp-project.my-lakehouse-warehouse-bucket.my_namespace.my_iceberg_table \ + * --equalityColumns=employee_id \ + * --startTimestamp=2026-10-01T00:00:00Z \ + * --endTimestamp=2026-10-02T00:00:00Z" + * }</pre> + * + * <h2>Running the Example on Dataflow Streaming Java Runner</h2> + * + * <p>When running on the <b>Dataflow Streaming Java Runner</b> (which is Dataflow's default runner + * when {@code --experiments=use_runner_v2} is not enabled), {@code change_type_column} (and {@code + * change_type_map}) <b>can be skipped</b> on the Iceberg CDC sink, and {@code _change_type} does + * not need to be included in {@code include_metadata_columns} on the Delta Lake CDC source. The + * Delta Lake CDC reader ({@code Managed.DELTA_LAKE_CDC}) automatically sets each emitted {@link + * Row}'s native Beam {@code ValueKind} ({@code INSERT}, {@code DELETE}, {@code UPDATE_BEFORE}, + * {@code UPDATE_AFTER}), which the Dataflow Streaming Java Runner preserves and passes directly to + * the Iceberg CDC sink. + * + * <p>In your own pipelines running on the Dataflow Streaming Java Runner, you do not need to set + * any runner-version flag (since it is Dataflow's default when {@code use_runner_v2} is not added), + * and the {@code Managed} configurations only need {@code _commit_version} for ordering: + * + * <pre>{@code + * // Delta Lake CDC read config on Dataflow Streaming Java Runner (no _change_type column needed): + * Map<String, Object> readConfig = ImmutableMap.of( + * "table", deltaTable, + * "start_version", 2L, + * "end_version", 4L, + * "include_metadata_columns", ImmutableList.of("_commit_version"), + * "hadoop_config", hadoopConfig); + * + * // Iceberg CDC write config on Dataflow Streaming Java Runner (change_type_column can be skipped): + * Map<String, Object> writeConfig = ImmutableMap.of( + * "table", tableId, + * "catalog_name", "lakehouse", + * "catalog_properties", catalogProps, + * "mode", "merge-on-read", + * "sequence_number_column", "_commit_version", + * "equality_columns", ImmutableList.of("employee_id")); + * }</pre> + * + * <p>Because this example enables the Dataflow Portable Runner by default ({@code + * --usePortableRunner=true}), you can override it to run on the Dataflow Streaming Java Runner (and + * skip {@code change_type_column}) by passing {@code --usePortableRunner=false}: + * + * <pre>{@code + * export GOOGLE_APPLICATION_CREDENTIALS=/path/to/service-account-key.json + * + * ./gradlew :examples:java:execute \ + * -PmainClass=org.apache.beam.examples.DeltaLakeToLakehouseCdcExample \ + * -Pexec.args="--runner=DataflowRunner \ + * --usePortableRunner=false \ + * --project=my-gcp-project \ + * --region=us-central1 \ + * --tempLocation=gs://my-temp-bucket/temp \ + * --deltaTable=gs://my-delta-lake-bucket/delta_lake/employee_data/ \ + * --lakehouseTable=my-gcp-project.my-lakehouse-warehouse-bucket.my_namespace.my_iceberg_table \ + * --equalityColumns=employee_id \ + * --startVersion=2 \ + * --endVersion=4" + * }</pre> + */ +public class DeltaLakeToLakehouseCdcExample { + + private static final Logger LOG = LoggerFactory.getLogger(DeltaLakeToLakehouseCdcExample.class); + + /** Delta Lake CDC metadata column names produced by {@code Managed.DELTA_LAKE_CDC}. */ + public static final String CHANGE_TYPE_COLUMN = "_change_type"; + + public static final String COMMIT_VERSION_COLUMN = "_commit_version"; + + /** Default BigLake Iceberg REST Catalog endpoint URI. */ + public static final String DEFAULT_BIGLAKE_CATALOG_URI = Review Comment: ```suggestion /** Default Lakehouse Iceberg REST Catalog endpoint URI. */ public static final String DEFAULT_LAKEHOUSE_CATALOG_URI = ``` ########## examples/java/src/main/java/org/apache/beam/examples/DeltaLakeToLakehouseCdcExample.java: ########## @@ -0,0 +1,643 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.beam.examples; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import java.io.BufferedReader; +import java.io.IOException; +import java.io.InputStreamReader; +import java.nio.channels.Channels; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Comparator; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; +import org.apache.beam.sdk.Pipeline; +import org.apache.beam.sdk.PipelineResult; +import org.apache.beam.sdk.extensions.gcp.options.GcpOptions; +import org.apache.beam.sdk.io.FileSystems; +import org.apache.beam.sdk.io.fs.MatchResult; +import org.apache.beam.sdk.managed.Managed; +import org.apache.beam.sdk.options.Default; +import org.apache.beam.sdk.options.Description; +import org.apache.beam.sdk.options.ExperimentalOptions; +import org.apache.beam.sdk.options.PipelineOptionsFactory; +import org.apache.beam.sdk.options.Validation; +import org.apache.beam.sdk.values.PCollection; +import org.apache.beam.sdk.values.Row; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Splitter; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap; +import org.checkerframework.checker.nullness.qual.Nullable; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * An Apache Beam Java example pipeline that reads Change Data Feed (CDC) records from a Delta Lake + * table on Google Cloud Storage (GCS) and applies those row-level changes (`INSERT`, + * `UPDATE_BEFORE`, `UPDATE_AFTER`, `DELETE`) to a Google Cloud Platform (GCP) Lakehouse (BigLake + * Metastore Iceberg REST Catalog) table using Beam's {@link Managed} I/O connectors. + * + * <h2>Overview</h2> + * + * <p>This pipeline performs the following steps: + * + * <ol> + * <li><b>Validates Delta Lake Change Data Feed (CDF):</b> Inspects the input Delta Lake table's + * transaction log ({@code _delta_log/*.json}) to verify that the table property {@code + * delta.enableChangeDataFeed = true} is enabled, failing fast with an error if it is not. + * <li><b>Reads Delta Lake CDC Data:</b> Uses {@code Managed.read(Managed.DELTA_LAKE_CDC)} to read + * change events over either a commit version range ({@code --startVersion} / {@code + * --endVersion}) or an ISO-8601 timestamp range ({@code --startTimestamp} / {@code + * --endTimestamp}). Standard GCS Hadoop filesystem properties and Delta CDC metadata columns + * are configured by default. + * <li><b>Writes CDC Data to GCP Lakehouse:</b> Uses {@code Managed.write(Managed.ICEBERG)} in + * {@code merge-on-read} mode against the BigLake Iceberg REST catalog to apply inserts, + * updates, and deletes ordered by {@code _commit_version}. Standard BigLake REST catalog and + * Iceberg CDC sink properties are configured by default. + * </ol> + * + * <h2>Prerequisites</h2> + * + * <ul> + * <li><b>Java 17+:</b> Both the Delta Lake Kernel API and Iceberg require Java 17 or later. If + * using SDKMAN, switch to Java 17 before running: + * <pre>{@code + * sdk use java 17.0.15-tem + * }</pre> + * <li><b>Source Delta Lake Table on GCS:</b> Must have Change Data Feed enabled: + * <pre>{@code + * ALTER TABLE delta.`gs://my-bucket/path/to/delta_table` + * SET TBLPROPERTIES ('delta.enableChangeDataFeed' = 'true'); + * }</pre> + * <li><b>Target GCP Lakehouse (Iceberg V2) Table:</b> Must be an Iceberg format-version 2 table + * registered in BigLake Metastore REST Catalog. Its schema must match the source Delta Lake + * table's columns. If the Iceberg table does not declare identifier (primary-key) fields in + * its table schema, pass {@code --equalityColumns=<col1,col2>} to specify the primary-key + * columns used for equality deletes. + * <li><b>GCP Authentication:</b> Set {@code GOOGLE_APPLICATION_CREDENTIALS} to a service account + * key JSON file (or configure Application Default Credentials) with permissions to access the + * GCS buckets, BigLake Metastore ({@code roles/biglake.admin}), and Dataflow. + * </ul> + * + * <h2>Running the Example on Dataflow Portable Runner (Example Default)</h2> + * + * <p>This example configures the <b>Dataflow Portable Runner</b> ({@code + * --experiments=use_runner_v2}) by default ({@code --usePortableRunner=true}). Because the Dataflow + * Portable Runner does not yet propagate an element's native Beam {@code ValueKind} metadata across + * the FnAPI boundary, the source includes {@code _change_type} in {@code include_metadata_columns} + * and the Iceberg CDC sink is configured with {@code change_type_column: "_change_type"} and {@code + * change_type_map}. + * + * <h3>1. Reading by Commit Version Range</h3> + * + * <pre>{@code + * export GOOGLE_APPLICATION_CREDENTIALS=/path/to/service-account-key.json + * + * ./gradlew :examples:java:execute \ + * -PmainClass=org.apache.beam.examples.DeltaLakeToLakehouseCdcExample \ + * -Pexec.args="--runner=DataflowRunner \ + * --project=my-gcp-project \ + * --region=us-central1 \ + * --tempLocation=gs://my-temp-bucket/temp \ + * --deltaTable=gs://my-delta-lake-bucket/delta_lake/employee_data/ \ + * --lakehouseTable=my-gcp-project.my-lakehouse-warehouse-bucket.my_namespace.my_iceberg_table \ + * --equalityColumns=employee_id \ + * --startVersion=2 \ + * --endVersion=4" + * }</pre> + * + * <h3>2. Reading by Timestamp Range</h3> + * + * <pre>{@code + * export GOOGLE_APPLICATION_CREDENTIALS=/path/to/service-account-key.json + * + * ./gradlew :examples:java:execute \ + * -PmainClass=org.apache.beam.examples.DeltaLakeToLakehouseCdcExample \ + * -Pexec.args="--runner=DataflowRunner \ + * --project=my-gcp-project \ + * --region=us-central1 \ + * --tempLocation=gs://my-temp-bucket/temp \ + * --deltaTable=gs://my-delta-lake-bucket/delta_lake/employee_data/ \ + * --lakehouseTable=my-gcp-project.my-lakehouse-warehouse-bucket.my_namespace.my_iceberg_table \ + * --equalityColumns=employee_id \ + * --startTimestamp=2026-10-01T00:00:00Z \ + * --endTimestamp=2026-10-02T00:00:00Z" + * }</pre> + * + * <h2>Running the Example on Dataflow Streaming Java Runner</h2> + * + * <p>When running on the <b>Dataflow Streaming Java Runner</b> (which is Dataflow's default runner + * when {@code --experiments=use_runner_v2} is not enabled), {@code change_type_column} (and {@code + * change_type_map}) <b>can be skipped</b> on the Iceberg CDC sink, and {@code _change_type} does + * not need to be included in {@code include_metadata_columns} on the Delta Lake CDC source. The + * Delta Lake CDC reader ({@code Managed.DELTA_LAKE_CDC}) automatically sets each emitted {@link + * Row}'s native Beam {@code ValueKind} ({@code INSERT}, {@code DELETE}, {@code UPDATE_BEFORE}, + * {@code UPDATE_AFTER}), which the Dataflow Streaming Java Runner preserves and passes directly to + * the Iceberg CDC sink. + * + * <p>In your own pipelines running on the Dataflow Streaming Java Runner, you do not need to set + * any runner-version flag (since it is Dataflow's default when {@code use_runner_v2} is not added), + * and the {@code Managed} configurations only need {@code _commit_version} for ordering: + * + * <pre>{@code + * // Delta Lake CDC read config on Dataflow Streaming Java Runner (no _change_type column needed): + * Map<String, Object> readConfig = ImmutableMap.of( + * "table", deltaTable, + * "start_version", 2L, + * "end_version", 4L, + * "include_metadata_columns", ImmutableList.of("_commit_version"), + * "hadoop_config", hadoopConfig); + * + * // Iceberg CDC write config on Dataflow Streaming Java Runner (change_type_column can be skipped): + * Map<String, Object> writeConfig = ImmutableMap.of( + * "table", tableId, + * "catalog_name", "lakehouse", + * "catalog_properties", catalogProps, + * "mode", "merge-on-read", + * "sequence_number_column", "_commit_version", + * "equality_columns", ImmutableList.of("employee_id")); + * }</pre> + * + * <p>Because this example enables the Dataflow Portable Runner by default ({@code + * --usePortableRunner=true}), you can override it to run on the Dataflow Streaming Java Runner (and + * skip {@code change_type_column}) by passing {@code --usePortableRunner=false}: + * + * <pre>{@code + * export GOOGLE_APPLICATION_CREDENTIALS=/path/to/service-account-key.json + * + * ./gradlew :examples:java:execute \ + * -PmainClass=org.apache.beam.examples.DeltaLakeToLakehouseCdcExample \ + * -Pexec.args="--runner=DataflowRunner \ + * --usePortableRunner=false \ + * --project=my-gcp-project \ + * --region=us-central1 \ + * --tempLocation=gs://my-temp-bucket/temp \ + * --deltaTable=gs://my-delta-lake-bucket/delta_lake/employee_data/ \ + * --lakehouseTable=my-gcp-project.my-lakehouse-warehouse-bucket.my_namespace.my_iceberg_table \ + * --equalityColumns=employee_id \ + * --startVersion=2 \ + * --endVersion=4" + * }</pre> + */ +public class DeltaLakeToLakehouseCdcExample { + + private static final Logger LOG = LoggerFactory.getLogger(DeltaLakeToLakehouseCdcExample.class); + + /** Delta Lake CDC metadata column names produced by {@code Managed.DELTA_LAKE_CDC}. */ + public static final String CHANGE_TYPE_COLUMN = "_change_type"; + + public static final String COMMIT_VERSION_COLUMN = "_commit_version"; + + /** Default BigLake Iceberg REST Catalog endpoint URI. */ + public static final String DEFAULT_BIGLAKE_CATALOG_URI = + "https://biglake.googleapis.com/iceberg/v1/restcatalog"; + + /** + * Mapping from Delta Lake Change Data Feed {@code _change_type} values to the canonical change + * types expected by the Iceberg CDC sink. + */ + public static final Map<String, String> DELTA_TO_ICEBERG_CHANGE_TYPES = + ImmutableMap.of( + "insert", "INSERT", + "delete", "DELETE", + "update_preimage", "UPDATE_BEFORE", + "update_postimage", "UPDATE_AFTER"); + + /** Pipeline options for {@link DeltaLakeToLakehouseCdcExample}. */ + public interface Options extends GcpOptions { + + @Description( + "GCS path of the source Delta Lake table repository to read CDC data from " + + "(e.g. gs://my-bucket/delta_lake/my_table/).") + @Validation.Required + String getDeltaTable(); + + void setDeltaTable(String value); + + @Description( + "Target GCP Lakehouse Iceberg table identifier. Accepts either a 4-part BigLake table " Review Comment: ```suggestion "Target GCP Lakehouse Iceberg table identifier. Accepts either a 4-part Lakehouse table " ``` ########## examples/java/build.gradle: ########## @@ -111,6 +112,44 @@ dependencies { if (project.hasProperty("runnerDependency")) { runtimeOnly project(path: project.getProperty("runnerDependency")) } + + def javaVer = (project.findProperty('testJavaVersion') ?: JavaVersion.current().majorVersion) as int Review Comment: Maybe move this to `examples/java/iceberg` instead? ########## examples/java/src/main/java/org/apache/beam/examples/DeltaLakeToLakehouseCdcExample.java: ########## @@ -0,0 +1,643 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.beam.examples; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import java.io.BufferedReader; +import java.io.IOException; +import java.io.InputStreamReader; +import java.nio.channels.Channels; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Comparator; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; +import org.apache.beam.sdk.Pipeline; +import org.apache.beam.sdk.PipelineResult; +import org.apache.beam.sdk.extensions.gcp.options.GcpOptions; +import org.apache.beam.sdk.io.FileSystems; +import org.apache.beam.sdk.io.fs.MatchResult; +import org.apache.beam.sdk.managed.Managed; +import org.apache.beam.sdk.options.Default; +import org.apache.beam.sdk.options.Description; +import org.apache.beam.sdk.options.ExperimentalOptions; +import org.apache.beam.sdk.options.PipelineOptionsFactory; +import org.apache.beam.sdk.options.Validation; +import org.apache.beam.sdk.values.PCollection; +import org.apache.beam.sdk.values.Row; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Splitter; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap; +import org.checkerframework.checker.nullness.qual.Nullable; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * An Apache Beam Java example pipeline that reads Change Data Feed (CDC) records from a Delta Lake + * table on Google Cloud Storage (GCS) and applies those row-level changes (`INSERT`, + * `UPDATE_BEFORE`, `UPDATE_AFTER`, `DELETE`) to a Google Cloud Platform (GCP) Lakehouse (BigLake + * Metastore Iceberg REST Catalog) table using Beam's {@link Managed} I/O connectors. + * + * <h2>Overview</h2> + * + * <p>This pipeline performs the following steps: + * + * <ol> + * <li><b>Validates Delta Lake Change Data Feed (CDF):</b> Inspects the input Delta Lake table's + * transaction log ({@code _delta_log/*.json}) to verify that the table property {@code + * delta.enableChangeDataFeed = true} is enabled, failing fast with an error if it is not. + * <li><b>Reads Delta Lake CDC Data:</b> Uses {@code Managed.read(Managed.DELTA_LAKE_CDC)} to read + * change events over either a commit version range ({@code --startVersion} / {@code + * --endVersion}) or an ISO-8601 timestamp range ({@code --startTimestamp} / {@code + * --endTimestamp}). Standard GCS Hadoop filesystem properties and Delta CDC metadata columns + * are configured by default. + * <li><b>Writes CDC Data to GCP Lakehouse:</b> Uses {@code Managed.write(Managed.ICEBERG)} in + * {@code merge-on-read} mode against the BigLake Iceberg REST catalog to apply inserts, + * updates, and deletes ordered by {@code _commit_version}. Standard BigLake REST catalog and + * Iceberg CDC sink properties are configured by default. + * </ol> + * + * <h2>Prerequisites</h2> + * + * <ul> + * <li><b>Java 17+:</b> Both the Delta Lake Kernel API and Iceberg require Java 17 or later. If + * using SDKMAN, switch to Java 17 before running: + * <pre>{@code + * sdk use java 17.0.15-tem + * }</pre> + * <li><b>Source Delta Lake Table on GCS:</b> Must have Change Data Feed enabled: + * <pre>{@code + * ALTER TABLE delta.`gs://my-bucket/path/to/delta_table` + * SET TBLPROPERTIES ('delta.enableChangeDataFeed' = 'true'); + * }</pre> + * <li><b>Target GCP Lakehouse (Iceberg V2) Table:</b> Must be an Iceberg format-version 2 table + * registered in BigLake Metastore REST Catalog. Its schema must match the source Delta Lake + * table's columns. If the Iceberg table does not declare identifier (primary-key) fields in + * its table schema, pass {@code --equalityColumns=<col1,col2>} to specify the primary-key + * columns used for equality deletes. + * <li><b>GCP Authentication:</b> Set {@code GOOGLE_APPLICATION_CREDENTIALS} to a service account + * key JSON file (or configure Application Default Credentials) with permissions to access the + * GCS buckets, BigLake Metastore ({@code roles/biglake.admin}), and Dataflow. + * </ul> + * + * <h2>Running the Example on Dataflow Portable Runner (Example Default)</h2> + * + * <p>This example configures the <b>Dataflow Portable Runner</b> ({@code + * --experiments=use_runner_v2}) by default ({@code --usePortableRunner=true}). Because the Dataflow + * Portable Runner does not yet propagate an element's native Beam {@code ValueKind} metadata across + * the FnAPI boundary, the source includes {@code _change_type} in {@code include_metadata_columns} + * and the Iceberg CDC sink is configured with {@code change_type_column: "_change_type"} and {@code + * change_type_map}. + * + * <h3>1. Reading by Commit Version Range</h3> + * + * <pre>{@code + * export GOOGLE_APPLICATION_CREDENTIALS=/path/to/service-account-key.json + * + * ./gradlew :examples:java:execute \ + * -PmainClass=org.apache.beam.examples.DeltaLakeToLakehouseCdcExample \ + * -Pexec.args="--runner=DataflowRunner \ + * --project=my-gcp-project \ + * --region=us-central1 \ + * --tempLocation=gs://my-temp-bucket/temp \ + * --deltaTable=gs://my-delta-lake-bucket/delta_lake/employee_data/ \ + * --lakehouseTable=my-gcp-project.my-lakehouse-warehouse-bucket.my_namespace.my_iceberg_table \ + * --equalityColumns=employee_id \ + * --startVersion=2 \ + * --endVersion=4" + * }</pre> + * + * <h3>2. Reading by Timestamp Range</h3> + * + * <pre>{@code + * export GOOGLE_APPLICATION_CREDENTIALS=/path/to/service-account-key.json + * + * ./gradlew :examples:java:execute \ + * -PmainClass=org.apache.beam.examples.DeltaLakeToLakehouseCdcExample \ + * -Pexec.args="--runner=DataflowRunner \ + * --project=my-gcp-project \ + * --region=us-central1 \ + * --tempLocation=gs://my-temp-bucket/temp \ + * --deltaTable=gs://my-delta-lake-bucket/delta_lake/employee_data/ \ + * --lakehouseTable=my-gcp-project.my-lakehouse-warehouse-bucket.my_namespace.my_iceberg_table \ + * --equalityColumns=employee_id \ + * --startTimestamp=2026-10-01T00:00:00Z \ + * --endTimestamp=2026-10-02T00:00:00Z" + * }</pre> + * + * <h2>Running the Example on Dataflow Streaming Java Runner</h2> + * + * <p>When running on the <b>Dataflow Streaming Java Runner</b> (which is Dataflow's default runner + * when {@code --experiments=use_runner_v2} is not enabled), {@code change_type_column} (and {@code + * change_type_map}) <b>can be skipped</b> on the Iceberg CDC sink, and {@code _change_type} does + * not need to be included in {@code include_metadata_columns} on the Delta Lake CDC source. The + * Delta Lake CDC reader ({@code Managed.DELTA_LAKE_CDC}) automatically sets each emitted {@link + * Row}'s native Beam {@code ValueKind} ({@code INSERT}, {@code DELETE}, {@code UPDATE_BEFORE}, + * {@code UPDATE_AFTER}), which the Dataflow Streaming Java Runner preserves and passes directly to + * the Iceberg CDC sink. + * + * <p>In your own pipelines running on the Dataflow Streaming Java Runner, you do not need to set + * any runner-version flag (since it is Dataflow's default when {@code use_runner_v2} is not added), + * and the {@code Managed} configurations only need {@code _commit_version} for ordering: + * + * <pre>{@code + * // Delta Lake CDC read config on Dataflow Streaming Java Runner (no _change_type column needed): + * Map<String, Object> readConfig = ImmutableMap.of( + * "table", deltaTable, + * "start_version", 2L, + * "end_version", 4L, + * "include_metadata_columns", ImmutableList.of("_commit_version"), + * "hadoop_config", hadoopConfig); + * + * // Iceberg CDC write config on Dataflow Streaming Java Runner (change_type_column can be skipped): + * Map<String, Object> writeConfig = ImmutableMap.of( + * "table", tableId, + * "catalog_name", "lakehouse", + * "catalog_properties", catalogProps, + * "mode", "merge-on-read", + * "sequence_number_column", "_commit_version", + * "equality_columns", ImmutableList.of("employee_id")); + * }</pre> + * + * <p>Because this example enables the Dataflow Portable Runner by default ({@code + * --usePortableRunner=true}), you can override it to run on the Dataflow Streaming Java Runner (and + * skip {@code change_type_column}) by passing {@code --usePortableRunner=false}: + * + * <pre>{@code + * export GOOGLE_APPLICATION_CREDENTIALS=/path/to/service-account-key.json + * + * ./gradlew :examples:java:execute \ + * -PmainClass=org.apache.beam.examples.DeltaLakeToLakehouseCdcExample \ + * -Pexec.args="--runner=DataflowRunner \ + * --usePortableRunner=false \ + * --project=my-gcp-project \ + * --region=us-central1 \ + * --tempLocation=gs://my-temp-bucket/temp \ + * --deltaTable=gs://my-delta-lake-bucket/delta_lake/employee_data/ \ + * --lakehouseTable=my-gcp-project.my-lakehouse-warehouse-bucket.my_namespace.my_iceberg_table \ + * --equalityColumns=employee_id \ + * --startVersion=2 \ + * --endVersion=4" + * }</pre> + */ +public class DeltaLakeToLakehouseCdcExample { + + private static final Logger LOG = LoggerFactory.getLogger(DeltaLakeToLakehouseCdcExample.class); + + /** Delta Lake CDC metadata column names produced by {@code Managed.DELTA_LAKE_CDC}. */ + public static final String CHANGE_TYPE_COLUMN = "_change_type"; + + public static final String COMMIT_VERSION_COLUMN = "_commit_version"; + + /** Default BigLake Iceberg REST Catalog endpoint URI. */ + public static final String DEFAULT_BIGLAKE_CATALOG_URI = + "https://biglake.googleapis.com/iceberg/v1/restcatalog"; + + /** + * Mapping from Delta Lake Change Data Feed {@code _change_type} values to the canonical change + * types expected by the Iceberg CDC sink. + */ + public static final Map<String, String> DELTA_TO_ICEBERG_CHANGE_TYPES = + ImmutableMap.of( + "insert", "INSERT", + "delete", "DELETE", + "update_preimage", "UPDATE_BEFORE", + "update_postimage", "UPDATE_AFTER"); + + /** Pipeline options for {@link DeltaLakeToLakehouseCdcExample}. */ + public interface Options extends GcpOptions { + + @Description( + "GCS path of the source Delta Lake table repository to read CDC data from " + + "(e.g. gs://my-bucket/delta_lake/my_table/).") + @Validation.Required + String getDeltaTable(); + + void setDeltaTable(String value); + + @Description( + "Target GCP Lakehouse Iceberg table identifier. Accepts either a 4-part BigLake table " + + "identifier '<project>.<warehouse_bucket>.<namespace>.<table>' " + + "(e.g. 'my-gcp-project.my-lakehouse-warehouse-bucket.my_namespace.my_iceberg_table'), " + + "a 3-part identifier '<warehouse_bucket>.<namespace>.<table>', or a 2-part Iceberg " + + "table identifier '<namespace>.<table>' (when --warehouse is also specified).") + @Validation.Required + String getLakehouseTable(); + + void setLakehouseTable(String value); + + @Description( + "Starting Delta Lake commit version (inclusive) to read changes from. " + + "Either --startVersion or --startTimestamp must be provided.") + @Nullable Long getStartVersion(); + + void setStartVersion(@Nullable Long value); + + @Description( + "Ending Delta Lake commit version (inclusive) to read changes up to. Optional; defaults " + + "to the latest commit version if omitted.") + @Nullable Long getEndVersion(); + + void setEndVersion(@Nullable Long value); + + @Description( + "Starting timestamp in ISO-8601 format (e.g. '2026-10-01T00:00:00Z') to read Delta Lake " + + "changes from. Either --startVersion or --startTimestamp must be provided.") + @Nullable String getStartTimestamp(); + + void setStartTimestamp(@Nullable String value); + + @Description( + "Ending timestamp in ISO-8601 format (e.g. '2026-10-01T23:59:59Z') to read Delta Lake " + + "changes up to. Optional.") + @Nullable String getEndTimestamp(); + + void setEndTimestamp(@Nullable String value); + + @Description( + "GCS warehouse location for the GCP Lakehouse BigLake catalog (e.g. 'gs://my-warehouse-bucket'). " Review Comment: ```suggestion "GCS warehouse location for the GCP Lakehouse catalog (e.g. 'gs://my-warehouse-bucket'). " ``` ########## examples/java/src/main/java/org/apache/beam/examples/DeltaLakeToLakehouseCdcExample.java: ########## @@ -0,0 +1,643 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.beam.examples; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import java.io.BufferedReader; +import java.io.IOException; +import java.io.InputStreamReader; +import java.nio.channels.Channels; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Comparator; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; +import org.apache.beam.sdk.Pipeline; +import org.apache.beam.sdk.PipelineResult; +import org.apache.beam.sdk.extensions.gcp.options.GcpOptions; +import org.apache.beam.sdk.io.FileSystems; +import org.apache.beam.sdk.io.fs.MatchResult; +import org.apache.beam.sdk.managed.Managed; +import org.apache.beam.sdk.options.Default; +import org.apache.beam.sdk.options.Description; +import org.apache.beam.sdk.options.ExperimentalOptions; +import org.apache.beam.sdk.options.PipelineOptionsFactory; +import org.apache.beam.sdk.options.Validation; +import org.apache.beam.sdk.values.PCollection; +import org.apache.beam.sdk.values.Row; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Splitter; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap; +import org.checkerframework.checker.nullness.qual.Nullable; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * An Apache Beam Java example pipeline that reads Change Data Feed (CDC) records from a Delta Lake + * table on Google Cloud Storage (GCS) and applies those row-level changes (`INSERT`, + * `UPDATE_BEFORE`, `UPDATE_AFTER`, `DELETE`) to a Google Cloud Platform (GCP) Lakehouse (BigLake + * Metastore Iceberg REST Catalog) table using Beam's {@link Managed} I/O connectors. + * + * <h2>Overview</h2> + * + * <p>This pipeline performs the following steps: + * + * <ol> + * <li><b>Validates Delta Lake Change Data Feed (CDF):</b> Inspects the input Delta Lake table's + * transaction log ({@code _delta_log/*.json}) to verify that the table property {@code + * delta.enableChangeDataFeed = true} is enabled, failing fast with an error if it is not. + * <li><b>Reads Delta Lake CDC Data:</b> Uses {@code Managed.read(Managed.DELTA_LAKE_CDC)} to read + * change events over either a commit version range ({@code --startVersion} / {@code + * --endVersion}) or an ISO-8601 timestamp range ({@code --startTimestamp} / {@code + * --endTimestamp}). Standard GCS Hadoop filesystem properties and Delta CDC metadata columns + * are configured by default. + * <li><b>Writes CDC Data to GCP Lakehouse:</b> Uses {@code Managed.write(Managed.ICEBERG)} in + * {@code merge-on-read} mode against the BigLake Iceberg REST catalog to apply inserts, + * updates, and deletes ordered by {@code _commit_version}. Standard BigLake REST catalog and + * Iceberg CDC sink properties are configured by default. + * </ol> + * + * <h2>Prerequisites</h2> + * + * <ul> + * <li><b>Java 17+:</b> Both the Delta Lake Kernel API and Iceberg require Java 17 or later. If + * using SDKMAN, switch to Java 17 before running: + * <pre>{@code + * sdk use java 17.0.15-tem + * }</pre> + * <li><b>Source Delta Lake Table on GCS:</b> Must have Change Data Feed enabled: + * <pre>{@code + * ALTER TABLE delta.`gs://my-bucket/path/to/delta_table` + * SET TBLPROPERTIES ('delta.enableChangeDataFeed' = 'true'); + * }</pre> + * <li><b>Target GCP Lakehouse (Iceberg V2) Table:</b> Must be an Iceberg format-version 2 table + * registered in BigLake Metastore REST Catalog. Its schema must match the source Delta Lake + * table's columns. If the Iceberg table does not declare identifier (primary-key) fields in + * its table schema, pass {@code --equalityColumns=<col1,col2>} to specify the primary-key + * columns used for equality deletes. + * <li><b>GCP Authentication:</b> Set {@code GOOGLE_APPLICATION_CREDENTIALS} to a service account + * key JSON file (or configure Application Default Credentials) with permissions to access the + * GCS buckets, BigLake Metastore ({@code roles/biglake.admin}), and Dataflow. + * </ul> + * + * <h2>Running the Example on Dataflow Portable Runner (Example Default)</h2> + * + * <p>This example configures the <b>Dataflow Portable Runner</b> ({@code + * --experiments=use_runner_v2}) by default ({@code --usePortableRunner=true}). Because the Dataflow + * Portable Runner does not yet propagate an element's native Beam {@code ValueKind} metadata across + * the FnAPI boundary, the source includes {@code _change_type} in {@code include_metadata_columns} + * and the Iceberg CDC sink is configured with {@code change_type_column: "_change_type"} and {@code + * change_type_map}. + * + * <h3>1. Reading by Commit Version Range</h3> + * + * <pre>{@code + * export GOOGLE_APPLICATION_CREDENTIALS=/path/to/service-account-key.json + * + * ./gradlew :examples:java:execute \ + * -PmainClass=org.apache.beam.examples.DeltaLakeToLakehouseCdcExample \ + * -Pexec.args="--runner=DataflowRunner \ + * --project=my-gcp-project \ + * --region=us-central1 \ + * --tempLocation=gs://my-temp-bucket/temp \ + * --deltaTable=gs://my-delta-lake-bucket/delta_lake/employee_data/ \ + * --lakehouseTable=my-gcp-project.my-lakehouse-warehouse-bucket.my_namespace.my_iceberg_table \ + * --equalityColumns=employee_id \ + * --startVersion=2 \ + * --endVersion=4" + * }</pre> + * + * <h3>2. Reading by Timestamp Range</h3> + * + * <pre>{@code + * export GOOGLE_APPLICATION_CREDENTIALS=/path/to/service-account-key.json + * + * ./gradlew :examples:java:execute \ + * -PmainClass=org.apache.beam.examples.DeltaLakeToLakehouseCdcExample \ + * -Pexec.args="--runner=DataflowRunner \ + * --project=my-gcp-project \ + * --region=us-central1 \ + * --tempLocation=gs://my-temp-bucket/temp \ + * --deltaTable=gs://my-delta-lake-bucket/delta_lake/employee_data/ \ + * --lakehouseTable=my-gcp-project.my-lakehouse-warehouse-bucket.my_namespace.my_iceberg_table \ + * --equalityColumns=employee_id \ + * --startTimestamp=2026-10-01T00:00:00Z \ + * --endTimestamp=2026-10-02T00:00:00Z" + * }</pre> + * + * <h2>Running the Example on Dataflow Streaming Java Runner</h2> + * + * <p>When running on the <b>Dataflow Streaming Java Runner</b> (which is Dataflow's default runner + * when {@code --experiments=use_runner_v2} is not enabled), {@code change_type_column} (and {@code + * change_type_map}) <b>can be skipped</b> on the Iceberg CDC sink, and {@code _change_type} does + * not need to be included in {@code include_metadata_columns} on the Delta Lake CDC source. The + * Delta Lake CDC reader ({@code Managed.DELTA_LAKE_CDC}) automatically sets each emitted {@link + * Row}'s native Beam {@code ValueKind} ({@code INSERT}, {@code DELETE}, {@code UPDATE_BEFORE}, + * {@code UPDATE_AFTER}), which the Dataflow Streaming Java Runner preserves and passes directly to + * the Iceberg CDC sink. + * + * <p>In your own pipelines running on the Dataflow Streaming Java Runner, you do not need to set + * any runner-version flag (since it is Dataflow's default when {@code use_runner_v2} is not added), + * and the {@code Managed} configurations only need {@code _commit_version} for ordering: + * + * <pre>{@code + * // Delta Lake CDC read config on Dataflow Streaming Java Runner (no _change_type column needed): + * Map<String, Object> readConfig = ImmutableMap.of( + * "table", deltaTable, + * "start_version", 2L, + * "end_version", 4L, + * "include_metadata_columns", ImmutableList.of("_commit_version"), + * "hadoop_config", hadoopConfig); + * + * // Iceberg CDC write config on Dataflow Streaming Java Runner (change_type_column can be skipped): + * Map<String, Object> writeConfig = ImmutableMap.of( + * "table", tableId, + * "catalog_name", "lakehouse", + * "catalog_properties", catalogProps, + * "mode", "merge-on-read", + * "sequence_number_column", "_commit_version", + * "equality_columns", ImmutableList.of("employee_id")); + * }</pre> + * + * <p>Because this example enables the Dataflow Portable Runner by default ({@code + * --usePortableRunner=true}), you can override it to run on the Dataflow Streaming Java Runner (and + * skip {@code change_type_column}) by passing {@code --usePortableRunner=false}: + * + * <pre>{@code + * export GOOGLE_APPLICATION_CREDENTIALS=/path/to/service-account-key.json + * + * ./gradlew :examples:java:execute \ + * -PmainClass=org.apache.beam.examples.DeltaLakeToLakehouseCdcExample \ + * -Pexec.args="--runner=DataflowRunner \ + * --usePortableRunner=false \ + * --project=my-gcp-project \ + * --region=us-central1 \ + * --tempLocation=gs://my-temp-bucket/temp \ + * --deltaTable=gs://my-delta-lake-bucket/delta_lake/employee_data/ \ + * --lakehouseTable=my-gcp-project.my-lakehouse-warehouse-bucket.my_namespace.my_iceberg_table \ + * --equalityColumns=employee_id \ + * --startVersion=2 \ + * --endVersion=4" + * }</pre> + */ +public class DeltaLakeToLakehouseCdcExample { + + private static final Logger LOG = LoggerFactory.getLogger(DeltaLakeToLakehouseCdcExample.class); + + /** Delta Lake CDC metadata column names produced by {@code Managed.DELTA_LAKE_CDC}. */ + public static final String CHANGE_TYPE_COLUMN = "_change_type"; + + public static final String COMMIT_VERSION_COLUMN = "_commit_version"; + + /** Default BigLake Iceberg REST Catalog endpoint URI. */ + public static final String DEFAULT_BIGLAKE_CATALOG_URI = + "https://biglake.googleapis.com/iceberg/v1/restcatalog"; + + /** + * Mapping from Delta Lake Change Data Feed {@code _change_type} values to the canonical change + * types expected by the Iceberg CDC sink. + */ + public static final Map<String, String> DELTA_TO_ICEBERG_CHANGE_TYPES = + ImmutableMap.of( + "insert", "INSERT", + "delete", "DELETE", + "update_preimage", "UPDATE_BEFORE", + "update_postimage", "UPDATE_AFTER"); + + /** Pipeline options for {@link DeltaLakeToLakehouseCdcExample}. */ + public interface Options extends GcpOptions { + + @Description( + "GCS path of the source Delta Lake table repository to read CDC data from " + + "(e.g. gs://my-bucket/delta_lake/my_table/).") + @Validation.Required + String getDeltaTable(); + + void setDeltaTable(String value); + + @Description( + "Target GCP Lakehouse Iceberg table identifier. Accepts either a 4-part BigLake table " + + "identifier '<project>.<warehouse_bucket>.<namespace>.<table>' " + + "(e.g. 'my-gcp-project.my-lakehouse-warehouse-bucket.my_namespace.my_iceberg_table'), " + + "a 3-part identifier '<warehouse_bucket>.<namespace>.<table>', or a 2-part Iceberg " + + "table identifier '<namespace>.<table>' (when --warehouse is also specified).") + @Validation.Required + String getLakehouseTable(); + + void setLakehouseTable(String value); + + @Description( + "Starting Delta Lake commit version (inclusive) to read changes from. " + + "Either --startVersion or --startTimestamp must be provided.") + @Nullable Long getStartVersion(); + + void setStartVersion(@Nullable Long value); + + @Description( + "Ending Delta Lake commit version (inclusive) to read changes up to. Optional; defaults " + + "to the latest commit version if omitted.") + @Nullable Long getEndVersion(); + + void setEndVersion(@Nullable Long value); + + @Description( + "Starting timestamp in ISO-8601 format (e.g. '2026-10-01T00:00:00Z') to read Delta Lake " + + "changes from. Either --startVersion or --startTimestamp must be provided.") + @Nullable String getStartTimestamp(); + + void setStartTimestamp(@Nullable String value); + + @Description( + "Ending timestamp in ISO-8601 format (e.g. '2026-10-01T23:59:59Z') to read Delta Lake " + + "changes up to. Optional.") + @Nullable String getEndTimestamp(); + + void setEndTimestamp(@Nullable String value); + + @Description( + "GCS warehouse location for the GCP Lakehouse BigLake catalog (e.g. 'gs://my-warehouse-bucket'). " + + "Optional when --lakehouseTable is specified as a 3-part or 4-part identifier containing " + + "the warehouse bucket name.") + @Nullable String getWarehouse(); + + void setWarehouse(@Nullable String value); + + @Description( + "Comma-separated list of primary-key (equality-delete) column names that uniquely identify " + + "rows in the target Lakehouse Iceberg table (e.g. 'employee_id'). Required if the " + + "target Iceberg table does not declare identifier-field-ids in its schema or if the " + + "table does not exist yet.") + @Nullable String getEqualityColumns(); + + void setEqualityColumns(@Nullable String value); + + @Description( + "If true, applies changes in upsert mode (UPDATE_BEFORE records are dropped and " + + "INSERT/UPDATE_AFTER records are applied as upserts). Defaults to false.") + @Default.Boolean(false) + boolean getUpsert(); + + void setUpsert(boolean value); + + @Description( + "If true (default for this example), runs on the Dataflow Portable Runner ('use_runner_v2') " + + "and configures 'change_type_column' ('_change_type') on the Iceberg CDC sink. " + + "Set to false to override and run on the Dataflow Streaming Java Runner, which skips " + + "'change_type_column' and relies on the native Beam ValueKind metadata attached to " + + "each Row by the Delta Lake CDC reader.") + @Default.Boolean(true) + boolean getUsePortableRunner(); + + void setUsePortableRunner(boolean value); + + @Description("Name of the Iceberg catalog instance. Defaults to 'lakehouse'.") + @Default.String("lakehouse") + String getCatalogName(); + + void setCatalogName(String value); + + @Description( + "BigLake Iceberg REST Catalog endpoint URI. Defaults to " Review Comment: ```suggestion "Lakehouse Iceberg REST Catalog endpoint URI. Defaults to " ``` ########## examples/java/src/main/java/org/apache/beam/examples/DeltaLakeToLakehouseCdcExample.java: ########## @@ -0,0 +1,643 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.beam.examples; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import java.io.BufferedReader; +import java.io.IOException; +import java.io.InputStreamReader; +import java.nio.channels.Channels; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Comparator; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; +import org.apache.beam.sdk.Pipeline; +import org.apache.beam.sdk.PipelineResult; +import org.apache.beam.sdk.extensions.gcp.options.GcpOptions; +import org.apache.beam.sdk.io.FileSystems; +import org.apache.beam.sdk.io.fs.MatchResult; +import org.apache.beam.sdk.managed.Managed; +import org.apache.beam.sdk.options.Default; +import org.apache.beam.sdk.options.Description; +import org.apache.beam.sdk.options.ExperimentalOptions; +import org.apache.beam.sdk.options.PipelineOptionsFactory; +import org.apache.beam.sdk.options.Validation; +import org.apache.beam.sdk.values.PCollection; +import org.apache.beam.sdk.values.Row; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Splitter; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap; +import org.checkerframework.checker.nullness.qual.Nullable; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * An Apache Beam Java example pipeline that reads Change Data Feed (CDC) records from a Delta Lake + * table on Google Cloud Storage (GCS) and applies those row-level changes (`INSERT`, + * `UPDATE_BEFORE`, `UPDATE_AFTER`, `DELETE`) to a Google Cloud Platform (GCP) Lakehouse (BigLake + * Metastore Iceberg REST Catalog) table using Beam's {@link Managed} I/O connectors. + * + * <h2>Overview</h2> + * + * <p>This pipeline performs the following steps: + * + * <ol> + * <li><b>Validates Delta Lake Change Data Feed (CDF):</b> Inspects the input Delta Lake table's + * transaction log ({@code _delta_log/*.json}) to verify that the table property {@code + * delta.enableChangeDataFeed = true} is enabled, failing fast with an error if it is not. + * <li><b>Reads Delta Lake CDC Data:</b> Uses {@code Managed.read(Managed.DELTA_LAKE_CDC)} to read + * change events over either a commit version range ({@code --startVersion} / {@code + * --endVersion}) or an ISO-8601 timestamp range ({@code --startTimestamp} / {@code + * --endTimestamp}). Standard GCS Hadoop filesystem properties and Delta CDC metadata columns + * are configured by default. + * <li><b>Writes CDC Data to GCP Lakehouse:</b> Uses {@code Managed.write(Managed.ICEBERG)} in + * {@code merge-on-read} mode against the BigLake Iceberg REST catalog to apply inserts, + * updates, and deletes ordered by {@code _commit_version}. Standard BigLake REST catalog and + * Iceberg CDC sink properties are configured by default. + * </ol> + * + * <h2>Prerequisites</h2> + * + * <ul> + * <li><b>Java 17+:</b> Both the Delta Lake Kernel API and Iceberg require Java 17 or later. If + * using SDKMAN, switch to Java 17 before running: + * <pre>{@code + * sdk use java 17.0.15-tem + * }</pre> + * <li><b>Source Delta Lake Table on GCS:</b> Must have Change Data Feed enabled: + * <pre>{@code + * ALTER TABLE delta.`gs://my-bucket/path/to/delta_table` + * SET TBLPROPERTIES ('delta.enableChangeDataFeed' = 'true'); + * }</pre> + * <li><b>Target GCP Lakehouse (Iceberg V2) Table:</b> Must be an Iceberg format-version 2 table + * registered in BigLake Metastore REST Catalog. Its schema must match the source Delta Lake + * table's columns. If the Iceberg table does not declare identifier (primary-key) fields in + * its table schema, pass {@code --equalityColumns=<col1,col2>} to specify the primary-key + * columns used for equality deletes. + * <li><b>GCP Authentication:</b> Set {@code GOOGLE_APPLICATION_CREDENTIALS} to a service account + * key JSON file (or configure Application Default Credentials) with permissions to access the + * GCS buckets, BigLake Metastore ({@code roles/biglake.admin}), and Dataflow. + * </ul> + * + * <h2>Running the Example on Dataflow Portable Runner (Example Default)</h2> + * + * <p>This example configures the <b>Dataflow Portable Runner</b> ({@code + * --experiments=use_runner_v2}) by default ({@code --usePortableRunner=true}). Because the Dataflow + * Portable Runner does not yet propagate an element's native Beam {@code ValueKind} metadata across + * the FnAPI boundary, the source includes {@code _change_type} in {@code include_metadata_columns} + * and the Iceberg CDC sink is configured with {@code change_type_column: "_change_type"} and {@code + * change_type_map}. + * + * <h3>1. Reading by Commit Version Range</h3> + * + * <pre>{@code + * export GOOGLE_APPLICATION_CREDENTIALS=/path/to/service-account-key.json + * + * ./gradlew :examples:java:execute \ + * -PmainClass=org.apache.beam.examples.DeltaLakeToLakehouseCdcExample \ + * -Pexec.args="--runner=DataflowRunner \ + * --project=my-gcp-project \ + * --region=us-central1 \ + * --tempLocation=gs://my-temp-bucket/temp \ + * --deltaTable=gs://my-delta-lake-bucket/delta_lake/employee_data/ \ + * --lakehouseTable=my-gcp-project.my-lakehouse-warehouse-bucket.my_namespace.my_iceberg_table \ + * --equalityColumns=employee_id \ + * --startVersion=2 \ + * --endVersion=4" + * }</pre> + * + * <h3>2. Reading by Timestamp Range</h3> + * + * <pre>{@code + * export GOOGLE_APPLICATION_CREDENTIALS=/path/to/service-account-key.json + * + * ./gradlew :examples:java:execute \ + * -PmainClass=org.apache.beam.examples.DeltaLakeToLakehouseCdcExample \ + * -Pexec.args="--runner=DataflowRunner \ + * --project=my-gcp-project \ + * --region=us-central1 \ + * --tempLocation=gs://my-temp-bucket/temp \ + * --deltaTable=gs://my-delta-lake-bucket/delta_lake/employee_data/ \ + * --lakehouseTable=my-gcp-project.my-lakehouse-warehouse-bucket.my_namespace.my_iceberg_table \ + * --equalityColumns=employee_id \ + * --startTimestamp=2026-10-01T00:00:00Z \ + * --endTimestamp=2026-10-02T00:00:00Z" + * }</pre> + * + * <h2>Running the Example on Dataflow Streaming Java Runner</h2> + * + * <p>When running on the <b>Dataflow Streaming Java Runner</b> (which is Dataflow's default runner + * when {@code --experiments=use_runner_v2} is not enabled), {@code change_type_column} (and {@code + * change_type_map}) <b>can be skipped</b> on the Iceberg CDC sink, and {@code _change_type} does + * not need to be included in {@code include_metadata_columns} on the Delta Lake CDC source. The + * Delta Lake CDC reader ({@code Managed.DELTA_LAKE_CDC}) automatically sets each emitted {@link + * Row}'s native Beam {@code ValueKind} ({@code INSERT}, {@code DELETE}, {@code UPDATE_BEFORE}, + * {@code UPDATE_AFTER}), which the Dataflow Streaming Java Runner preserves and passes directly to + * the Iceberg CDC sink. + * + * <p>In your own pipelines running on the Dataflow Streaming Java Runner, you do not need to set + * any runner-version flag (since it is Dataflow's default when {@code use_runner_v2} is not added), + * and the {@code Managed} configurations only need {@code _commit_version} for ordering: + * + * <pre>{@code + * // Delta Lake CDC read config on Dataflow Streaming Java Runner (no _change_type column needed): + * Map<String, Object> readConfig = ImmutableMap.of( + * "table", deltaTable, + * "start_version", 2L, + * "end_version", 4L, + * "include_metadata_columns", ImmutableList.of("_commit_version"), + * "hadoop_config", hadoopConfig); + * + * // Iceberg CDC write config on Dataflow Streaming Java Runner (change_type_column can be skipped): + * Map<String, Object> writeConfig = ImmutableMap.of( + * "table", tableId, + * "catalog_name", "lakehouse", + * "catalog_properties", catalogProps, + * "mode", "merge-on-read", + * "sequence_number_column", "_commit_version", + * "equality_columns", ImmutableList.of("employee_id")); + * }</pre> + * + * <p>Because this example enables the Dataflow Portable Runner by default ({@code + * --usePortableRunner=true}), you can override it to run on the Dataflow Streaming Java Runner (and + * skip {@code change_type_column}) by passing {@code --usePortableRunner=false}: + * + * <pre>{@code + * export GOOGLE_APPLICATION_CREDENTIALS=/path/to/service-account-key.json + * + * ./gradlew :examples:java:execute \ + * -PmainClass=org.apache.beam.examples.DeltaLakeToLakehouseCdcExample \ + * -Pexec.args="--runner=DataflowRunner \ + * --usePortableRunner=false \ + * --project=my-gcp-project \ + * --region=us-central1 \ + * --tempLocation=gs://my-temp-bucket/temp \ + * --deltaTable=gs://my-delta-lake-bucket/delta_lake/employee_data/ \ + * --lakehouseTable=my-gcp-project.my-lakehouse-warehouse-bucket.my_namespace.my_iceberg_table \ + * --equalityColumns=employee_id \ + * --startVersion=2 \ + * --endVersion=4" + * }</pre> + */ +public class DeltaLakeToLakehouseCdcExample { + + private static final Logger LOG = LoggerFactory.getLogger(DeltaLakeToLakehouseCdcExample.class); + + /** Delta Lake CDC metadata column names produced by {@code Managed.DELTA_LAKE_CDC}. */ + public static final String CHANGE_TYPE_COLUMN = "_change_type"; + + public static final String COMMIT_VERSION_COLUMN = "_commit_version"; + + /** Default BigLake Iceberg REST Catalog endpoint URI. */ + public static final String DEFAULT_BIGLAKE_CATALOG_URI = + "https://biglake.googleapis.com/iceberg/v1/restcatalog"; + + /** + * Mapping from Delta Lake Change Data Feed {@code _change_type} values to the canonical change + * types expected by the Iceberg CDC sink. + */ + public static final Map<String, String> DELTA_TO_ICEBERG_CHANGE_TYPES = + ImmutableMap.of( + "insert", "INSERT", + "delete", "DELETE", + "update_preimage", "UPDATE_BEFORE", + "update_postimage", "UPDATE_AFTER"); + + /** Pipeline options for {@link DeltaLakeToLakehouseCdcExample}. */ + public interface Options extends GcpOptions { + + @Description( + "GCS path of the source Delta Lake table repository to read CDC data from " + + "(e.g. gs://my-bucket/delta_lake/my_table/).") + @Validation.Required + String getDeltaTable(); + + void setDeltaTable(String value); + + @Description( + "Target GCP Lakehouse Iceberg table identifier. Accepts either a 4-part BigLake table " + + "identifier '<project>.<warehouse_bucket>.<namespace>.<table>' " + + "(e.g. 'my-gcp-project.my-lakehouse-warehouse-bucket.my_namespace.my_iceberg_table'), " + + "a 3-part identifier '<warehouse_bucket>.<namespace>.<table>', or a 2-part Iceberg " + + "table identifier '<namespace>.<table>' (when --warehouse is also specified).") + @Validation.Required + String getLakehouseTable(); + + void setLakehouseTable(String value); + + @Description( + "Starting Delta Lake commit version (inclusive) to read changes from. " + + "Either --startVersion or --startTimestamp must be provided.") + @Nullable Long getStartVersion(); + + void setStartVersion(@Nullable Long value); + + @Description( + "Ending Delta Lake commit version (inclusive) to read changes up to. Optional; defaults " + + "to the latest commit version if omitted.") + @Nullable Long getEndVersion(); + + void setEndVersion(@Nullable Long value); + + @Description( + "Starting timestamp in ISO-8601 format (e.g. '2026-10-01T00:00:00Z') to read Delta Lake " + + "changes from. Either --startVersion or --startTimestamp must be provided.") + @Nullable String getStartTimestamp(); + + void setStartTimestamp(@Nullable String value); + + @Description( + "Ending timestamp in ISO-8601 format (e.g. '2026-10-01T23:59:59Z') to read Delta Lake " + + "changes up to. Optional.") + @Nullable String getEndTimestamp(); + + void setEndTimestamp(@Nullable String value); + + @Description( + "GCS warehouse location for the GCP Lakehouse BigLake catalog (e.g. 'gs://my-warehouse-bucket'). " + + "Optional when --lakehouseTable is specified as a 3-part or 4-part identifier containing " + + "the warehouse bucket name.") + @Nullable String getWarehouse(); + + void setWarehouse(@Nullable String value); + + @Description( + "Comma-separated list of primary-key (equality-delete) column names that uniquely identify " + + "rows in the target Lakehouse Iceberg table (e.g. 'employee_id'). Required if the " + + "target Iceberg table does not declare identifier-field-ids in its schema or if the " + + "table does not exist yet.") + @Nullable String getEqualityColumns(); + + void setEqualityColumns(@Nullable String value); + + @Description( + "If true, applies changes in upsert mode (UPDATE_BEFORE records are dropped and " + + "INSERT/UPDATE_AFTER records are applied as upserts). Defaults to false.") + @Default.Boolean(false) + boolean getUpsert(); + + void setUpsert(boolean value); + + @Description( + "If true (default for this example), runs on the Dataflow Portable Runner ('use_runner_v2') " + + "and configures 'change_type_column' ('_change_type') on the Iceberg CDC sink. " + + "Set to false to override and run on the Dataflow Streaming Java Runner, which skips " + + "'change_type_column' and relies on the native Beam ValueKind metadata attached to " + + "each Row by the Delta Lake CDC reader.") + @Default.Boolean(true) + boolean getUsePortableRunner(); + + void setUsePortableRunner(boolean value); + + @Description("Name of the Iceberg catalog instance. Defaults to 'lakehouse'.") + @Default.String("lakehouse") + String getCatalogName(); + + void setCatalogName(String value); + + @Description( + "BigLake Iceberg REST Catalog endpoint URI. Defaults to " + + DEFAULT_BIGLAKE_CATALOG_URI + + ".") + @Default.String(DEFAULT_BIGLAKE_CATALOG_URI) + String getCatalogUri(); + + void setCatalogUri(String value); + } + + /** Parsed GCP Lakehouse table coordinates (`project`, `warehouse`, and `namespace.table`). */ + @VisibleForTesting + static final class LakehouseTableSpec { + final String project; + final String warehouse; + final String tableId; + + LakehouseTableSpec(String project, String warehouse, String tableId) { + this.project = project; + this.warehouse = warehouse; + this.tableId = tableId; + } + } + + /** + * Resolves the GCP project, GCS warehouse URI, and 2-part Iceberg {@code namespace.table} + * identifier from {@link Options}. + */ + @VisibleForTesting + static LakehouseTableSpec resolveLakehouseTableSpec(Options options) { Review Comment: Can we pull this helper method and the validation methods (validateRangeOptions, verifyChangeDataFeedEnabled) to a separate file? Would help to simplify this example and make it more readable. The `build__Config` methods might be more relevant to keep here since they give insight into what the config should look like -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
