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 5abe02fc09c7..3a009261f4f9 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": 2 } 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 918455dbdbd2..debacc245d60 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 5ee5442ecd18..66bacf547d16 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 cf10adc9865c..414402429c38 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 @@ @DoFn.BoundedPerElement class DeltaCDCSourceDoFn extends DoFn { @Nullable Map hadoopConfig; + private final @Nullable List metadataColumns; private transient @Nullable Engine engine; private transient @Nullable Configuration conf; - public DeltaCDCSourceDoFn(@Nullable Map hadoopConfig) { + public DeltaCDCSourceDoFn( + @Nullable Map hadoopConfig, @Nullable List metadataColumns) { this.hadoopConfig = hadoopConfig; + this.metadataColumns = metadataColumns; } private synchronized Configuration getConfiguration() { @@ -117,7 +120,8 @@ public void processElement( 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 000000000000..f35a7a52b053 --- /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 { + static final String OUTPUT_TAG = "output"; + + @Override + protected SchemaTransform from(Configuration configuration) { + return new DeltaCdcReadSchemaTransform(configuration); + } + + @Override + public List 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 hadoopConfig = configuration.getHadoopConfig(); + if (hadoopConfig != null) { + read = read.withConfig(hadoopConfig); + } + List includeMetadataColumns = configuration.getIncludeMetadataColumns(); + if (includeMetadataColumns != null && !includeMetadataColumns.isEmpty()) { + read = read.withMetadataColumns(includeMetadataColumns.toArray(new String[0])); + } + + PCollection 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 getHadoopConfig(); + + @SchemaFieldDescription( + "Metadata columns to include in the output rows. Supported columns are: _change_type, _commit_version, and _commit_timestamp.") + @Nullable + abstract List 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 hadoopConfig); + + abstract Builder setIncludeMetadataColumns(List 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 3ac2c7a84a8f..8057332ddce4 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.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 @@ static Schema.FieldType convertToBeamFieldType(DataType deltaType) { } } + static Schema buildPublicBeamSchema(Schema baseSchema, @Nullable List 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> { public abstract @Nullable String getTablePath(); @@ -218,6 +240,8 @@ public abstract static class ReadChanges extends PTransform getMetadataColumns(); + public abstract @Nullable Map getHadoopConfig(); abstract Builder toBuilder(); @@ -234,6 +258,8 @@ abstract static class Builder { abstract Builder setEndTimestamp(@Nullable String endTimestamp); + abstract Builder setMetadataColumns(@Nullable List metadataColumns); + abstract Builder setHadoopConfig(@Nullable Map hadoopConfig); abstract ReadChanges build(); @@ -259,6 +285,20 @@ public ReadChanges withEndTimestamp(String endTimestamp) { 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 config) { return toBuilder().setHadoopConfig(config).build(); } @@ -310,7 +350,8 @@ public PCollection expand(PBegin input) { 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 PCollection expand(PBegin input) { 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 e0d35f30faa6..0e6d074e0428 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.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,16 @@ 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 +83,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 +134,11 @@ public void setup() throws Exception { 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 +220,33 @@ public String getString(int rowId) { CloseableIterator dataActions = Transaction.generateAppendActions(engine, txnState, dataFiles, writeContext); + List 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 dataActionsIterable = - CloseableIterable.inMemoryIterable(dataActions); + CloseableIterable.inMemoryIterable( + io.delta.kernel.internal.util.Utils.toCloseableIterator(addActionsList.iterator())); TransactionCommitResult commitResult = txn.commit(engine, dataActionsIterable); @@ -236,12 +272,56 @@ public void teardown() { } } + // @Test + // public void testReadDeltaLakeTable() { + // ExperimentalOptions options = + // readPipeline.getOptions().as(ExperimentalOptions.class); + // ExperimentalOptions.addExperiment(options, "use_runner_v2"); + + // Map 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"); + // 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); + // } + + // PCollection output = + // readPipeline + // .apply( + // Managed.read(Managed.DELTA_LAKE) + // .withConfig(ImmutableMap.of("table", repoPath, "hadoop_config", + // hadoopConfig))) + // .getSinglePCollection(); + + // PAssert.that(output).containsInAnyOrder(TEST_ROWS); + // readPipeline.run().waitUntilFinish(); + // } + @Test - public void testReadDeltaLakeTable() { + public void testReadChangesDeltaLake() throws Exception { + ExperimentalOptions options = readPipeline.getOptions().as(ExperimentalOptions.class); + List experiments = options.getExperiments(); + if (experiments != null) { + List 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 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() @@ -251,14 +331,101 @@ public void testReadDeltaLakeTable() { hadoopConfig.put("fs.gs.project.id", project); } + org.apache.hadoop.conf.Configuration conf = new org.apache.hadoop.conf.Configuration(); + for (Map.Entry 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 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 output = readPipeline - .apply( - Managed.read(Managed.DELTA_LAKE) - .withConfig(ImmutableMap.of("table", repoPath, "hadoop_config", hadoopConfig))) + .apply(Managed.read(Managed.DELTA_LAKE_CDC).withConfig(readConfig)) .getSinglePCollection(); - PAssert.that(output).containsInAnyOrder(TEST_ROWS); + PCollection formattedOutput = + output.apply("Format Row with Metadata", ParDo.of(new FormatITRowWithMetadata())); + + // Generate expected outputs for version 0 (inserts of id 0-99) + List 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 { + @ProcessElement + public void process(@Element Row row, OutputReceiver 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 f00b34be4609..97534edf79e3 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.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.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 void testManagedDeltaRead() throws Exception { 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 void testReadChanges() throws Exception { 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 void testReadChanges() throws Exception { .addValues("row-2", "delete", 1L, new Instant(123456789000L)) .build(); - writeCdcCommit( + DeltaWriteTestUtils.writeCdcCommit( engine, tableDir.getAbsolutePath(), 1L, @@ -857,6 +838,474 @@ public void testReadChanges() throws Exception { 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 outputCDC = + readPipeline.apply( + "Read Changes", + DeltaIO.readChanges().from(tableDir.getAbsolutePath()).withStartVersion(0L)); + + PCollection 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 outputNormal = + writePipeline.apply("Read Normal", DeltaIO.readRows().from(tableDir.getAbsolutePath())); + + PCollection 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 output = + PCollectionRowTuple.empty(readPipeline) + .apply(new DeltaCdcReadSchemaTransformProvider().from(config)) + .get(DeltaCdcReadSchemaTransformProvider.OUTPUT_TAG); + + PCollection 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 output = + PCollectionRowTuple.empty(readPipeline) + .apply(new DeltaCdcReadSchemaTransformProvider().from(config)) + .get(DeltaCdcReadSchemaTransformProvider.OUTPUT_TAG); + + PCollection 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 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 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 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 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 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 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 config = new HashMap<>(); + config.put("table", tableDir.getAbsolutePath()); + config.put("start_version", 0L); + + PCollection output = + readPipeline + .apply(Managed.read(Managed.DELTA_LAKE_CDC).withConfig(config)) + .getSinglePCollection(); + + PCollection 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 void testReadChangesRanges() throws Exception { // 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 void testReadChangesRanges() throws Exception { .addValues("row-2", "delete", 1L, new Instant(200000000000L)) .build(); - writeCdcCommit( + DeltaWriteTestUtils.writeCdcCommit( engine, tableDir.getAbsolutePath(), 1L, @@ -917,7 +1366,7 @@ public void testReadChangesRanges() throws Exception { // 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 void testReadChangesPartialRange() throws Exception { // 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 void testReadChangesPartialRange() throws Exception { .addValues("row-2", "delete", 1L, new Instant(200000000000L)) .build(); - writeCdcCommit( + DeltaWriteTestUtils.writeCdcCommit( engine, tableDir.getAbsolutePath(), 1L, @@ -1030,7 +1479,7 @@ public void testReadChangesPartialRange() throws Exception { // 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 void testReadChangesPartialRange() throws Exception { .addValues("row-1-updated", "delete", 3L, new Instant(400000000000L)) .build(); - writeCdcCommit( + DeltaWriteTestUtils.writeCdcCommit( engine, tableDir.getAbsolutePath(), 3L, @@ -1065,7 +1514,7 @@ public void testReadChangesPartialRange() throws Exception { // 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 void process( } } - private List writeAppendCommit( - Engine engine, - String tablePath, - long expectedVersion, - long timestamp, - StructType deltaSchema, - List 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 data = - io.delta.kernel.internal.util.Utils.toCloseableIterator( - Collections.singletonList(filteredBatch).iterator()); - - CloseableIterator physicalData = - Transaction.transformLogicalData(engine, txnState, data, Collections.emptyMap()); - - DataWriteContext writeContext = - Transaction.getWriteContext(engine, txnState, Collections.emptyMap()); - - CloseableIterator dataFiles = - engine - .getParquetHandler() - .writeParquetFiles( - writeContext.getTargetDirectory(), - physicalData, - writeContext.getStatisticsColumns()); - - List writtenFiles = new ArrayList<>(); - List 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 { + @ProcessElement + public void process(@Element Row row, OutputReceiver 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 dataFilesCopy = - io.delta.kernel.internal.util.Utils.toCloseableIterator(filesList.iterator()); - - CloseableIterator 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 addBeamRows, - @Nullable String removePath, - @Nullable List 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 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 data = - io.delta.kernel.internal.util.Utils.toCloseableIterator( - Collections.singletonList(filteredBatch).iterator()); - CloseableIterator physicalData = - Transaction.transformLogicalData(engine, txnState, data, Collections.emptyMap()); - DataWriteContext writeContext = - Transaction.getWriteContext(engine, txnState, Collections.emptyMap()); - CloseableIterator dataFiles = - engine - .getParquetHandler() - .writeParquetFiles( - writeContext.getTargetDirectory(), - physicalData, - writeContext.getStatisticsColumns()); - CloseableIterator 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 data = - io.delta.kernel.internal.util.Utils.toCloseableIterator( - Collections.singletonList(filteredBatch).iterator()); - - String cdcDir = new File(tablePath, "_change_data").getAbsolutePath(); - - CloseableIterator 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 { + @ProcessElement + public void process(@Element Row row, OutputReceiver 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 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 { + @ProcessElement + public void process(@Element Row row, OutputReceiver 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 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 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 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 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 values; - - public TestRow(StructType schema, Map 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 { + @ProcessElement + public void process(@Element Row row, OutputReceiver 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 { + @ProcessElement + public void process(@Element Row row, OutputReceiver 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 000000000000..4ae75bcd47cd --- /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 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 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 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 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 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 writeAppendCommit( + Engine engine, + String tablePath, + long expectedVersion, + long timestamp, + StructType deltaSchema, + List 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 data = + io.delta.kernel.internal.util.Utils.toCloseableIterator( + Collections.singletonList(filteredBatch).iterator()); + CloseableIterator physicalData = + Transaction.transformLogicalData(engine, txnState, data, Collections.emptyMap()); + DataWriteContext writeContext = + Transaction.getWriteContext(engine, txnState, Collections.emptyMap()); + CloseableIterator dataFiles = + engine + .getParquetHandler() + .writeParquetFiles( + writeContext.getTargetDirectory(), + physicalData, + writeContext.getStatisticsColumns()); + + List writtenFiles = new ArrayList<>(); + List filesList = new ArrayList<>(); + while (dataFiles.hasNext()) { + DataFileStatus file = dataFiles.next(); + filesList.add(file); + writtenFiles.add(new File(file.getPath()).getName()); + } + CloseableIterator dataFilesCopy = + io.delta.kernel.internal.util.Utils.toCloseableIterator(filesList.iterator()); + + CloseableIterator 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). + * + *

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. + * + *

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 addBeamRows, + @Nullable String removePath, + @Nullable List 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 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 data = + io.delta.kernel.internal.util.Utils.toCloseableIterator( + Collections.singletonList(filteredBatch).iterator()); + CloseableIterator physicalData = + Transaction.transformLogicalData(engine, txnState, data, Collections.emptyMap()); + DataWriteContext writeContext = + Transaction.getWriteContext(engine, txnState, Collections.emptyMap()); + CloseableIterator dataFiles = + engine + .getParquetHandler() + .writeParquetFiles( + writeContext.getTargetDirectory(), + physicalData, + writeContext.getStatisticsColumns()); + CloseableIterator 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 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 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 9589992e079a..27c647478e17 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 *

  • {@link Managed#ICEBERG_CDC} : CDC Read from Apache Iceberg tables using IcebergIO + *
  • {@link Managed#DELTA_LAKE_CDC} : CDC Read from Delta Lake tables using DeltaIO *
  • {@link Managed#KAFKA} : Read from Apache Kafka topics using KafkaIO *
  • {@link Managed#BIGQUERY} : Read from GCP BigQuery tables using + + DELTA_CDC + + table (str)
    + start_version (int64)
    + start_timestamp (str)
    + end_version (int64)
    + end_timestamp (str)
    + hadoop_config (map[str, str])
    + include_metadata_columns (list[str])
    + + + Unavailable + + ICEBERG @@ -306,6 +321,95 @@ and Beam SQL is invoked via the Managed API under the hood. +### `DELTA_CDC` Read + +
    + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + +
    ConfigurationTypeDescription
    + table + + str + + Identifier of the Delta Lake table. +
    + start_version + + int64 + + Start version of the Delta Lake table to read changes from. Either start_version or start_timestamp must be set. +
    + start_timestamp + + str + + Start timestamp of the Delta Lake table to read changes from. Either start_version or start_timestamp must be set. +
    + end_version + + int64 + + End version of the Delta Lake table to read changes up to. +
    + end_timestamp + + str + + End timestamp of the Delta Lake table to read changes up to. +
    + hadoop_config + + map[str, str] + + Properties passed to the Hadoop Configuration. +
    + include_metadata_columns + + list[str] + + Metadata columns to include in the output rows. Supported columns are: _change_type, _commit_version, and _commit_timestamp. +
    +
    + ### `ICEBERG` Read