chamikaramj commented on code in PR #40387: URL: https://github.com/apache/beam/pull/40387#discussion_r4228269353
########## 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: Done. -- 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]
