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]

Reply via email to