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">