This is an automated email from the ASF dual-hosted git repository.

chamikaramj pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/master by this push:
     new 86ca05321d5 Adds the Delta Lake CDC read transforms to the Managed I/O 
API (#39599)
86ca05321d5 is described below

commit 86ca05321d51cd8c8719c1f6d8fd10d6c5ee7c0a
Author: Chamikara Jayalath <[email protected]>
AuthorDate: Tue Aug 4 23:42:18 2026 -0700

    Adds the Delta Lake CDC read transforms to the Managed I/O API (#39599)
---
 .../beam_PostCommit_Java_Delta_IO_Dataflow.json    |   2 +-
 .../model/pipeline/v1/external_transforms.proto    |   2 +
 sdks/java/io/delta/build.gradle                    |   2 +-
 .../beam/sdk/io/delta/DeltaCDCSourceDoFn.java      |   8 +-
 .../delta/DeltaCdcReadSchemaTransformProvider.java | 178 ++++
 .../java/org/apache/beam/sdk/io/delta/DeltaIO.java |  48 +-
 .../org/apache/beam/sdk/io/delta/DeltaIOIT.java    | 168 +++-
 .../org/apache/beam/sdk/io/delta/DeltaIOTest.java  | 941 +++++++++++----------
 .../beam/sdk/io/delta/DeltaWriteTestUtils.java     | 371 ++++++++
 .../java/org/apache/beam/sdk/managed/Managed.java  |   5 +
 sdks/standard_expansion_services.yaml              |   1 +
 .../site/content/en/documentation/io/managed-io.md | 104 +++
 12 files changed, 1390 insertions(+), 440 deletions(-)

diff --git a/.github/trigger_files/beam_PostCommit_Java_Delta_IO_Dataflow.json 
b/.github/trigger_files/beam_PostCommit_Java_Delta_IO_Dataflow.json
index 5abe02fc09c..ab4daeae234 100644
--- a/.github/trigger_files/beam_PostCommit_Java_Delta_IO_Dataflow.json
+++ b/.github/trigger_files/beam_PostCommit_Java_Delta_IO_Dataflow.json
@@ -1,4 +1,4 @@
 {
   "comment": "Modify this file in a trivial way to cause this test suite to 
run.",
-  "modification": 1
+  "modification": 3
 }
diff --git 
a/model/pipeline/src/main/proto/org/apache/beam/model/pipeline/v1/external_transforms.proto
 
b/model/pipeline/src/main/proto/org/apache/beam/model/pipeline/v1/external_transforms.proto
index 918455dbdbd..debacc245d6 100644
--- 
a/model/pipeline/src/main/proto/org/apache/beam/model/pipeline/v1/external_transforms.proto
+++ 
b/model/pipeline/src/main/proto/org/apache/beam/model/pipeline/v1/external_transforms.proto
@@ -107,6 +107,8 @@ message ManagedTransforms {
       "beam:schematransform:org.apache.beam:sql_server_write:v1"];
     DELTA_LAKE_READ = 13 [(org.apache.beam.model.pipeline.v1.beam_urn) =
       "beam:schematransform:org.apache.beam:delta_lake_read:v1"];
+    DELTA_LAKE_CDC_READ = 14 [(org.apache.beam.model.pipeline.v1.beam_urn) =
+      "beam:schematransform:org.apache.beam:delta_lake_cdc_read:v1"];
   }
 }
 
