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 ce45298a609 Fixes to delta CDC read (#39713)
ce45298a609 is described below
commit ce45298a609064e6d5d4264088c2e658a028a22a
Author: Chamikara Jayalath <[email protected]>
AuthorDate: Tue Aug 11 10:41:51 2026 -0700
Fixes to delta CDC read (#39713)
---
.../beam/sdk/io/delta/DeltaCDCSourceDoFn.java | 24 +++---
.../java/org/apache/beam/sdk/io/delta/DeltaIO.java | 4 +-
.../org/apache/beam/sdk/io/delta/DeltaIOTest.java | 98 ++++++++++++++++++++++
3 files changed, 115 insertions(+), 11 deletions(-)
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 414402429c3..97d9c914a08 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
@@ -231,7 +231,7 @@ class DeltaCDCSourceDoFn extends DoFn<DeltaCDCReadTask,
Row> {
"Field " + DeltaIO.CHANGE_TYPE_COLUMN + " must not be
null.");
}
ValueKind kind = getValueKind(changeType);
- Row publicRow = projectRow(beamRow, publicBeamSchema);
+ Row publicRow = projectRow(beamRow, publicBeamSchema, task);
out.builder(publicRow).setValueKind(kind).output();
}
}
@@ -240,14 +240,18 @@ class DeltaCDCSourceDoFn extends DoFn<DeltaCDCReadTask,
Row> {
}
}
- private static Row projectRow(Row row, Schema targetSchema) {
- if (row.getSchema().equals(targetSchema)) {
- // We can return the original Row since schemas are the same.
- return row;
- }
+ private static Row projectRow(Row row, Schema targetSchema, DeltaCDCReadTask
task) {
Row.Builder builder = Row.withSchema(targetSchema);
for (Schema.Field field : targetSchema.getFields()) {
- builder.addValue(row.getValue(field.getName()));
+ Object value = row.getValue(field.getName());
+ if (value == null) {
+ if (field.getName().equals(DeltaIO.COMMIT_VERSION_COLUMN)) {
+ value = task.getVersion();
+ } else if (field.getName().equals(DeltaIO.COMMIT_TIMESTAMP_COLUMN)) {
+ value = new org.joda.time.Instant(task.getTimestamp());
+ }
+ }
+ builder.addValue(value);
}
return builder.build();
}
@@ -271,9 +275,9 @@ class DeltaCDCSourceDoFn extends DoFn<DeltaCDCReadTask,
Row> {
private static StructType appendCDFColumns(StructType schema) {
return schema
- .add(DeltaIO.CHANGE_TYPE_COLUMN, StringType.STRING, false)
- .add(DeltaIO.COMMIT_VERSION_COLUMN, LongType.LONG, false)
- .add(DeltaIO.COMMIT_TIMESTAMP_COLUMN, TimestampType.TIMESTAMP, false);
+ .add(DeltaIO.CHANGE_TYPE_COLUMN, StringType.STRING, true)
+ .add(DeltaIO.COMMIT_VERSION_COLUMN, LongType.LONG, true)
+ .add(DeltaIO.COMMIT_TIMESTAMP_COLUMN, TimestampType.TIMESTAMP, true);
}
private ColumnarBatch appendConstantCDFColumns(
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 8057332ddce..3a53b5c7620 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
@@ -168,7 +168,9 @@ public class DeltaIO {
static Schema convertToBeamSchema(StructType deltaSchema) {
Schema.Builder builder = Schema.builder();
for (StructField field : deltaSchema.fields()) {
- builder.addField(field.getName(),
convertToBeamFieldType(field.getDataType()));
+ builder.addField(
+ Schema.Field.of(field.getName(),
convertToBeamFieldType(field.getDataType()))
+ .withNullable(field.isNullable()));
}
return builder.build();
}
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 b9b2a604a4b..0db0aef9e08 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
@@ -451,6 +451,26 @@ public class DeltaIOTest {
org.junit.Assert.assertEquals(expectedSchema, actualSchema);
}
+ @Test
+ public void testConvertToBeamSchemaPreservesNullability() {
+ StructType deltaSchema =
+ new StructType(
+ java.util.Arrays.asList(
+ new StructField("nullable_string", StringType.STRING, true),
+ new StructField("non_nullable_integer", IntegerType.INTEGER,
false)));
+
+ Schema expectedSchema =
+ Schema.builder()
+ .addField(
+ Schema.Field.of("nullable_string",
Schema.FieldType.STRING).withNullable(true))
+ .addField(
+ Schema.Field.of("non_nullable_integer",
Schema.FieldType.INT32).withNullable(false))
+ .build();
+
+ Schema actualSchema = DeltaIO.ReadRows.convertToBeamSchema(deltaSchema);
+ org.junit.Assert.assertEquals(expectedSchema, actualSchema);
+ }
+
@Test
public void testDeltaReadTaskTracker() {
java.util.List<Long> sizes = java.util.Arrays.asList(100L, 200L, 300L);
@@ -1090,6 +1110,84 @@ public class DeltaIOTest {
readPipeline.run().waitUntilFinish();
}
+ @Test
+ public void testReadChangesWithMissingMetadataColumns() throws Exception {
+ File tableDir =
tempFolder.newFolder("delta-table-changes-missing-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), but
OMIT version and
+ // timestamp columns!
+ Schema cdcWriteSchema =
+ Schema.builder()
+ .addField("name", Schema.FieldType.STRING)
+ .addField(DeltaIO.CHANGE_TYPE_COLUMN, Schema.FieldType.STRING)
+ .build();
+ StructType cdcWriteDeltaSchema =
+ new StructType()
+ .add("name", StringType.STRING)
+ .add(DeltaIO.CHANGE_TYPE_COLUMN, StringType.STRING);
+
+ Row cdcRow1 = Row.withSchema(cdcWriteSchema).addValues("row-1",
"update_preimage").build();
+ Row cdcRow2 =
+ Row.withSchema(cdcWriteSchema).addValues("row-1-updated",
"update_postimage").build();
+ Row cdcRow3 = Row.withSchema(cdcWriteSchema).addValues("row-2",
"delete").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 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:t200000000000",
+ "row-1-updated:update_postimage:v1:t200000000000",
+ "row-2:delete:v1:t200000000000");
+
+ readPipeline.run().waitUntilFinish();
+ }
+
@Test
public void testReadChangesWithSubsetOfMetadataColumns() throws Exception {
File tableDir =
tempFolder.newFolder("delta-table-changes-subset-metadata");