diff --git a/sdks/java/io/delta/build.gradle b/sdks/java/io/delta/build.gradle
index 5ee5442ecd1..66bacf547d1 100644
--- a/sdks/java/io/delta/build.gradle
+++ b/sdks/java/io/delta/build.gradle
@@ -101,7 +101,7 @@ task dataflowIntegrationTest(type: Test) {
     def dockerJavaImageName = 
project.project(':runners:google-cloud-dataflow-java').ext.dockerJavaImageName
 
     def args = [
-            "--runner=DataflowRunner",
+            "--runner=TestDataflowRunner",
             "--region=us-central1",
             "--project=${gcpProject}",
             "--tempLocation=${gcpTempLocation}",
diff --git 
a/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaCDCSourceDoFn.java
 
b/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaCDCSourceDoFn.java
index cf10adc9865..414402429c3 100644
--- 
a/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaCDCSourceDoFn.java
+++ 
b/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaCDCSourceDoFn.java
@@ -61,11 +61,14 @@ import org.checkerframework.checker.nullness.qual.Nullable;
 @DoFn.BoundedPerElement
 class DeltaCDCSourceDoFn extends DoFn<DeltaCDCReadTask, Row> {
   @Nullable Map<String, String> hadoopConfig;
+  private final @Nullable List<String> metadataColumns;
   private transient @Nullable Engine engine;
   private transient @Nullable Configuration conf;
 
-  public DeltaCDCSourceDoFn(@Nullable Map<String, String> hadoopConfig) {
+  public DeltaCDCSourceDoFn(
+      @Nullable Map<String, String> hadoopConfig, @Nullable List<String> 
metadataColumns) {
     this.hadoopConfig = hadoopConfig;
+    this.metadataColumns = metadataColumns;
   }
 
   private synchronized Configuration getConfiguration() {
@@ -117,7 +120,8 @@ class DeltaCDCSourceDoFn extends DoFn<DeltaCDCReadTask, 
Row> {
 
     SerializableRow originalScanStateRow = task.getScanStateRow();
     StructType logicalTableSchema = 
ScanStateRow.getLogicalSchema(originalScanStateRow);
-    Schema publicBeamSchema = 
DeltaIO.ReadRows.convertToBeamSchema(logicalTableSchema);
+    Schema baseSchema = 
DeltaIO.ReadRows.convertToBeamSchema(logicalTableSchema);
+    Schema publicBeamSchema = DeltaIO.buildPublicBeamSchema(baseSchema, 
metadataColumns);
     StructType physicalTableSchema = 
ScanStateRow.getPhysicalDataReadSchema(originalScanStateRow);
 
     StructType scanStateSchema = originalScanStateRow.getSchema();
diff --git 
a/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaCdcReadSchemaTransformProvider.java
 
b/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaCdcReadSchemaTransformProvider.java
new file mode 100644
index 00000000000..f35a7a52b05
--- /dev/null
+++ 
b/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaCdcReadSchemaTransformProvider.java
@@ -0,0 +1,178 @@
+/*
+ * 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.sdk.io.delta;
+
+import static 
org.apache.beam.sdk.io.delta.DeltaCdcReadSchemaTransformProvider.Configuration;
+import static org.apache.beam.sdk.util.construction.BeamUrns.getUrn;
+
+import com.google.auto.service.AutoService;
+import com.google.auto.value.AutoValue;
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+import org.apache.beam.model.pipeline.v1.ExternalTransforms;
+import org.apache.beam.sdk.schemas.AutoValueSchema;
+import org.apache.beam.sdk.schemas.NoSuchSchemaException;
+import org.apache.beam.sdk.schemas.SchemaRegistry;
+import org.apache.beam.sdk.schemas.annotations.DefaultSchema;
+import org.apache.beam.sdk.schemas.annotations.SchemaFieldDescription;
+import org.apache.beam.sdk.schemas.transforms.SchemaTransform;
+import org.apache.beam.sdk.schemas.transforms.SchemaTransformProvider;
+import org.apache.beam.sdk.schemas.transforms.TypedSchemaTransformProvider;
+import org.apache.beam.sdk.values.PCollection;
+import org.apache.beam.sdk.values.PCollectionRowTuple;
+import org.apache.beam.sdk.values.Row;
+import org.checkerframework.checker.nullness.qual.Nullable;
+
+/**
+ * SchemaTransform implementation for {@link DeltaIO#readChanges}. Reads 
change records from Delta
+ * Lake and outputs a {@link org.apache.beam.sdk.values.PCollection} of Beam 
{@link
+ * org.apache.beam.sdk.values.Row}s.
+ */
+@AutoService(SchemaTransformProvider.class)
+public class DeltaCdcReadSchemaTransformProvider
+    extends TypedSchemaTransformProvider<Configuration> {
+  static final String OUTPUT_TAG = "output";
+
+  @Override
+  protected SchemaTransform from(Configuration configuration) {
+    return new DeltaCdcReadSchemaTransform(configuration);
+  }
+
+  @Override
+  public List<String> outputCollectionNames() {
+    return Collections.singletonList(OUTPUT_TAG);
+  }
+
+  @Override
+  public String identifier() {
+    return 
getUrn(ExternalTransforms.ManagedTransforms.Urns.DELTA_LAKE_CDC_READ);
+  }
+
+  static class DeltaCdcReadSchemaTransform extends SchemaTransform {
+    private final Configuration configuration;
+
+    DeltaCdcReadSchemaTransform(Configuration configuration) {
+      this.configuration =
+          java.util.Objects.requireNonNull(configuration, "configuration 
cannot be null");
+    }
+
+    Row getConfigurationRow() {
+      try {
+        return SchemaRegistry.createDefault()
+            .getToRowFunction(Configuration.class)
+            .apply(configuration)
+            .sorted()
+            .toSnakeCase();
+      } catch (NoSuchSchemaException e) {
+        throw new RuntimeException(e);
+      }
+    }
+
+    @Override
+    public PCollectionRowTuple expand(PCollectionRowTuple input) {
+      DeltaIO.ReadChanges read = 
DeltaIO.readChanges().from(configuration.getTable());
+      Long startVersion = configuration.getStartVersion();
+      if (startVersion != null) {
+        read = read.withStartVersion(startVersion);
+      }
+      String startTimestamp = configuration.getStartTimestamp();
+      if (startTimestamp != null) {
+        read = read.withStartTimestamp(startTimestamp);
+      }
+      Long endVersion = configuration.getEndVersion();
+      if (endVersion != null) {
+        read = read.withEndVersion(endVersion);
+      }
+      String endTimestamp = configuration.getEndTimestamp();
+      if (endTimestamp != null) {
+        read = read.withEndTimestamp(endTimestamp);
+      }
+      Map<String, String> hadoopConfig = configuration.getHadoopConfig();
+      if (hadoopConfig != null) {
+        read = read.withConfig(hadoopConfig);
+      }
+      List<String> includeMetadataColumns = 
configuration.getIncludeMetadataColumns();
+      if (includeMetadataColumns != null && !includeMetadataColumns.isEmpty()) 
{
+        read = read.withMetadataColumns(includeMetadataColumns.toArray(new 
String[0]));
+      }
+
+      PCollection<Row> output = input.getPipeline().apply(read);
+
+      return PCollectionRowTuple.of(OUTPUT_TAG, output);
+    }
+  }
+
+  @DefaultSchema(AutoValueSchema.class)
+  @AutoValue
+  public abstract static class Configuration {
+    static Builder builder() {
+      return new 
AutoValue_DeltaCdcReadSchemaTransformProvider_Configuration.Builder();
+    }
+
+    @SchemaFieldDescription("Identifier of the Delta Lake table.")
+    abstract String getTable();
+
+    @SchemaFieldDescription(
+        "Start version of the Delta Lake table to read changes from. Either 
this or the start timestamp has to be provided.")
+    @Nullable
+    abstract Long getStartVersion();
+
+    @SchemaFieldDescription(
+        "Start timestamp of the Delta Lake table to read changes from. Should 
be specified in the ISO 8601 standard. Either this or the start version has to 
be provided.")
+    @Nullable
+    abstract String getStartTimestamp();
+
+    @SchemaFieldDescription("End version of the Delta Lake table to read 
changes up to.")
+    @Nullable
+    abstract Long getEndVersion();
+
+    @SchemaFieldDescription(
+        "End timestamp of the Delta Lake table to read changes up to. Should 
be specified in the ISO 8601 standard.")
+    @Nullable
+    abstract String getEndTimestamp();
+
+    @SchemaFieldDescription("Properties passed to the Hadoop Configuration.")
+    @Nullable
+    abstract Map<String, String> getHadoopConfig();
+
+    @SchemaFieldDescription(
+        "Metadata columns to include in the output rows. Supported columns 
are: _change_type, _commit_version, and _commit_timestamp.")
+    @Nullable
+    abstract List<String> getIncludeMetadataColumns();
+
+    @AutoValue.Builder
+    abstract static class Builder {
+      abstract Builder setTable(String table);
+
+      abstract Builder setStartVersion(Long startVersion);
+
+      abstract Builder setStartTimestamp(String startTimestamp);
+
+      abstract Builder setEndVersion(Long endVersion);
+
+      abstract Builder setEndTimestamp(String endTimestamp);
+
+      abstract Builder setHadoopConfig(Map<String, String> hadoopConfig);
+
+      abstract Builder setIncludeMetadataColumns(List<String> 
includeMetadataColumns);
+
+      abstract Configuration build();
+    }
+  }
+}
diff --git 
a/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaIO.java 
b/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaIO.java
index 3ac2c7a84a8..8057332ddce 100644
--- a/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaIO.java
+++ b/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaIO.java
@@ -37,6 +37,8 @@ import io.delta.kernel.types.StringType;
 import io.delta.kernel.types.StructField;
 import io.delta.kernel.types.StructType;
 import io.delta.kernel.types.TimestampType;
+import java.util.Arrays;
+import java.util.List;
 import java.util.Map;
 import org.apache.beam.sdk.annotations.Internal;
 import org.apache.beam.sdk.schemas.Schema;
@@ -206,6 +208,26 @@ public class DeltaIO {
     }
   }
 
+  static Schema buildPublicBeamSchema(Schema baseSchema, @Nullable 
List<String> metadataColumns) {
+    if (metadataColumns == null || metadataColumns.isEmpty()) {
+      return baseSchema;
+    }
+    Schema.Builder builder = Schema.builder();
+    for (Schema.Field field : baseSchema.getFields()) {
+      builder.addField(field);
+    }
+    for (String col : metadataColumns) {
+      if (col.equals(CHANGE_TYPE_COLUMN)) {
+        builder.addField(CHANGE_TYPE_COLUMN, Schema.FieldType.STRING);
+      } else if (col.equals(COMMIT_VERSION_COLUMN)) {
+        builder.addField(COMMIT_VERSION_COLUMN, Schema.FieldType.INT64);
+      } else if (col.equals(COMMIT_TIMESTAMP_COLUMN)) {
+        builder.addField(COMMIT_TIMESTAMP_COLUMN, Schema.FieldType.DATETIME);
+      }
+    }
+    return builder.build();
+  }
+
   @AutoValue
   public abstract static class ReadChanges extends PTransform<PBegin, 
PCollection<Row>> {
     public abstract @Nullable String getTablePath();
@@ -218,6 +240,8 @@ public class DeltaIO {
 
     public abstract @Nullable String getEndTimestamp();
 
+    public abstract @Nullable List<String> getMetadataColumns();
+
     public abstract @Nullable Map<String, String> getHadoopConfig();
 
     abstract Builder toBuilder();
@@ -234,6 +258,8 @@ public class DeltaIO {
 
       abstract Builder setEndTimestamp(@Nullable String endTimestamp);
 
+      abstract Builder setMetadataColumns(@Nullable List<String> 
metadataColumns);
+
       abstract Builder setHadoopConfig(@Nullable Map<String, String> 
hadoopConfig);
 
       abstract ReadChanges build();
@@ -259,6 +285,20 @@ public class DeltaIO {
       return toBuilder().setEndTimestamp(endTimestamp).build();
     }
 
+    public ReadChanges withMetadataColumns(String... metadataColumns) {
+      for (String col : metadataColumns) {
+        if (!col.equals(CHANGE_TYPE_COLUMN)
+            && !col.equals(COMMIT_VERSION_COLUMN)
+            && !col.equals(COMMIT_TIMESTAMP_COLUMN)) {
+          throw new IllegalArgumentException(
+              String.format(
+                  "Unsupported metadata column %s. Supported columns are: %s, 
%s, and %s.",
+                  col, CHANGE_TYPE_COLUMN, COMMIT_VERSION_COLUMN, 
COMMIT_TIMESTAMP_COLUMN));
+        }
+      }
+      return 
toBuilder().setMetadataColumns(Arrays.asList(metadataColumns)).build();
+    }
+
     public ReadChanges withConfig(Map<String, String> config) {
       return toBuilder().setHadoopConfig(config).build();
     }
@@ -310,7 +350,8 @@ public class DeltaIO {
       if (deltaSchema == null) {
         throw new IllegalStateException("Table schema is null.");
       }
-      Schema beamSchema = ReadRows.convertToBeamSchema(deltaSchema);
+      Schema baseSchema = ReadRows.convertToBeamSchema(deltaSchema);
+      Schema publicBeamSchema = buildPublicBeamSchema(baseSchema, 
getMetadataColumns());
 
       return input
           .apply("Create Path", Create.of(path))
@@ -323,8 +364,9 @@ public class DeltaIO {
                       getStartTimestamp(),
                       getEndVersion(),
                       getEndTimestamp())))
-          .apply("Read CDF Data", ParDo.of(new 
DeltaCDCSourceDoFn(hadoopConfig)))
-          .setRowSchema(beamSchema);
+          .apply(
+              "Read CDF Data", ParDo.of(new DeltaCDCSourceDoFn(hadoopConfig, 
getMetadataColumns())))
+          .setRowSchema(publicBeamSchema);
     }
   }
 }
diff --git 
a/sdks/java/io/delta/src/test/java/org/apache/beam/sdk/io/delta/DeltaIOIT.java 
b/sdks/java/io/delta/src/test/java/org/apache/beam/sdk/io/delta/DeltaIOIT.java
index e0d35f30faa..ad526008b20 100644
--- 
a/sdks/java/io/delta/src/test/java/org/apache/beam/sdk/io/delta/DeltaIOIT.java
+++ 
b/sdks/java/io/delta/src/test/java/org/apache/beam/sdk/io/delta/DeltaIOIT.java
@@ -34,11 +34,14 @@ import 
io.delta.kernel.defaults.internal.data.DefaultColumnarBatch;
 import io.delta.kernel.engine.Engine;
 import io.delta.kernel.types.DataType;
 import io.delta.kernel.types.IntegerType;
+import io.delta.kernel.types.LongType;
 import io.delta.kernel.types.StringType;
 import io.delta.kernel.types.StructType;
+import io.delta.kernel.types.TimestampType;
 import io.delta.kernel.utils.CloseableIterable;
 import io.delta.kernel.utils.CloseableIterator;
 import io.delta.kernel.utils.DataFileStatus;
+import java.util.ArrayList;
 import java.util.Collections;
 import java.util.HashMap;
 import java.util.List;
@@ -47,13 +50,17 @@ import java.util.Optional;
 import java.util.stream.Collectors;
 import java.util.stream.IntStream;
 import org.apache.beam.sdk.managed.Managed;
+import org.apache.beam.sdk.options.ExperimentalOptions;
 import org.apache.beam.sdk.schemas.Schema;
 import org.apache.beam.sdk.testing.PAssert;
 import org.apache.beam.sdk.testing.TestPipeline;
+import org.apache.beam.sdk.transforms.DoFn;
+import org.apache.beam.sdk.transforms.ParDo;
 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.collect.ImmutableMap;
 import org.apache.hadoop.conf.Configuration;
+import org.joda.time.Instant;
 import org.junit.After;
 import org.junit.Before;
 import org.junit.Rule;
@@ -77,6 +84,7 @@ public class DeltaIOIT {
   private String repoPath;
   private String repoPrefix;
   private Storage storage;
+  private String version0FilePath;
 
   private static final Schema ROW_SCHEMA =
       Schema.builder().addInt32Field("id").addStringField("name").build();
@@ -127,7 +135,11 @@ public class DeltaIOIT {
 
     TransactionBuilder txnBuilder =
         table.createTransactionBuilder(engine, "DeltaIOIT", 
Operation.CREATE_TABLE);
-    txnBuilder = txnBuilder.withSchema(engine, deltaSchema);
+    txnBuilder =
+        txnBuilder
+            .withSchema(engine, deltaSchema)
+            .withTableProperties(
+                engine, Collections.singletonMap("delta.enableChangeDataFeed", 
"true"));
     Transaction txn = txnBuilder.build(engine);
     io.delta.kernel.data.Row txnState = txn.getTransactionState(engine);
 
@@ -209,8 +221,33 @@ public class DeltaIOIT {
     CloseableIterator<io.delta.kernel.data.Row> dataActions =
         Transaction.generateAppendActions(engine, txnState, dataFiles, 
writeContext);
 
+    List<io.delta.kernel.data.Row> addActionsList = new ArrayList<>();
+    while (dataActions.hasNext()) {
+      addActionsList.add(dataActions.next());
+    }
+
+    if (!addActionsList.isEmpty()) {
+      io.delta.kernel.data.Row action = addActionsList.get(0);
+      int addOrdinal = action.getSchema().indexOf("add");
+      if (addOrdinal < 0) {
+        throw new IllegalStateException(
+            "Expected append action to contain 'add' field, but it didn't: " + 
action.getSchema());
+      }
+      io.delta.kernel.data.Row addAction = action.getStruct(addOrdinal);
+      if (addAction == null) {
+        throw new IllegalStateException("Action 'add' struct is null");
+      }
+      int pathOrdinal = addAction.getSchema().indexOf("path");
+      if (pathOrdinal < 0) {
+        throw new IllegalStateException(
+            "'add' action schema does not contain 'path': " + 
addAction.getSchema());
+      }
+      version0FilePath = addAction.getString(pathOrdinal);
+    }
+
     CloseableIterable<io.delta.kernel.data.Row> dataActionsIterable =
-        CloseableIterable.inMemoryIterable(dataActions);
+        CloseableIterable.inMemoryIterable(
+            
io.delta.kernel.internal.util.Utils.toCloseableIterator(addActionsList.iterator()));
 
     TransactionCommitResult commitResult = txn.commit(engine, 
dataActionsIterable);
 
@@ -238,6 +275,9 @@ public class DeltaIOIT {
 
   @Test
   public void testReadDeltaLakeTable() {
+    ExperimentalOptions options = 
readPipeline.getOptions().as(ExperimentalOptions.class);
+    ExperimentalOptions.addExperiment(options, "use_runner_v2");
+
     Map<String, String> hadoopConfig = new HashMap<>();
     hadoopConfig.put("fs.gs.impl", 
"com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem");
     hadoopConfig.put(
@@ -261,4 +301,128 @@ public class DeltaIOIT {
     PAssert.that(output).containsInAnyOrder(TEST_ROWS);
     readPipeline.run().waitUntilFinish();
   }
+
+  @Test
+  public void testReadChangesDeltaLake() throws Exception {
+    ExperimentalOptions options = 
readPipeline.getOptions().as(ExperimentalOptions.class);
+    List<String> experiments = options.getExperiments();
+    if (experiments != null) {
+      List<String> modifiableExperiments = new 
java.util.ArrayList<>(experiments);
+      // TODO: remove this when Runner v2 supports elements that includes CDC 
metadata
+      // (ValueKind).
+      modifiableExperiments.remove("use_runner_v2");
+      options.setExperiments(modifiableExperiments);
+    }
+
+    Map<String, String> hadoopConfig = new HashMap<>();
+    hadoopConfig.put("fs.gs.impl", 
"com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem");
+    hadoopConfig.put(
+        "fs.AbstractFileSystem.gs.impl", 
"com.google.cloud.hadoop.fs.gcs.GoogleHadoopFS");
+    hadoopConfig.put("fs.gs.auth.type", "APPLICATION_DEFAULT");
+    String project =
+        readPipeline
+            .getOptions()
+            .as(org.apache.beam.sdk.extensions.gcp.options.GcpOptions.class)
+            .getProject();
+    if (project != null) {
+      hadoopConfig.put("fs.gs.project.id", project);
+    }
+
+    org.apache.hadoop.conf.Configuration conf = new 
org.apache.hadoop.conf.Configuration();
+    for (Map.Entry<String, String> entry : hadoopConfig.entrySet()) {
+      conf.set(entry.getKey(), entry.getValue());
+    }
+    Engine engine = DefaultEngine.create(conf);
+
+    StructType deltaSchema =
+        new StructType().add("id", IntegerType.INTEGER).add("name", 
StringType.STRING);
+
+    // 1. Write version 1 containing cdc actions for testing updates and 
deletes
+    Schema cdcWriteSchema =
+        Schema.builder()
+            .addField("id", Schema.FieldType.INT32)
+            .addField("name", Schema.FieldType.STRING)
+            .addField(DeltaIO.CHANGE_TYPE_COLUMN, Schema.FieldType.STRING)
+            .addField(DeltaIO.COMMIT_VERSION_COLUMN, Schema.FieldType.INT64)
+            .addField(DeltaIO.COMMIT_TIMESTAMP_COLUMN, 
Schema.FieldType.DATETIME)
+            .build();
+    StructType cdcWriteDeltaSchema =
+        new StructType()
+            .add("id", IntegerType.INTEGER)
+            .add("name", StringType.STRING)
+            .add(DeltaIO.CHANGE_TYPE_COLUMN, StringType.STRING)
+            .add(DeltaIO.COMMIT_VERSION_COLUMN, LongType.LONG)
+            .add(DeltaIO.COMMIT_TIMESTAMP_COLUMN, TimestampType.TIMESTAMP);
+
+    Row cdcRow1 =
+        Row.withSchema(cdcWriteSchema)
+            .addValues(0, "name_0", "delete", 1L, new Instant(123456789000L))
+            .build();
+    Row cdcRow2 =
+        Row.withSchema(cdcWriteSchema)
+            .addValues(1, "name_1", "update_preimage", 1L, new 
Instant(123456789000L))
+            .build();
+    Row cdcRow3 =
+        Row.withSchema(cdcWriteSchema)
+            .addValues(1, "name_1_updated", "update_postimage", 1L, new 
Instant(123456789000L))
+            .build();
+
+    DeltaWriteTestUtils.writeCdcCommit(
+        engine,
+        repoPath,
+        1L,
+        System.currentTimeMillis(),
+        deltaSchema,
+        null,
+        version0FilePath,
+        java.util.Arrays.asList(cdcRow1, cdcRow2, cdcRow3),
+        cdcWriteDeltaSchema);
+
+    // 2. Read CDF data from table using Managed.read(Managed.DELTA_LAKE_CDC)
+    Map<String, Object> readConfig = new HashMap<>();
+    readConfig.put("table", repoPath);
+    readConfig.put("start_version", 0L);
+    readConfig.put("hadoop_config", hadoopConfig);
+    readConfig.put(
+        "include_metadata_columns",
+        java.util.Arrays.asList(
+            DeltaIO.CHANGE_TYPE_COLUMN,
+            DeltaIO.COMMIT_VERSION_COLUMN,
+            DeltaIO.COMMIT_TIMESTAMP_COLUMN));
+
+    PCollection<Row> output =
+        readPipeline
+            .apply(Managed.read(Managed.DELTA_LAKE_CDC).withConfig(readConfig))
+            .getSinglePCollection();
+
+    PCollection<String> formattedOutput =
+        output.apply("Format Row with Metadata", ParDo.of(new 
FormatITRowWithMetadata()));
+
+    // Generate expected outputs for version 0 (inserts of id 0-99)
+    List<String> expectedOutputs = new ArrayList<>();
+    for (int i = 0; i < 100; i++) {
+      expectedOutputs.add(String.format("%d:name_%d:insert:v0", i, i));
+    }
+    // Expected outputs for version 1
+    expectedOutputs.add("0:name_0:delete:v1");
+    expectedOutputs.add("1:name_1:update_preimage:v1");
+    expectedOutputs.add("1:name_1_updated:update_postimage:v1");
+
+    PAssert.that(formattedOutput).containsInAnyOrder(expectedOutputs);
+
+    readPipeline.run().waitUntilFinish();
+  }
+
+  private static final class FormatITRowWithMetadata extends DoFn<Row, String> 
{
+    @ProcessElement
+    public void process(@Element Row row, OutputReceiver<String> out) {
+      out.output(
+          String.format(
+              "%d:%s:%s:v%d",
+              row.getInt32("id"),
+              row.getString("name"),
+              row.getString(DeltaIO.CHANGE_TYPE_COLUMN),
+              row.getInt64(DeltaIO.COMMIT_VERSION_COLUMN)));
+    }
+  }
 }
diff --git 
a/sdks/java/io/delta/src/test/java/org/apache/beam/sdk/io/delta/DeltaIOTest.java
 
b/sdks/java/io/delta/src/test/java/org/apache/beam/sdk/io/delta/DeltaIOTest.java
index f00b34be460..97534edf79e 100644
--- 
a/sdks/java/io/delta/src/test/java/org/apache/beam/sdk/io/delta/DeltaIOTest.java
+++ 
b/sdks/java/io/delta/src/test/java/org/apache/beam/sdk/io/delta/DeltaIOTest.java
@@ -17,23 +17,11 @@
  */
 package org.apache.beam.sdk.io.delta;
 
-import io.delta.kernel.DataWriteContext;
-import io.delta.kernel.Operation;
-import io.delta.kernel.Table;
-import io.delta.kernel.Transaction;
-import io.delta.kernel.TransactionBuilder;
-import io.delta.kernel.TransactionCommitResult;
-import io.delta.kernel.data.ColumnVector;
-import io.delta.kernel.data.ColumnarBatch;
-import io.delta.kernel.data.FilteredColumnarBatch;
-import io.delta.kernel.data.MapValue;
 import io.delta.kernel.defaults.engine.DefaultEngine;
-import io.delta.kernel.defaults.internal.data.DefaultColumnarBatch;
 import io.delta.kernel.engine.Engine;
 import io.delta.kernel.types.ArrayType;
 import io.delta.kernel.types.BinaryType;
 import io.delta.kernel.types.BooleanType;
-import io.delta.kernel.types.DataType;
 import io.delta.kernel.types.DateType;
 import io.delta.kernel.types.DoubleType;
 import io.delta.kernel.types.FloatType;
@@ -44,19 +32,12 @@ import io.delta.kernel.types.StringType;
 import io.delta.kernel.types.StructField;
 import io.delta.kernel.types.StructType;
 import io.delta.kernel.types.TimestampType;
-import io.delta.kernel.utils.CloseableIterable;
-import io.delta.kernel.utils.CloseableIterator;
-import io.delta.kernel.utils.DataFileStatus;
 import java.io.File;
-import java.math.BigDecimal;
 import java.nio.charset.StandardCharsets;
 import java.nio.file.Files;
-import java.util.ArrayList;
 import java.util.Collections;
 import java.util.HashMap;
-import java.util.List;
 import java.util.Map;
-import java.util.Optional;
 import org.apache.avro.generic.GenericRecord;
 import org.apache.beam.sdk.extensions.avro.coders.AvroCoder;
 import org.apache.beam.sdk.extensions.avro.schemas.utils.AvroUtils;
@@ -75,10 +56,10 @@ import org.apache.beam.sdk.transforms.ParDo;
 import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
 import org.apache.beam.sdk.transforms.windowing.PaneInfo;
 import org.apache.beam.sdk.values.PCollection;
+import org.apache.beam.sdk.values.PCollectionRowTuple;
 import org.apache.beam.sdk.values.Row;
 import org.apache.beam.sdk.values.ValueKind;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
-import org.checkerframework.checker.nullness.qual.Nullable;
 import org.joda.time.Instant;
 import org.junit.Assert;
 import org.junit.Rule;
@@ -376,7 +357,7 @@ public class DeltaIOTest {
     Row row = Row.withSchema(schema).addValues("test-name").build();
     StructType deltaSchema = new StructType().add("name", StringType.STRING);
 
-    writeAppendCommit(
+    DeltaWriteTestUtils.writeAppendCommit(
         engine,
         tableDir.getAbsolutePath(),
         0L,
@@ -791,7 +772,7 @@ public class DeltaIOTest {
     Row tableRow2 = Row.withSchema(tableSchema).addValues("row-2").build();
     StructType deltaSchema = new StructType().add("name", StringType.STRING);
 
-    writeAppendCommit(
+    DeltaWriteTestUtils.writeAppendCommit(
         engine,
         tableDir.getAbsolutePath(),
         0L,
@@ -827,7 +808,7 @@ public class DeltaIOTest {
             .addValues("row-2", "delete", 1L, new Instant(123456789000L))
             .build();
 
-    writeCdcCommit(
+    DeltaWriteTestUtils.writeCdcCommit(
         engine,
         tableDir.getAbsolutePath(),
         1L,
@@ -857,6 +838,474 @@ public class DeltaIOTest {
     readPipeline.run().waitUntilFinish();
   }
 
+  @Test
+  public void testReadChangesAndNormalReadWithCDCAndAppend() throws Exception {
+    File tableDir = tempFolder.newFolder("delta-table-cdc-and-append");
+    Engine engine = DefaultEngine.create(new 
org.apache.hadoop.conf.Configuration());
+
+    // 1. Write parquet files for Version 0 (insert-only commit)
+    Schema tableSchema = Schema.builder().addField("name", 
Schema.FieldType.STRING).build();
+    Row tableRow1 = Row.withSchema(tableSchema).addValues("row-1").build();
+    Row tableRow2 = Row.withSchema(tableSchema).addValues("row-2").build();
+    StructType deltaSchema = new StructType().add("name", StringType.STRING);
+
+    DeltaWriteTestUtils.writeAppendCommit(
+        engine,
+        tableDir.getAbsolutePath(),
+        0L,
+        100000000000L,
+        deltaSchema,
+        java.util.Arrays.asList(tableRow1, tableRow2));
+
+    // 2. Write cdc and append parquet files for Version 1 (commit with cdc 
and add actions)
+    Schema cdcWriteSchema =
+        Schema.builder()
+            .addField("name", Schema.FieldType.STRING)
+            .addField(DeltaIO.CHANGE_TYPE_COLUMN, Schema.FieldType.STRING)
+            .addField(DeltaIO.COMMIT_VERSION_COLUMN, Schema.FieldType.INT64)
+            .addField(DeltaIO.COMMIT_TIMESTAMP_COLUMN, 
Schema.FieldType.DATETIME)
+            .build();
+    StructType cdcWriteDeltaSchema =
+        new StructType()
+            .add("name", StringType.STRING)
+            .add(DeltaIO.CHANGE_TYPE_COLUMN, StringType.STRING)
+            .add(DeltaIO.COMMIT_VERSION_COLUMN, LongType.LONG)
+            .add(DeltaIO.COMMIT_TIMESTAMP_COLUMN, TimestampType.TIMESTAMP);
+
+    Row cdcRow =
+        Row.withSchema(cdcWriteSchema)
+            .addValues("row-3", "insert", 1L, new Instant(123456789000L))
+            .build();
+
+    Row appendRow = Row.withSchema(tableSchema).addValues("row-3").build();
+
+    DeltaWriteTestUtils.writeCdcCommit(
+        engine,
+        tableDir.getAbsolutePath(),
+        1L,
+        200000000000L,
+        deltaSchema,
+        java.util.Arrays.asList(appendRow),
+        null,
+        java.util.Arrays.asList(cdcRow),
+        cdcWriteDeltaSchema);
+
+    // 3. Read CDF data from table using ReadChanges
+    PCollection<Row> outputCDC =
+        readPipeline.apply(
+            "Read Changes",
+            
DeltaIO.readChanges().from(tableDir.getAbsolutePath()).withStartVersion(0L));
+
+    PCollection<String> formattedOutputCDC =
+        outputCDC.apply("Format CDC Row", ParDo.of(new 
FormatValueKindAndRow()));
+
+    PAssert.that(formattedOutputCDC)
+        .containsInAnyOrder("INSERT:row-1", "INSERT:row-2", "INSERT:row-3");
+
+    // 4. Read latest snapshot using normal read via writePipeline
+    PCollection<Row> outputNormal =
+        writePipeline.apply("Read Normal", 
DeltaIO.readRows().from(tableDir.getAbsolutePath()));
+
+    PCollection<String> formattedOutputNormal =
+        outputNormal.apply("Format Normal Row", ParDo.of(new FormatRowName()));
+
+    PAssert.that(formattedOutputNormal).containsInAnyOrder("row-1", "row-2", 
"row-3");
+
+    readPipeline.run().waitUntilFinish();
+    writePipeline.run().waitUntilFinish();
+  }
+
+  @Test
+  public void testReadChangesWithSchemaTransformProvider() throws Exception {
+    File tableDir = tempFolder.newFolder("delta-table-changes-provider");
+    Engine engine = DefaultEngine.create(new 
org.apache.hadoop.conf.Configuration());
+
+    // 1. Write parquet files for Version 0 (insert-only commit)
+    Schema tableSchema = Schema.builder().addField("name", 
Schema.FieldType.STRING).build();
+    Row tableRow1 = Row.withSchema(tableSchema).addValues("row-1").build();
+    Row tableRow2 = Row.withSchema(tableSchema).addValues("row-2").build();
+    StructType deltaSchema = new StructType().add("name", StringType.STRING);
+
+    DeltaWriteTestUtils.writeAppendCommit(
+        engine,
+        tableDir.getAbsolutePath(),
+        0L,
+        100000000000L,
+        deltaSchema,
+        java.util.Arrays.asList(tableRow1, tableRow2));
+
+    // 2. Write cdc parquet file for Version 1 (commit with cdc actions)
+    Schema cdcWriteSchema =
+        Schema.builder()
+            .addField("name", Schema.FieldType.STRING)
+            .addField(DeltaIO.CHANGE_TYPE_COLUMN, Schema.FieldType.STRING)
+            .addField(DeltaIO.COMMIT_VERSION_COLUMN, Schema.FieldType.INT64)
+            .addField(DeltaIO.COMMIT_TIMESTAMP_COLUMN, 
Schema.FieldType.DATETIME)
+            .build();
+    StructType cdcWriteDeltaSchema =
+        new StructType()
+            .add("name", StringType.STRING)
+            .add(DeltaIO.CHANGE_TYPE_COLUMN, StringType.STRING)
+            .add(DeltaIO.COMMIT_VERSION_COLUMN, LongType.LONG)
+            .add(DeltaIO.COMMIT_TIMESTAMP_COLUMN, TimestampType.TIMESTAMP);
+
+    Row cdcRow1 =
+        Row.withSchema(cdcWriteSchema)
+            .addValues("row-1", "update_preimage", 1L, new 
Instant(123456789000L))
+            .build();
+    Row cdcRow2 =
+        Row.withSchema(cdcWriteSchema)
+            .addValues("row-1-updated", "update_postimage", 1L, new 
Instant(123456789000L))
+            .build();
+    Row cdcRow3 =
+        Row.withSchema(cdcWriteSchema)
+            .addValues("row-2", "delete", 1L, new Instant(123456789000L))
+            .build();
+
+    DeltaWriteTestUtils.writeCdcCommit(
+        engine,
+        tableDir.getAbsolutePath(),
+        1L,
+        200000000000L,
+        deltaSchema,
+        null,
+        null,
+        java.util.Arrays.asList(cdcRow1, cdcRow2, cdcRow3),
+        cdcWriteDeltaSchema);
+
+    // 3. Read CDF data from table using DeltaCdcReadSchemaTransformProvider
+    DeltaCdcReadSchemaTransformProvider.Configuration config =
+        DeltaCdcReadSchemaTransformProvider.Configuration.builder()
+            .setTable(tableDir.getAbsolutePath())
+            .setStartVersion(0L)
+            .build();
+
+    PCollection<Row> output =
+        PCollectionRowTuple.empty(readPipeline)
+            .apply(new DeltaCdcReadSchemaTransformProvider().from(config))
+            .get(DeltaCdcReadSchemaTransformProvider.OUTPUT_TAG);
+
+    PCollection<String> formattedOutput =
+        output.apply("Format ValueKind and Row", ParDo.of(new 
FormatValueKindAndRow()));
+
+    PAssert.that(formattedOutput)
+        .containsInAnyOrder(
+            "INSERT:row-1",
+            "INSERT:row-2",
+            "UPDATE_BEFORE:row-1",
+            "UPDATE_AFTER:row-1-updated",
+            "DELETE:row-2");
+
+    readPipeline.run().waitUntilFinish();
+  }
+
+  @Test
+  public void testReadChangesWithMetadataColumns() throws Exception {
+    File tableDir = tempFolder.newFolder("delta-table-changes-metadata");
+    Engine engine = DefaultEngine.create(new 
org.apache.hadoop.conf.Configuration());
+
+    // 1. Write parquet files for Version 0 (insert-only commit)
+    Schema tableSchema = Schema.builder().addField("name", 
Schema.FieldType.STRING).build();
+    Row tableRow1 = Row.withSchema(tableSchema).addValues("row-1").build();
+    Row tableRow2 = Row.withSchema(tableSchema).addValues("row-2").build();
+    StructType deltaSchema = new StructType().add("name", StringType.STRING);
+
+    DeltaWriteTestUtils.writeAppendCommit(
+        engine,
+        tableDir.getAbsolutePath(),
+        0L,
+        100000000000L,
+        deltaSchema,
+        java.util.Arrays.asList(tableRow1, tableRow2));
+
+    // 2. Write cdc parquet file for Version 1 (commit with cdc actions)
+    Schema cdcWriteSchema =
+        Schema.builder()
+            .addField("name", Schema.FieldType.STRING)
+            .addField(DeltaIO.CHANGE_TYPE_COLUMN, Schema.FieldType.STRING)
+            .addField(DeltaIO.COMMIT_VERSION_COLUMN, Schema.FieldType.INT64)
+            .addField(DeltaIO.COMMIT_TIMESTAMP_COLUMN, 
Schema.FieldType.DATETIME)
+            .build();
+    StructType cdcWriteDeltaSchema =
+        new StructType()
+            .add("name", StringType.STRING)
+            .add(DeltaIO.CHANGE_TYPE_COLUMN, StringType.STRING)
+            .add(DeltaIO.COMMIT_VERSION_COLUMN, LongType.LONG)
+            .add(DeltaIO.COMMIT_TIMESTAMP_COLUMN, TimestampType.TIMESTAMP);
+
+    Row cdcRow1 =
+        Row.withSchema(cdcWriteSchema)
+            .addValues("row-1", "update_preimage", 1L, new 
Instant(123456789000L))
+            .build();
+    Row cdcRow2 =
+        Row.withSchema(cdcWriteSchema)
+            .addValues("row-1-updated", "update_postimage", 1L, new 
Instant(123456789000L))
+            .build();
+    Row cdcRow3 =
+        Row.withSchema(cdcWriteSchema)
+            .addValues("row-2", "delete", 1L, new Instant(123456789000L))
+            .build();
+
+    DeltaWriteTestUtils.writeCdcCommit(
+        engine,
+        tableDir.getAbsolutePath(),
+        1L,
+        200000000000L,
+        deltaSchema,
+        null,
+        null,
+        java.util.Arrays.asList(cdcRow1, cdcRow2, cdcRow3),
+        cdcWriteDeltaSchema);
+
+    // 3. Read CDF data from table using DeltaCdcReadSchemaTransformProvider 
requesting metadata
+    // columns
+    DeltaCdcReadSchemaTransformProvider.Configuration config =
+        DeltaCdcReadSchemaTransformProvider.Configuration.builder()
+            .setTable(tableDir.getAbsolutePath())
+            .setStartVersion(0L)
+            .setIncludeMetadataColumns(
+                java.util.Arrays.asList(
+                    DeltaIO.CHANGE_TYPE_COLUMN,
+                    DeltaIO.COMMIT_VERSION_COLUMN,
+                    DeltaIO.COMMIT_TIMESTAMP_COLUMN))
+            .build();
+
+    PCollection<Row> output =
+        PCollectionRowTuple.empty(readPipeline)
+            .apply(new DeltaCdcReadSchemaTransformProvider().from(config))
+            .get(DeltaCdcReadSchemaTransformProvider.OUTPUT_TAG);
+
+    PCollection<String> formattedOutput =
+        output.apply("Format Row with Metadata", ParDo.of(new 
FormatRowWithMetadata()));
+
+    PAssert.that(formattedOutput)
+        .containsInAnyOrder(
+            "row-1:insert:v0:t100000000000",
+            "row-2:insert:v0:t100000000000",
+            "row-1:update_preimage:v1:t123456789000",
+            "row-1-updated:update_postimage:v1:t123456789000",
+            "row-2:delete:v1:t123456789000");
+
+    readPipeline.run().waitUntilFinish();
+  }
+
+  @Test
+  public void testReadChangesWithSubsetOfMetadataColumns() throws Exception {
+    File tableDir = 
tempFolder.newFolder("delta-table-changes-subset-metadata");
+    Engine engine = DefaultEngine.create(new 
org.apache.hadoop.conf.Configuration());
+
+    // 1. Write parquet files for Version 0 (insert-only commit)
+    Schema tableSchema = Schema.builder().addField("name", 
Schema.FieldType.STRING).build();
+    Row tableRow1 = Row.withSchema(tableSchema).addValues("row-1").build();
+    StructType deltaSchema = new StructType().add("name", StringType.STRING);
+
+    DeltaWriteTestUtils.writeAppendCommit(
+        engine,
+        tableDir.getAbsolutePath(),
+        0L,
+        100000000000L,
+        deltaSchema,
+        java.util.Arrays.asList(tableRow1));
+
+    // 2. Read CDF data from table requesting ONLY _change_type
+    DeltaCdcReadSchemaTransformProvider.Configuration config =
+        DeltaCdcReadSchemaTransformProvider.Configuration.builder()
+            .setTable(tableDir.getAbsolutePath())
+            .setStartVersion(0L)
+            .setIncludeMetadataColumns(
+                
java.util.Collections.singletonList(DeltaIO.CHANGE_TYPE_COLUMN))
+            .build();
+
+    PCollection<Row> output =
+        PCollectionRowTuple.empty(readPipeline)
+            .apply(new DeltaCdcReadSchemaTransformProvider().from(config))
+            .get(DeltaCdcReadSchemaTransformProvider.OUTPUT_TAG);
+
+    // Verify schema does not contain version or timestamp
+    
org.junit.Assert.assertTrue(output.getSchema().hasField(DeltaIO.CHANGE_TYPE_COLUMN));
+    
org.junit.Assert.assertFalse(output.getSchema().hasField(DeltaIO.COMMIT_VERSION_COLUMN));
+    
org.junit.Assert.assertFalse(output.getSchema().hasField(DeltaIO.COMMIT_TIMESTAMP_COLUMN));
+
+    PCollection<String> formattedOutput =
+        output.apply("Format Row", ParDo.of(new FormatRowSubsetMetadata()));
+
+    PAssert.that(formattedOutput).containsInAnyOrder("row-1:insert");
+
+    readPipeline.run().waitUntilFinish();
+  }
+
+  @Test
+  public void testReadChangesWithCommitVersionMetadataColumn() throws 
Exception {
+    File tableDir = 
tempFolder.newFolder("delta-table-changes-version-metadata");
+    Engine engine = DefaultEngine.create(new 
org.apache.hadoop.conf.Configuration());
+
+    // 1. Write parquet files for Version 0 (insert-only commit)
+    Schema tableSchema = Schema.builder().addField("name", 
Schema.FieldType.STRING).build();
+    Row tableRow1 = Row.withSchema(tableSchema).addValues("row-1").build();
+    StructType deltaSchema = new StructType().add("name", StringType.STRING);
+
+    DeltaWriteTestUtils.writeAppendCommit(
+        engine,
+        tableDir.getAbsolutePath(),
+        0L,
+        100000000000L,
+        deltaSchema,
+        java.util.Arrays.asList(tableRow1));
+
+    // 2. Read CDF data from table requesting ONLY _commit_version
+    DeltaCdcReadSchemaTransformProvider.Configuration config =
+        DeltaCdcReadSchemaTransformProvider.Configuration.builder()
+            .setTable(tableDir.getAbsolutePath())
+            .setStartVersion(0L)
+            .setIncludeMetadataColumns(
+                
java.util.Collections.singletonList(DeltaIO.COMMIT_VERSION_COLUMN))
+            .build();
+
+    PCollection<Row> output =
+        PCollectionRowTuple.empty(readPipeline)
+            .apply(new DeltaCdcReadSchemaTransformProvider().from(config))
+            .get(DeltaCdcReadSchemaTransformProvider.OUTPUT_TAG);
+
+    // Verify schema contains version but not change type or timestamp
+    
org.junit.Assert.assertFalse(output.getSchema().hasField(DeltaIO.CHANGE_TYPE_COLUMN));
+    
org.junit.Assert.assertTrue(output.getSchema().hasField(DeltaIO.COMMIT_VERSION_COLUMN));
+    
org.junit.Assert.assertFalse(output.getSchema().hasField(DeltaIO.COMMIT_TIMESTAMP_COLUMN));
+
+    PCollection<String> formattedOutput =
+        output.apply("Format Row", ParDo.of(new FormatRowVersionMetadata()));
+
+    PAssert.that(formattedOutput).containsInAnyOrder("row-1:0");
+
+    readPipeline.run().waitUntilFinish();
+  }
+
+  @Test
+  public void testReadChangesWithCommitTimestampMetadataColumn() throws 
Exception {
+    File tableDir = 
tempFolder.newFolder("delta-table-changes-timestamp-metadata");
+    Engine engine = DefaultEngine.create(new 
org.apache.hadoop.conf.Configuration());
+
+    // 1. Write parquet files for Version 0 (insert-only commit)
+    Schema tableSchema = Schema.builder().addField("name", 
Schema.FieldType.STRING).build();
+    Row tableRow1 = Row.withSchema(tableSchema).addValues("row-1").build();
+    StructType deltaSchema = new StructType().add("name", StringType.STRING);
+
+    DeltaWriteTestUtils.writeAppendCommit(
+        engine,
+        tableDir.getAbsolutePath(),
+        0L,
+        100000000000L,
+        deltaSchema,
+        java.util.Arrays.asList(tableRow1));
+
+    // 2. Read CDF data from table requesting ONLY _commit_timestamp
+    DeltaCdcReadSchemaTransformProvider.Configuration config =
+        DeltaCdcReadSchemaTransformProvider.Configuration.builder()
+            .setTable(tableDir.getAbsolutePath())
+            .setStartVersion(0L)
+            .setIncludeMetadataColumns(
+                
java.util.Collections.singletonList(DeltaIO.COMMIT_TIMESTAMP_COLUMN))
+            .build();
+
+    PCollection<Row> output =
+        PCollectionRowTuple.empty(readPipeline)
+            .apply(new DeltaCdcReadSchemaTransformProvider().from(config))
+            .get(DeltaCdcReadSchemaTransformProvider.OUTPUT_TAG);
+
+    // Verify schema contains timestamp but not change type or version
+    
org.junit.Assert.assertFalse(output.getSchema().hasField(DeltaIO.CHANGE_TYPE_COLUMN));
+    
org.junit.Assert.assertFalse(output.getSchema().hasField(DeltaIO.COMMIT_VERSION_COLUMN));
+    
org.junit.Assert.assertTrue(output.getSchema().hasField(DeltaIO.COMMIT_TIMESTAMP_COLUMN));
+
+    PCollection<String> formattedOutput =
+        output.apply("Format Row", ParDo.of(new FormatRowTimestampMetadata()));
+
+    PAssert.that(formattedOutput).containsInAnyOrder("row-1:100000000000");
+
+    readPipeline.run().waitUntilFinish();
+  }
+
+  @Test
+  public void testReadChangesWithManaged() throws Exception {
+    File tableDir = tempFolder.newFolder("delta-table-changes-managed");
+    Engine engine = DefaultEngine.create(new 
org.apache.hadoop.conf.Configuration());
+
+    // 1. Write parquet files for Version 0 (insert-only commit)
+    Schema tableSchema = Schema.builder().addField("name", 
Schema.FieldType.STRING).build();
+    Row tableRow1 = Row.withSchema(tableSchema).addValues("row-1").build();
+    Row tableRow2 = Row.withSchema(tableSchema).addValues("row-2").build();
+    StructType deltaSchema = new StructType().add("name", StringType.STRING);
+
+    DeltaWriteTestUtils.writeAppendCommit(
+        engine,
+        tableDir.getAbsolutePath(),
+        0L,
+        100000000000L,
+        deltaSchema,
+        java.util.Arrays.asList(tableRow1, tableRow2));
+
+    // 2. Write cdc parquet file for Version 1 (commit with cdc actions)
+    Schema cdcWriteSchema =
+        Schema.builder()
+            .addField("name", Schema.FieldType.STRING)
+            .addField(DeltaIO.CHANGE_TYPE_COLUMN, Schema.FieldType.STRING)
+            .addField(DeltaIO.COMMIT_VERSION_COLUMN, Schema.FieldType.INT64)
+            .addField(DeltaIO.COMMIT_TIMESTAMP_COLUMN, 
Schema.FieldType.DATETIME)
+            .build();
+    StructType cdcWriteDeltaSchema =
+        new StructType()
+            .add("name", StringType.STRING)
+            .add(DeltaIO.CHANGE_TYPE_COLUMN, StringType.STRING)
+            .add(DeltaIO.COMMIT_VERSION_COLUMN, LongType.LONG)
+            .add(DeltaIO.COMMIT_TIMESTAMP_COLUMN, TimestampType.TIMESTAMP);
+
+    Row cdcRow1 =
+        Row.withSchema(cdcWriteSchema)
+            .addValues("row-1", "update_preimage", 1L, new 
Instant(123456789000L))
+            .build();
+    Row cdcRow2 =
+        Row.withSchema(cdcWriteSchema)
+            .addValues("row-1-updated", "update_postimage", 1L, new 
Instant(123456789000L))
+            .build();
+    Row cdcRow3 =
+        Row.withSchema(cdcWriteSchema)
+            .addValues("row-2", "delete", 1L, new Instant(123456789000L))
+            .build();
+
+    DeltaWriteTestUtils.writeCdcCommit(
+        engine,
+        tableDir.getAbsolutePath(),
+        1L,
+        200000000000L,
+        deltaSchema,
+        null,
+        null,
+        java.util.Arrays.asList(cdcRow1, cdcRow2, cdcRow3),
+        cdcWriteDeltaSchema);
+
+    // 3. Read CDF data from table using Managed.read(Managed.DELTA_LAKE_CDC)
+    Map<String, Object> config = new HashMap<>();
+    config.put("table", tableDir.getAbsolutePath());
+    config.put("start_version", 0L);
+
+    PCollection<Row> output =
+        readPipeline
+            .apply(Managed.read(Managed.DELTA_LAKE_CDC).withConfig(config))
+            .getSinglePCollection();
+
+    PCollection<String> formattedOutput =
+        output.apply("Format ValueKind and Row", ParDo.of(new 
FormatValueKindAndRow()));
+
+    PAssert.that(formattedOutput)
+        .containsInAnyOrder(
+            "INSERT:row-1",
+            "INSERT:row-2",
+            "UPDATE_BEFORE:row-1",
+            "UPDATE_AFTER:row-1-updated",
+            "DELETE:row-2");
+
+    readPipeline.run().waitUntilFinish();
+  }
+
   @Test
   public void testReadChangesRanges() throws Exception {
     File tableDir = tempFolder.newFolder("delta-table-changes-ranges");
@@ -868,7 +1317,7 @@ public class DeltaIOTest {
     // 1. Write parquet files for Version 0 (insert-only commit)
     Row tableRow1 = Row.withSchema(tableSchema).addValues("row-1").build();
     Row tableRow2 = Row.withSchema(tableSchema).addValues("row-2").build();
-    writeAppendCommit(
+    DeltaWriteTestUtils.writeAppendCommit(
         engine,
         tableDir.getAbsolutePath(),
         0L,
@@ -904,7 +1353,7 @@ public class DeltaIOTest {
             .addValues("row-2", "delete", 1L, new Instant(200000000000L))
             .build();
 
-    writeCdcCommit(
+    DeltaWriteTestUtils.writeCdcCommit(
         engine,
         tableDir.getAbsolutePath(),
         1L,
@@ -917,7 +1366,7 @@ public class DeltaIOTest {
 
     // 3. Write parquet files for Version 2 (insert-only commit)
     Row tableRow3 = Row.withSchema(tableSchema).addValues("row-3").build();
-    writeAppendCommit(
+    DeltaWriteTestUtils.writeAppendCommit(
         engine,
         tableDir.getAbsolutePath(),
         2L,
@@ -981,7 +1430,7 @@ public class DeltaIOTest {
     // 1. Write parquet files for Version 0 (insert-only commit)
     Row tableRow1 = Row.withSchema(tableSchema).addValues("row-1").build();
     Row tableRow2 = Row.withSchema(tableSchema).addValues("row-2").build();
-    writeAppendCommit(
+    DeltaWriteTestUtils.writeAppendCommit(
         engine,
         tableDir.getAbsolutePath(),
         0L,
@@ -1017,7 +1466,7 @@ public class DeltaIOTest {
             .addValues("row-2", "delete", 1L, new Instant(200000000000L))
             .build();
 
-    writeCdcCommit(
+    DeltaWriteTestUtils.writeCdcCommit(
         engine,
         tableDir.getAbsolutePath(),
         1L,
@@ -1030,7 +1479,7 @@ public class DeltaIOTest {
 
     // 3. Write parquet files for Version 2 (insert-only commit)
     Row tableRow3 = Row.withSchema(tableSchema).addValues("row-3").build();
-    writeAppendCommit(
+    DeltaWriteTestUtils.writeAppendCommit(
         engine,
         tableDir.getAbsolutePath(),
         2L,
@@ -1052,7 +1501,7 @@ public class DeltaIOTest {
             .addValues("row-1-updated", "delete", 3L, new 
Instant(400000000000L))
             .build();
 
-    writeCdcCommit(
+    DeltaWriteTestUtils.writeCdcCommit(
         engine,
         tableDir.getAbsolutePath(),
         3L,
@@ -1065,7 +1514,7 @@ public class DeltaIOTest {
 
     // 5. Write parquet files for Version 4 (insert-only commit)
     Row tableRow4 = Row.withSchema(tableSchema).addValues("row-4").build();
-    writeAppendCommit(
+    DeltaWriteTestUtils.writeAppendCommit(
         engine,
         tableDir.getAbsolutePath(),
         4L,
@@ -1106,417 +1555,47 @@ public class DeltaIOTest {
     }
   }
 
-  private List<String> writeAppendCommit(
-      Engine engine,
-      String tablePath,
-      long expectedVersion,
-      long timestamp,
-      StructType deltaSchema,
-      List<Row> beamRows)
-      throws Exception {
-
-    Table table = Table.forPath(engine, tablePath);
-    TransactionBuilder txnBuilder =
-        table.createTransactionBuilder(engine, "DeltaIOTest", Operation.WRITE);
-    if (expectedVersion == 0) {
-      txnBuilder =
-          txnBuilder
-              .withSchema(engine, deltaSchema)
-              .withTableProperties(
-                  engine, 
Collections.singletonMap("delta.enableChangeDataFeed", "true"));
-    }
-    Transaction txn = txnBuilder.build(engine);
-    io.delta.kernel.data.Row txnState = txn.getTransactionState(engine);
-
-    ColumnVector[] vectors = new ColumnVector[deltaSchema.fields().size()];
-    for (int i = 0; i < deltaSchema.fields().size(); i++) {
-      StructField field = deltaSchema.fields().get(i);
-      vectors[i] = createColumnVector(beamRows, i, field.getDataType());
-    }
-
-    ColumnarBatch columnarBatch = new DefaultColumnarBatch(beamRows.size(), 
deltaSchema, vectors);
-    FilteredColumnarBatch filteredBatch =
-        new FilteredColumnarBatch(columnarBatch, Optional.empty());
-
-    CloseableIterator<FilteredColumnarBatch> data =
-        io.delta.kernel.internal.util.Utils.toCloseableIterator(
-            Collections.singletonList(filteredBatch).iterator());
-
-    CloseableIterator<FilteredColumnarBatch> physicalData =
-        Transaction.transformLogicalData(engine, txnState, data, 
Collections.emptyMap());
-
-    DataWriteContext writeContext =
-        Transaction.getWriteContext(engine, txnState, Collections.emptyMap());
-
-    CloseableIterator<DataFileStatus> dataFiles =
-        engine
-            .getParquetHandler()
-            .writeParquetFiles(
-                writeContext.getTargetDirectory(),
-                physicalData,
-                writeContext.getStatisticsColumns());
-
-    List<String> writtenFiles = new ArrayList<>();
-    List<DataFileStatus> filesList = new ArrayList<>();
-    while (dataFiles.hasNext()) {
-      DataFileStatus file = dataFiles.next();
-      filesList.add(file);
-      writtenFiles.add(new File(file.getPath()).getName());
+  private static final class FormatRowWithMetadata extends DoFn<Row, String> {
+    @ProcessElement
+    public void process(@Element Row row, OutputReceiver<String> out) {
+      out.output(
+          String.format(
+              "%s:%s:v%d:t%d",
+              row.getString("name"),
+              row.getString(DeltaIO.CHANGE_TYPE_COLUMN),
+              row.getInt64(DeltaIO.COMMIT_VERSION_COLUMN),
+              row.getDateTime(DeltaIO.COMMIT_TIMESTAMP_COLUMN).getMillis()));
     }
-    CloseableIterator<DataFileStatus> dataFilesCopy =
-        
io.delta.kernel.internal.util.Utils.toCloseableIterator(filesList.iterator());
-
-    CloseableIterator<io.delta.kernel.data.Row> dataActions =
-        Transaction.generateAppendActions(engine, txnState, dataFilesCopy, 
writeContext);
-
-    TransactionCommitResult result =
-        txn.commit(engine, CloseableIterable.inMemoryIterable(dataActions));
-    org.junit.Assert.assertEquals(expectedVersion, result.getVersion());
-    File commitFile =
-        new File(new File(tablePath, "_delta_log"), 
String.format("%020d.json", expectedVersion));
-    commitFile.setLastModified(timestamp);
-    return writtenFiles;
   }
 
-  private void writeCdcCommit(
-      Engine engine,
-      String tablePath,
-      long expectedVersion,
-      long timestamp,
-      StructType deltaSchema,
-      @Nullable List<Row> addBeamRows,
-      @Nullable String removePath,
-      @Nullable List<Row> cdcBeamRows,
-      StructType cdcWriteSchema)
-      throws Exception {
-
-    Table table = Table.forPath(engine, tablePath);
-    TransactionBuilder txnBuilder =
-        table.createTransactionBuilder(engine, "DeltaIOTest", Operation.WRITE);
-    Transaction txn = txnBuilder.build(engine);
-    io.delta.kernel.data.Row txnState = txn.getTransactionState(engine);
-
-    StructType customSingleActionSchema = getCustomSingleActionSchema();
-    List<io.delta.kernel.data.Row> commitActions = new ArrayList<>();
-
-    if (addBeamRows != null && !addBeamRows.isEmpty()) {
-      ColumnVector[] vectors = new ColumnVector[deltaSchema.fields().size()];
-      for (int i = 0; i < deltaSchema.fields().size(); i++) {
-        StructField field = deltaSchema.fields().get(i);
-        vectors[i] = createColumnVector(addBeamRows, i, field.getDataType());
-      }
-      ColumnarBatch columnarBatch =
-          new DefaultColumnarBatch(addBeamRows.size(), deltaSchema, vectors);
-      FilteredColumnarBatch filteredBatch =
-          new FilteredColumnarBatch(columnarBatch, Optional.empty());
-      CloseableIterator<FilteredColumnarBatch> data =
-          io.delta.kernel.internal.util.Utils.toCloseableIterator(
-              Collections.singletonList(filteredBatch).iterator());
-      CloseableIterator<FilteredColumnarBatch> physicalData =
-          Transaction.transformLogicalData(engine, txnState, data, 
Collections.emptyMap());
-      DataWriteContext writeContext =
-          Transaction.getWriteContext(engine, txnState, 
Collections.emptyMap());
-      CloseableIterator<DataFileStatus> dataFiles =
-          engine
-              .getParquetHandler()
-              .writeParquetFiles(
-                  writeContext.getTargetDirectory(),
-                  physicalData,
-                  writeContext.getStatisticsColumns());
-      CloseableIterator<io.delta.kernel.data.Row> addActions =
-          Transaction.generateAppendActions(engine, txnState, dataFiles, 
writeContext);
-      while (addActions.hasNext()) {
-        commitActions.add(addActions.next());
-      }
-    }
-
-    if (removePath != null) {
-      StructType removeSchema =
-          (StructType)
-              io.delta.kernel.internal.actions.SingleAction.FULL_SCHEMA
-                  .fields()
-                  
.get(io.delta.kernel.internal.actions.SingleAction.REMOVE_FILE_ORDINAL)
-                  .getDataType();
-      io.delta.kernel.data.Row removeAction =
-          createRemoveAction(removeSchema, removePath, timestamp);
-      commitActions.add(createSingleAction(customSingleActionSchema, "remove", 
removeAction));
-    }
-
-    if (cdcBeamRows != null && !cdcBeamRows.isEmpty()) {
-      ColumnVector[] vectors = new 
ColumnVector[cdcWriteSchema.fields().size()];
-      for (int i = 0; i < cdcWriteSchema.fields().size(); i++) {
-        StructField field = cdcWriteSchema.fields().get(i);
-        vectors[i] = createColumnVector(cdcBeamRows, i, field.getDataType());
-      }
-      ColumnarBatch columnarBatch =
-          new DefaultColumnarBatch(cdcBeamRows.size(), cdcWriteSchema, 
vectors);
-      FilteredColumnarBatch filteredBatch =
-          new FilteredColumnarBatch(columnarBatch, Optional.empty());
-      CloseableIterator<FilteredColumnarBatch> data =
-          io.delta.kernel.internal.util.Utils.toCloseableIterator(
-              Collections.singletonList(filteredBatch).iterator());
-
-      String cdcDir = new File(tablePath, "_change_data").getAbsolutePath();
-
-      CloseableIterator<DataFileStatus> cdcFiles =
-          engine.getParquetHandler().writeParquetFiles(cdcDir, data, 
Collections.emptyList());
-
-      StructType cdcActionSchema = CDC_ACTION_SCHEMA;
-      while (cdcFiles.hasNext()) {
-        DataFileStatus cdcFile = cdcFiles.next();
-        String relativeCdcPath = "_change_data/" + new 
File(cdcFile.getPath()).getName();
-        io.delta.kernel.data.Row cdcAction =
-            createCdcAction(cdcActionSchema, relativeCdcPath, 
cdcFile.getSize());
-        commitActions.add(createSingleAction(customSingleActionSchema, "cdc", 
cdcAction));
-      }
+  private static final class FormatRowSubsetMetadata extends DoFn<Row, String> 
{
+    @ProcessElement
+    public void process(@Element Row row, OutputReceiver<String> out) {
+      out.output(row.getString("name") + ":" + 
row.getString(DeltaIO.CHANGE_TYPE_COLUMN));
     }
-
-    TransactionCommitResult result =
-        txn.commit(
-            engine,
-            CloseableIterable.inMemoryIterable(
-                
io.delta.kernel.internal.util.Utils.toCloseableIterator(commitActions.iterator())));
-    org.junit.Assert.assertEquals(expectedVersion, result.getVersion());
-    File commitFile =
-        new File(new File(tablePath, "_delta_log"), 
String.format("%020d.json", expectedVersion));
-    commitFile.setLastModified(timestamp);
   }
 
-  private static final StructType CDC_ACTION_SCHEMA =
-      new StructType()
-          .add("path", StringType.STRING, false)
-          .add("partitionValues", new MapType(StringType.STRING, 
StringType.STRING, false), false)
-          .add("size", LongType.LONG, false)
-          .add("dataChange", BooleanType.BOOLEAN, false);
-
-  private static StructType getCustomSingleActionSchema() {
-    StructType originalSchema = 
io.delta.kernel.internal.actions.SingleAction.FULL_SCHEMA;
-    List<StructField> fields = new ArrayList<>();
-    for (StructField field : originalSchema.fields()) {
-      if (field.getName().equals("cdc")) {
-        fields.add(new StructField("cdc", CDC_ACTION_SCHEMA, true));
-      } else {
-        fields.add(field);
-      }
+  private static final class FormatRowName extends DoFn<Row, String> {
+    @ProcessElement
+    public void process(@Element Row row, OutputReceiver<String> out) {
+      out.output(row.getString("name"));
     }
-    return new StructType(fields);
-  }
-
-  private static io.delta.kernel.data.Row createSingleAction(
-      StructType customSingleActionSchema, String actionName, 
io.delta.kernel.data.Row actionRow) {
-    Map<String, Object> values = new HashMap<>();
-    values.put(actionName, actionRow);
-    return new TestRow(customSingleActionSchema, values);
-  }
-
-  private static final MapValue EMPTY_MAP_VALUE =
-      new MapValue() {
-        @Override
-        public int getSize() {
-          return 0;
-        }
-
-        @Override
-        public ColumnVector getKeys() {
-          return new ColumnVector() {
-            @Override
-            public DataType getDataType() {
-              return StringType.STRING;
-            }
-
-            @Override
-            public int getSize() {
-              return 0;
-            }
-
-            @Override
-            public void close() {}
-
-            @Override
-            public boolean isNullAt(int rowId) {
-              return true;
-            }
-          };
-        }
-
-        @Override
-        public ColumnVector getValues() {
-          return new ColumnVector() {
-            @Override
-            public DataType getDataType() {
-              return StringType.STRING;
-            }
-
-            @Override
-            public int getSize() {
-              return 0;
-            }
-
-            @Override
-            public void close() {}
-
-            @Override
-            public boolean isNullAt(int rowId) {
-              return true;
-            }
-          };
-        }
-      };
-
-  private static io.delta.kernel.data.Row createRemoveAction(
-      StructType removeSchema, String path, long deletionTimestamp) {
-    Map<String, Object> values = new HashMap<>();
-    values.put("path", path);
-    values.put("deletionTimestamp", deletionTimestamp);
-    values.put("dataChange", true);
-    values.put("size", 100L);
-    return new TestRow(removeSchema, values);
-  }
-
-  private static io.delta.kernel.data.Row createCdcAction(
-      StructType cdcSchema, String path, long size) {
-    Map<String, Object> values = new HashMap<>();
-    values.put("path", path);
-    values.put("partitionValues", EMPTY_MAP_VALUE);
-    values.put("size", size);
-    values.put("dataChange", true);
-    return new TestRow(cdcSchema, values);
-  }
-
-  private static ColumnVector createColumnVector(
-      List<Row> rows, int fieldIndex, DataType dataType) {
-    return new ColumnVector() {
-      @Override
-      public DataType getDataType() {
-        return dataType;
-      }
-
-      @Override
-      public int getSize() {
-        return rows.size();
-      }
-
-      @Override
-      public void close() {}
-
-      @Override
-      public boolean isNullAt(int rowId) {
-        return rows.get(rowId).getValue(fieldIndex) == null;
-      }
-
-      @Override
-      public boolean getBoolean(int rowId) {
-        return rows.get(rowId).getBoolean(fieldIndex);
-      }
-
-      @Override
-      public int getInt(int rowId) {
-        return rows.get(rowId).getInt32(fieldIndex);
-      }
-
-      @Override
-      public long getLong(int rowId) {
-        if (dataType instanceof TimestampType) {
-          org.joda.time.Instant instant = 
rows.get(rowId).getDateTime(fieldIndex).toInstant();
-          return instant.getMillis() * 1000L;
-        }
-        return rows.get(rowId).getInt64(fieldIndex);
-      }
-
-      @Override
-      public String getString(int rowId) {
-        return rows.get(rowId).getString(fieldIndex);
-      }
-    };
   }
 
-  private static class TestRow implements io.delta.kernel.data.Row {
-    private final StructType schema;
-    private final Map<String, Object> values;
-
-    public TestRow(StructType schema, Map<String, Object> values) {
-      this.schema = schema;
-      this.values = values;
-    }
-
-    @Override
-    public StructType getSchema() {
-      return schema;
-    }
-
-    private Object getVal(int ord) {
-      String name = schema.fields().get(ord).getName();
-      return values.get(name);
-    }
-
-    @Override
-    public boolean isNullAt(int ord) {
-      return getVal(ord) == null;
-    }
-
-    @Override
-    public boolean getBoolean(int ord) {
-      return (Boolean) getVal(ord);
-    }
-
-    @Override
-    public byte getByte(int ord) {
-      return (Byte) getVal(ord);
-    }
-
-    @Override
-    public short getShort(int ord) {
-      return (Short) getVal(ord);
-    }
-
-    @Override
-    public int getInt(int ord) {
-      return (Integer) getVal(ord);
-    }
-
-    @Override
-    public long getLong(int ord) {
-      return (Long) getVal(ord);
-    }
-
-    @Override
-    public float getFloat(int ord) {
-      return (Float) getVal(ord);
-    }
-
-    @Override
-    public double getDouble(int ord) {
-      return (Double) getVal(ord);
-    }
-
-    @Override
-    public String getString(int ord) {
-      return (String) getVal(ord);
-    }
-
-    @Override
-    public byte[] getBinary(int ord) {
-      return (byte[]) getVal(ord);
-    }
-
-    @Override
-    public BigDecimal getDecimal(int ord) {
-      return (BigDecimal) getVal(ord);
-    }
-
-    @Override
-    public io.delta.kernel.data.Row getStruct(int ord) {
-      return (io.delta.kernel.data.Row) getVal(ord);
-    }
-
-    @Override
-    public io.delta.kernel.data.ArrayValue getArray(int ord) {
-      return (io.delta.kernel.data.ArrayValue) getVal(ord);
+  private static final class FormatRowVersionMetadata extends DoFn<Row, 
String> {
+    @ProcessElement
+    public void process(@Element Row row, OutputReceiver<String> out) {
+      out.output(row.getString("name") + ":" + 
row.getInt64(DeltaIO.COMMIT_VERSION_COLUMN));
     }
+  }
 
-    @Override
-    public io.delta.kernel.data.MapValue getMap(int ord) {
-      return (io.delta.kernel.data.MapValue) getVal(ord);
+  private static final class FormatRowTimestampMetadata extends DoFn<Row, 
String> {
+    @ProcessElement
+    public void process(@Element Row row, OutputReceiver<String> out) {
+      out.output(
+          row.getString("name")
+              + ":"
+              + row.getDateTime(DeltaIO.COMMIT_TIMESTAMP_COLUMN).getMillis());
     }
   }
 }
diff --git 
a/sdks/java/io/delta/src/test/java/org/apache/beam/sdk/io/delta/DeltaWriteTestUtils.java
 
b/sdks/java/io/delta/src/test/java/org/apache/beam/sdk/io/delta/DeltaWriteTestUtils.java
new file mode 100644
index 00000000000..4ae75bcd47c
--- /dev/null
+++ 
b/sdks/java/io/delta/src/test/java/org/apache/beam/sdk/io/delta/DeltaWriteTestUtils.java
@@ -0,0 +1,371 @@
+/*
+ * 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.sdk.io.delta;
+
+import io.delta.kernel.DataWriteContext;
+import io.delta.kernel.Operation;
+import io.delta.kernel.Table;
+import io.delta.kernel.Transaction;
+import io.delta.kernel.TransactionBuilder;
+import io.delta.kernel.TransactionCommitResult;
+import io.delta.kernel.data.ColumnVector;
+import io.delta.kernel.data.ColumnarBatch;
+import io.delta.kernel.data.FilteredColumnarBatch;
+import io.delta.kernel.defaults.internal.data.DefaultColumnarBatch;
+import io.delta.kernel.engine.Engine;
+import io.delta.kernel.internal.data.GenericRow;
+import io.delta.kernel.types.BooleanType;
+import io.delta.kernel.types.DataType;
+import io.delta.kernel.types.LongType;
+import io.delta.kernel.types.MapType;
+import io.delta.kernel.types.StringType;
+import io.delta.kernel.types.StructField;
+import io.delta.kernel.types.StructType;
+import io.delta.kernel.types.TimestampType;
+import io.delta.kernel.utils.CloseableIterable;
+import io.delta.kernel.utils.CloseableIterator;
+import io.delta.kernel.utils.DataFileStatus;
+import java.io.File;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import javax.annotation.Nullable;
+import org.apache.beam.sdk.values.Row;
+import org.joda.time.Instant;
+
+/** Utility class for writing test commits (appends and CDC actions) to Delta 
tables in tests. */
+final class DeltaWriteTestUtils {
+
+  private DeltaWriteTestUtils() {}
+
+  private static final StructType CDC_ACTION_SCHEMA =
+      new StructType()
+          .add("path", StringType.STRING, false)
+          .add("partitionValues", new MapType(StringType.STRING, 
StringType.STRING, false), false)
+          .add("size", LongType.LONG, false)
+          .add("dataChange", BooleanType.BOOLEAN, false);
+
+  private static StructType getCustomSingleActionSchema() {
+    StructType originalSchema = 
io.delta.kernel.internal.actions.SingleAction.FULL_SCHEMA;
+    List<StructField> fields = new ArrayList<>();
+    for (StructField field : originalSchema.fields()) {
+      if (field.getName().equals("cdc")) {
+        fields.add(new StructField("cdc", CDC_ACTION_SCHEMA, true));
+      } else {
+        fields.add(field);
+      }
+    }
+    return new StructType(fields);
+  }
+
+  private static io.delta.kernel.data.Row createSingleAction(
+      StructType customSingleActionSchema, String actionName, 
io.delta.kernel.data.Row actionRow) {
+    Map<Integer, Object> values = new HashMap<>();
+    values.put(customSingleActionSchema.indexOf(actionName), actionRow);
+    return new GenericRow(customSingleActionSchema, values);
+  }
+
+  private static io.delta.kernel.data.Row createRemoveAction(
+      StructType removeSchema, String path, long deletionTimestamp) {
+    Map<Integer, Object> values = new HashMap<>();
+    values.put(removeSchema.indexOf("path"), path);
+    values.put(removeSchema.indexOf("deletionTimestamp"), deletionTimestamp);
+    values.put(removeSchema.indexOf("dataChange"), true);
+    values.put(removeSchema.indexOf("size"), 100L);
+    return new GenericRow(removeSchema, values);
+  }
+
+  private static io.delta.kernel.data.Row createCdcAction(
+      StructType cdcSchema, String path, long size) {
+    Map<Integer, Object> values = new HashMap<>();
+    values.put(cdcSchema.indexOf("path"), path);
+    values.put(
+        cdcSchema.indexOf("partitionValues"),
+        
io.delta.kernel.internal.util.VectorUtils.stringStringMapValue(Collections.emptyMap()));
+    values.put(cdcSchema.indexOf("size"), size);
+    values.put(cdcSchema.indexOf("dataChange"), true);
+    return new GenericRow(cdcSchema, values);
+  }
+
+  private static ColumnVector createColumnVector(
+      List<Row> rows, int fieldIndex, DataType dataType) {
+    return new ColumnVector() {
+      @Override
+      public DataType getDataType() {
+        return dataType;
+      }
+
+      @Override
+      public int getSize() {
+        return rows.size();
+      }
+
+      @Override
+      public void close() {}
+
+      @Override
+      public boolean isNullAt(int rowId) {
+        return rows.get(rowId).getValue(fieldIndex) == null;
+      }
+
+      @Override
+      public boolean getBoolean(int rowId) {
+        return rows.get(rowId).getBoolean(fieldIndex);
+      }
+
+      @Override
+      public int getInt(int rowId) {
+        return rows.get(rowId).getInt32(fieldIndex);
+      }
+
+      @Override
+      public long getLong(int rowId) {
+        if (dataType instanceof TimestampType) {
+          Instant instant = 
rows.get(rowId).getDateTime(fieldIndex).toInstant();
+          return instant.getMillis() * 1000L;
+        }
+        return rows.get(rowId).getInt64(fieldIndex);
+      }
+
+      @Override
+      public String getString(int rowId) {
+        return rows.get(rowId).getString(fieldIndex);
+      }
+    };
+  }
+
+  /**
+   * Writes a Delta commit containing append actions.
+   *
+   * @param engine the Delta Lake {@link Engine} instance to use
+   * @param tablePath the path of the Delta table to write to
+   * @param expectedVersion the expected version of the commit to be created
+   * @param timestamp the timestamp of the commit file
+   * @param deltaSchema the schema of the Delta table
+   * @param beamRows the rows to write
+   * @return the list of names of the written Parquet data files
+   * @throws Exception if any error occurs during write or commit
+   */
+  static List<String> writeAppendCommit(
+      Engine engine,
+      String tablePath,
+      long expectedVersion,
+      long timestamp,
+      StructType deltaSchema,
+      List<Row> beamRows)
+      throws Exception {
+
+    Table table = Table.forPath(engine, tablePath);
+    TransactionBuilder txnBuilder =
+        table.createTransactionBuilder(engine, "DeltaTestUtils", 
Operation.WRITE);
+    if (expectedVersion == 0) {
+      txnBuilder =
+          txnBuilder
+              .withSchema(engine, deltaSchema)
+              .withTableProperties(
+                  engine, 
Collections.singletonMap("delta.enableChangeDataFeed", "true"));
+    }
+    Transaction txn = txnBuilder.build(engine);
+    io.delta.kernel.data.Row txnState = txn.getTransactionState(engine);
+
+    ColumnVector[] vectors = new ColumnVector[deltaSchema.fields().size()];
+    for (int i = 0; i < deltaSchema.fields().size(); i++) {
+      StructField field = deltaSchema.fields().get(i);
+      vectors[i] = createColumnVector(beamRows, i, field.getDataType());
+    }
+    ColumnarBatch columnarBatch = new DefaultColumnarBatch(beamRows.size(), 
deltaSchema, vectors);
+    FilteredColumnarBatch filteredBatch =
+        new FilteredColumnarBatch(columnarBatch, Optional.empty());
+    CloseableIterator<FilteredColumnarBatch> data =
+        io.delta.kernel.internal.util.Utils.toCloseableIterator(
+            Collections.singletonList(filteredBatch).iterator());
+    CloseableIterator<FilteredColumnarBatch> physicalData =
+        Transaction.transformLogicalData(engine, txnState, data, 
Collections.emptyMap());
+    DataWriteContext writeContext =
+        Transaction.getWriteContext(engine, txnState, Collections.emptyMap());
+    CloseableIterator<DataFileStatus> dataFiles =
+        engine
+            .getParquetHandler()
+            .writeParquetFiles(
+                writeContext.getTargetDirectory(),
+                physicalData,
+                writeContext.getStatisticsColumns());
+
+    List<String> writtenFiles = new ArrayList<>();
+    List<DataFileStatus> filesList = new ArrayList<>();
+    while (dataFiles.hasNext()) {
+      DataFileStatus file = dataFiles.next();
+      filesList.add(file);
+      writtenFiles.add(new File(file.getPath()).getName());
+    }
+    CloseableIterator<DataFileStatus> dataFilesCopy =
+        
io.delta.kernel.internal.util.Utils.toCloseableIterator(filesList.iterator());
+
+    CloseableIterator<io.delta.kernel.data.Row> dataActions =
+        Transaction.generateAppendActions(engine, txnState, dataFilesCopy, 
writeContext);
+
+    TransactionCommitResult result =
+        txn.commit(engine, CloseableIterable.inMemoryIterable(dataActions));
+    org.junit.Assert.assertEquals(expectedVersion, result.getVersion());
+    if (!tablePath.startsWith("gs://")
+        && !tablePath.startsWith("s3://")
+        && !tablePath.startsWith("hdfs://")) {
+      File commitFile =
+          new File(new File(tablePath, "_delta_log"), 
String.format("%020d.json", expectedVersion));
+      commitFile.setLastModified(timestamp);
+    }
+    return writtenFiles;
+  }
+
+  /**
+   * Writes a Delta commit containing CDC actions (simulating updates/deletes).
+   *
+   * <p>Note on why this is manual: In a standard Spark or Flink writer, 
setting the table property
+   * {@code "delta.enableChangeDataFeed" = "true"} automatically instructs the 
engine to compute and
+   * write the change data files to {@code _change_data/} and append the 
{@code cdc} actions to the
+   * commit log whenever DML statements (like UPDATE/DELETE) are executed.
+   *
+   * <p>However, we are using the Delta Lake Kernel API which does not contain 
an SQL execution
+   * engine or a DML parser. Thus, it cannot automatically compute which rows 
were deleted or
+   * updated. To generate a realistic integration test dataset, we must 
manually construct these
+   * change records, write them into the GCS {@code _change_data/} directory 
using the low-level
+   * parquet handler, and manually register them as {@code cdc} actions in the 
committed
+   * transaction.
+   *
+   * @param engine the Delta Lake {@link Engine} instance to use
+   * @param tablePath the path of the Delta table to write to
+   * @param expectedVersion the expected version of the commit to be created
+   * @param timestamp the timestamp of the commit file
+   * @param deltaSchema the schema of the Delta table
+   * @param addBeamRows the optional list of rows to add in this commit
+   * @param removePath the optional path of the file to remove in this commit
+   * @param cdcBeamRows the optional list of CDC rows to write
+   * @param cdcWriteSchema the schema used for writing the CDC files
+   * @throws Exception if any error occurs during write or commit
+   */
+  static void writeCdcCommit(
+      Engine engine,
+      String tablePath,
+      long expectedVersion,
+      long timestamp,
+      StructType deltaSchema,
+      @Nullable List<Row> addBeamRows,
+      @Nullable String removePath,
+      @Nullable List<Row> cdcBeamRows,
+      StructType cdcWriteSchema)
+      throws Exception {
+
+    Table table = Table.forPath(engine, tablePath);
+    TransactionBuilder txnBuilder =
+        table.createTransactionBuilder(engine, "DeltaTestUtils", 
Operation.WRITE);
+    Transaction txn = txnBuilder.build(engine);
+    io.delta.kernel.data.Row txnState = txn.getTransactionState(engine);
+
+    StructType customSingleActionSchema = getCustomSingleActionSchema();
+    List<io.delta.kernel.data.Row> commitActions = new ArrayList<>();
+
+    if (addBeamRows != null && !addBeamRows.isEmpty()) {
+      ColumnVector[] vectors = new ColumnVector[deltaSchema.fields().size()];
+      for (int i = 0; i < deltaSchema.fields().size(); i++) {
+        StructField field = deltaSchema.fields().get(i);
+        vectors[i] = createColumnVector(addBeamRows, i, field.getDataType());
+      }
+      ColumnarBatch columnarBatch =
+          new DefaultColumnarBatch(addBeamRows.size(), deltaSchema, vectors);
+      FilteredColumnarBatch filteredBatch =
+          new FilteredColumnarBatch(columnarBatch, Optional.empty());
+      CloseableIterator<FilteredColumnarBatch> data =
+          io.delta.kernel.internal.util.Utils.toCloseableIterator(
+              Collections.singletonList(filteredBatch).iterator());
+      CloseableIterator<FilteredColumnarBatch> physicalData =
+          Transaction.transformLogicalData(engine, txnState, data, 
Collections.emptyMap());
+      DataWriteContext writeContext =
+          Transaction.getWriteContext(engine, txnState, 
Collections.emptyMap());
+      CloseableIterator<DataFileStatus> dataFiles =
+          engine
+              .getParquetHandler()
+              .writeParquetFiles(
+                  writeContext.getTargetDirectory(),
+                  physicalData,
+                  writeContext.getStatisticsColumns());
+      CloseableIterator<io.delta.kernel.data.Row> addActions =
+          Transaction.generateAppendActions(engine, txnState, dataFiles, 
writeContext);
+      while (addActions.hasNext()) {
+        commitActions.add(addActions.next());
+      }
+    }
+
+    if (removePath != null) {
+      StructType removeSchema =
+          (StructType)
+              io.delta.kernel.internal.actions.SingleAction.FULL_SCHEMA
+                  .fields()
+                  
.get(io.delta.kernel.internal.actions.SingleAction.REMOVE_FILE_ORDINAL)
+                  .getDataType();
+      io.delta.kernel.data.Row removeAction =
+          createRemoveAction(removeSchema, removePath, timestamp);
+      commitActions.add(createSingleAction(customSingleActionSchema, "remove", 
removeAction));
+    }
+
+    if (cdcBeamRows != null && !cdcBeamRows.isEmpty()) {
+      ColumnVector[] vectors = new 
ColumnVector[cdcWriteSchema.fields().size()];
+      for (int i = 0; i < cdcWriteSchema.fields().size(); i++) {
+        StructField field = cdcWriteSchema.fields().get(i);
+        vectors[i] = createColumnVector(cdcBeamRows, i, field.getDataType());
+      }
+      ColumnarBatch columnarBatch =
+          new DefaultColumnarBatch(cdcBeamRows.size(), cdcWriteSchema, 
vectors);
+      FilteredColumnarBatch filteredBatch =
+          new FilteredColumnarBatch(columnarBatch, Optional.empty());
+      CloseableIterator<FilteredColumnarBatch> data =
+          io.delta.kernel.internal.util.Utils.toCloseableIterator(
+              Collections.singletonList(filteredBatch).iterator());
+
+      String cdcDir = new org.apache.hadoop.fs.Path(tablePath, 
"_change_data").toString();
+
+      CloseableIterator<DataFileStatus> cdcFiles =
+          engine.getParquetHandler().writeParquetFiles(cdcDir, data, 
Collections.emptyList());
+
+      StructType cdcActionSchema = CDC_ACTION_SCHEMA;
+      while (cdcFiles.hasNext()) {
+        DataFileStatus cdcFile = cdcFiles.next();
+        String relativeCdcPath = "_change_data/" + new 
File(cdcFile.getPath()).getName();
+        io.delta.kernel.data.Row cdcAction =
+            createCdcAction(cdcActionSchema, relativeCdcPath, 
cdcFile.getSize());
+        commitActions.add(createSingleAction(customSingleActionSchema, "cdc", 
cdcAction));
+      }
+    }
+
+    TransactionCommitResult result =
+        txn.commit(
+            engine,
+            CloseableIterable.inMemoryIterable(
+                
io.delta.kernel.internal.util.Utils.toCloseableIterator(commitActions.iterator())));
+    org.junit.Assert.assertEquals(expectedVersion, result.getVersion());
+    if (!tablePath.startsWith("gs://")
+        && !tablePath.startsWith("s3://")
+        && !tablePath.startsWith("hdfs://")) {
+      File commitFile =
+          new File(new File(tablePath, "_delta_log"), 
String.format("%020d.json", expectedVersion));
+      commitFile.setLastModified(timestamp);
+    }
+  }
+}
diff --git 
a/sdks/java/managed/src/main/java/org/apache/beam/sdk/managed/Managed.java 
b/sdks/java/managed/src/main/java/org/apache/beam/sdk/managed/Managed.java
index 9589992e079..27c647478e1 100644
--- a/sdks/java/managed/src/main/java/org/apache/beam/sdk/managed/Managed.java
+++ b/sdks/java/managed/src/main/java/org/apache/beam/sdk/managed/Managed.java
@@ -95,6 +95,7 @@ public class Managed {
   public static final String ICEBERG = "iceberg";
   public static final String DELTA_LAKE = "delta";
   public static final String ICEBERG_CDC = "iceberg_cdc";
+  public static final String DELTA_LAKE_CDC = "delta_cdc";
   public static final String KAFKA = "kafka";
   public static final String BIGQUERY = "bigquery";
   public static final String POSTGRES = "postgres";
@@ -107,6 +108,8 @@ public class Managed {
           .put(ICEBERG, 
getUrn(ExternalTransforms.ManagedTransforms.Urns.ICEBERG_READ))
           .put(DELTA_LAKE, 
getUrn(ExternalTransforms.ManagedTransforms.Urns.DELTA_LAKE_READ))
           .put(ICEBERG_CDC, 
getUrn(ExternalTransforms.ManagedTransforms.Urns.ICEBERG_CDC_READ))
+          .put(
+              DELTA_LAKE_CDC, 
getUrn(ExternalTransforms.ManagedTransforms.Urns.DELTA_LAKE_CDC_READ))
           .put(KAFKA, 
getUrn(ExternalTransforms.ManagedTransforms.Urns.KAFKA_READ))
           .put(BIGQUERY, 
getUrn(ExternalTransforms.ManagedTransforms.Urns.BIGQUERY_READ))
           .put(POSTGRES, 
getUrn(ExternalTransforms.ManagedTransforms.Urns.POSTGRES_READ))
@@ -134,6 +137,8 @@ public class Managed {
    *       
href="https://beam.apache.org/releases/javadoc/current/org/apache/beam/sdk/io/delta/DeltaIO.html";>DeltaIO</a>
    *   <li>{@link Managed#ICEBERG_CDC} : CDC Read from Apache Iceberg tables 
using <a
    *       
href="https://beam.apache.org/releases/javadoc/current/org/apache/beam/sdk/io/iceberg/IcebergIO.html";>IcebergIO</a>
+   *   <li>{@link Managed#DELTA_LAKE_CDC} : CDC Read from Delta Lake tables 
using <a
+   *       
href="https://beam.apache.org/releases/javadoc/current/org/apache/beam/sdk/io/delta/DeltaIO.html";>DeltaIO</a>
    *   <li>{@link Managed#KAFKA} : Read from Apache Kafka topics using <a
    *       
href="https://beam.apache.org/releases/javadoc/current/org/apache/beam/sdk/io/kafka/KafkaIO.html";>KafkaIO</a>
    *   <li>{@link Managed#BIGQUERY} : Read from GCP BigQuery tables using <a
diff --git a/sdks/standard_expansion_services.yaml 
b/sdks/standard_expansion_services.yaml
index cb617d3eba1..c6e1654bcd1 100644
--- a/sdks/standard_expansion_services.yaml
+++ b/sdks/standard_expansion_services.yaml
@@ -53,6 +53,7 @@
     - 'beam:schematransform:org.apache.beam:iceberg_read:v1'
     - 'beam:schematransform:org.apache.beam:iceberg_cdc_read:v1'
     - 'beam:schematransform:org.apache.beam:delta_lake_read:v1'
+    - 'beam:schematransform:org.apache.beam:delta_lake_cdc_read:v1'
 
 - gradle_target: 'sdks:java:io:messaging-expansion-service:shadowJar'
   destinations:
diff --git a/website/www/site/content/en/documentation/io/managed-io.md 
b/website/www/site/content/en/documentation/io/managed-io.md
index 7239832937a..70b5efe9ed1 100644
--- a/website/www/site/content/en/documentation/io/managed-io.md
+++ b/website/www/site/content/en/documentation/io/managed-io.md
@@ -70,6 +70,21 @@ and Beam SQL is invoked via the Managed API under the hood.
         Unavailable
       </td>
     </tr>
+    <tr>
+      <td><strong>DELTA_CDC</strong></td>
+      <td>
+        <strong>table</strong> (<code style="color: green">str</code>)<br>
+        start_version (<code style="color: #f54251">int64</code>)<br>
+        start_timestamp (<code style="color: green">str</code>)<br>
+        end_version (<code style="color: #f54251">int64</code>)<br>
+        end_timestamp (<code style="color: green">str</code>)<br>
+        hadoop_config (<code>map[<span style="color: green;">str</span>, <span 
style="color: green;">str</span>]</code>)<br>
+        include_metadata_columns (<code>list[<span style="color: 
green;">str</span>]</code>)<br>
+      </td>
+      <td>
+        Unavailable
+      </td>
+    </tr>
     <tr>
       <td><strong>ICEBERG</strong></td>
       <td>
@@ -306,6 +321,95 @@ and Beam SQL is invoked via the Managed API under the hood.
   </table>
 </div>
 
+### `DELTA_CDC` Read
+
+<div class="table-container-wrapper">
+  <table class="table table-bordered">
+    <tr>
+      <th>Configuration</th>
+      <th>Type</th>
+      <th>Description</th>
+    </tr>
+    <tr>
+      <td>
+        <strong>table</strong>
+      </td>
+      <td>
+        <code style="color: green">str</code>
+      </td>
+      <td>
+        Identifier of the Delta Lake table.
+      </td>
+    </tr>
+    <tr>
+      <td>
+        start_version
+      </td>
+      <td>
+        <code style="color: #f54251">int64</code>
+      </td>
+      <td>
+        Start version of the Delta Lake table to read changes from. Either 
start_version or start_timestamp must be set.
+      </td>
+    </tr>
+    <tr>
+      <td>
+        start_timestamp
+      </td>
+      <td>
+        <code style="color: green">str</code>
+      </td>
+      <td>
+        Start timestamp of the Delta Lake table to read changes from. Either 
start_version or start_timestamp must be set.
+      </td>
+    </tr>
+    <tr>
+      <td>
+        end_version
+      </td>
+      <td>
+        <code style="color: #f54251">int64</code>
+      </td>
+      <td>
+        End version of the Delta Lake table to read changes up to.
+      </td>
+    </tr>
+    <tr>
+      <td>
+        end_timestamp
+      </td>
+      <td>
+        <code style="color: green">str</code>
+      </td>
+      <td>
+        End timestamp of the Delta Lake table to read changes up to.
+      </td>
+    </tr>
+    <tr>
+      <td>
+        hadoop_config
+      </td>
+      <td>
+        <code>map[<span style="color: green;">str</span>, <span style="color: 
green;">str</span>]</code>
+      </td>
+      <td>
+        Properties passed to the Hadoop Configuration.
+      </td>
+    </tr>
+    <tr>
+      <td>
+        include_metadata_columns
+      </td>
+      <td>
+        <code>list[<span style="color: green;">str</span>]</code>
+      </td>
+      <td>
+        Metadata columns to include in the output rows. Supported columns are: 
_change_type, _commit_version, and _commit_timestamp.
+      </td>
+    </tr>
+  </table>
+</div>
+
 ### `ICEBERG` Read
 
 <div class="table-container-wrapper">

Reply via email to