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 f65e0e0b0a8 Updates the Delta Lake source to support reading bounded
change data (#39426)
f65e0e0b0a8 is described below
commit f65e0e0b0a8e0e17e18080b703ddc45bb5970abb
Author: Chamikara Jayalath <[email protected]>
AuthorDate: Thu Jul 30 08:47:23 2026 -0700
Updates the Delta Lake source to support reading bounded change data
(#39426)
---
.../beam/sdk/io/delta/CreateCDCReadTasksDoFn.java | 296 ++++++++
.../apache/beam/sdk/io/delta/DeltaCDCReadTask.java | 125 +++
.../beam/sdk/io/delta/DeltaCDCSourceDoFn.java | 359 +++++++++
.../java/org/apache/beam/sdk/io/delta/DeltaIO.java | 150 +++-
.../io/delta/DeltaReadSchemaTransformProvider.java | 6 +-
.../apache/beam/sdk/io/delta/DeltaSourceDoFn.java | 2 +-
.../org/apache/beam/sdk/io/delta/DeltaIOTest.java | 841 +++++++++++++++++++--
7 files changed, 1692 insertions(+), 87 deletions(-)
diff --git
a/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/CreateCDCReadTasksDoFn.java
b/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/CreateCDCReadTasksDoFn.java
new file mode 100644
index 00000000000..9a0049a8651
--- /dev/null
+++
b/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/CreateCDCReadTasksDoFn.java
@@ -0,0 +1,296 @@
+/*
+ * 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.CommitRange;
+import io.delta.kernel.CommitRangeBuilder;
+import io.delta.kernel.Scan;
+import io.delta.kernel.Snapshot;
+import io.delta.kernel.Table;
+import io.delta.kernel.TableManager;
+import io.delta.kernel.data.ColumnarBatch;
+import io.delta.kernel.data.Row;
+import io.delta.kernel.defaults.engine.DefaultEngine;
+import io.delta.kernel.engine.Engine;
+import io.delta.kernel.internal.DeltaLogActionUtils.DeltaAction;
+import io.delta.kernel.internal.TableImpl;
+import io.delta.kernel.internal.actions.AddCDCFile;
+import io.delta.kernel.internal.actions.AddFile;
+import io.delta.kernel.internal.util.VectorUtils;
+import io.delta.kernel.utils.CloseableIterator;
+import java.time.Instant;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import org.apache.beam.sdk.transforms.DoFn;
+import org.apache.hadoop.conf.Configuration;
+import org.checkerframework.checker.nullness.qual.Nullable;
+
+/** A DoFn that reads the Delta log and plans read tasks for Change Data Feed.
*/
+class CreateCDCReadTasksDoFn extends DoFn<String, DeltaCDCReadTask> {
+ private static final long MAX_TASK_SIZE_BYTES = 1024L * 1024L * 1024L; // 1
GB
+ private final @Nullable Map<String, String> hadoopConfig;
+ private final @Nullable Long startVersion;
+ private final @Nullable String startTimestamp;
+ private final @Nullable Long endVersion;
+ private final @Nullable String endTimestamp;
+
+ public CreateCDCReadTasksDoFn(
+ @Nullable Map<String, String> hadoopConfig,
+ @Nullable Long startVersion,
+ @Nullable String startTimestamp,
+ @Nullable Long endVersion,
+ @Nullable String endTimestamp) {
+ this.hadoopConfig = hadoopConfig;
+ this.startVersion = startVersion;
+ this.startTimestamp = startTimestamp;
+ this.endVersion = endVersion;
+ this.endTimestamp = endTimestamp;
+ }
+
+ @ProcessElement
+ public void processElement(@Element String tablePath,
OutputReceiver<DeltaCDCReadTask> out)
+ throws Exception {
+ Configuration conf = new Configuration();
+ if (hadoopConfig != null) {
+ for (Map.Entry<String, String> entry : hadoopConfig.entrySet()) {
+ conf.set(entry.getKey(), entry.getValue());
+ }
+ }
+ Engine engine = DefaultEngine.create(conf);
+ Table table = Table.forPath(engine, tablePath);
+ TableImpl tableImpl = (TableImpl) table;
+
+ // 1. Resolve starting and ending versions
+ long resolvedStartVersion;
+ if (startVersion != null) {
+ resolvedStartVersion = startVersion;
+ } else if (startTimestamp != null) {
+ long startMillis = Instant.parse(startTimestamp).toEpochMilli();
+ resolvedStartVersion = tableImpl.getVersionAtOrAfterTimestamp(engine,
startMillis);
+ } else {
+ throw new IllegalArgumentException("Starting version or timestamp must
be specified.");
+ }
+
+ long resolvedEndVersion;
+ if (endVersion != null) {
+ resolvedEndVersion = endVersion;
+ } else if (endTimestamp != null) {
+ long endMillis = Instant.parse(endTimestamp).toEpochMilli();
+ resolvedEndVersion = tableImpl.getVersionBeforeOrAtTimestamp(engine,
endMillis);
+ } else {
+ resolvedEndVersion = table.getLatestSnapshot(engine).getVersion();
+ }
+
+ if (resolvedStartVersion > resolvedEndVersion) {
+ throw new IllegalArgumentException(
+ String.format(
+ "Resolved start version %d is greater than resolved end version
%d",
+ resolvedStartVersion, resolvedEndVersion));
+ }
+
+ // 2. Load snapshot at resolvedEndVersion to get the scanStateRow
+ // We use endVersion's schema because it represents the latest schema in
the
+ // read range
+ // which handles schema evolution (older files will just lack new columns).
+ Snapshot endSnapshot = table.getSnapshotAsOfVersion(engine,
resolvedEndVersion);
+ Scan scan = endSnapshot.getScanBuilder().build();
+ Row scanState = scan.getScanState(engine);
+ SerializableRow serializableScanState = new SerializableRow(scanState);
+
+ // 3. Load snapshot at resolvedStartVersion to initialize the CommitRange
+ Snapshot startSnapshot = table.getSnapshotAsOfVersion(engine,
resolvedStartVersion);
+
+ CommitRangeBuilder rangeBuilder =
+ TableManager.loadCommitRange(
+ tablePath,
CommitRangeBuilder.CommitBoundary.atVersion(resolvedStartVersion));
+
rangeBuilder.withEndBoundary(CommitRangeBuilder.CommitBoundary.atVersion(resolvedEndVersion));
+ CommitRange range = rangeBuilder.build(engine);
+
+ // We need both CDC and ADD actions.
+ // If a commit version has CDC files, we only read CDC files.
+ // If a commit version has no CDC files, we read ADD files (inserts).
+ Set<DeltaAction> actionSet = new HashSet<>();
+ actionSet.add(DeltaAction.CDC);
+ actionSet.add(DeltaAction.ADD);
+
+ // 4. Iterate over commits in the range and group actions by version
+ try (CloseableIterator<ColumnarBatch> batchIter =
+ range.getActions(engine, startSnapshot, actionSet)) {
+ Map<Long, CommitActionsInfo> commitActionsMap = new HashMap<>();
+
+ while (batchIter.hasNext()) {
+ ColumnarBatch batch = batchIter.next();
+ int versionIdx = batch.getSchema().indexOf("version");
+ int timestampIdx = batch.getSchema().indexOf("timestamp");
+ int cdcIdx = batch.getSchema().indexOf("cdc");
+ int addIdx = batch.getSchema().indexOf("add");
+
+ for (int i = 0; i < batch.getSize(); i++) {
+ long version = batch.getColumnVector(versionIdx).getLong(i);
+ long timestamp = batch.getColumnVector(timestampIdx).getLong(i);
+
+ CommitActionsInfo info =
+ commitActionsMap.computeIfAbsent(
+ version, k -> new CommitActionsInfo(version, timestamp));
+
+ if (cdcIdx >= 0 && !batch.getColumnVector(cdcIdx).isNullAt(i)) {
+ Row cdcRow =
+ (Row)
+ VectorUtils.getValueAsObject(
+ batch.getColumnVector(cdcIdx),
+ batch.getSchema().at(cdcIdx).getDataType(),
+ i);
+ info.cdcInfo.add(cdcRow);
+ }
+ if (addIdx >= 0 && !batch.getColumnVector(addIdx).isNullAt(i)) {
+ Row addRow =
+ (Row)
+ VectorUtils.getValueAsObject(
+ batch.getColumnVector(addIdx),
+ batch.getSchema().at(addIdx).getDataType(),
+ i);
+ AddFile addFile = new AddFile(addRow);
+ // Only consider add files that change data (ignore OPTIMIZE etc.)
+ if (addFile.getDataChange()) {
+ info.insertInfo.add(addRow);
+ }
+ }
+ }
+ }
+
+ // 5. Emit tasks for each version
+ List<DeltaCDCReadTask> currentGroup = new ArrayList<>();
+ long currentGroupSize = 0L;
+
+ // TODO: to prevent OOMs in true streaming executions, update
DeltaReadTask to include a group
+ // of files.
+
+ // Sort versions to process them in order
+ List<Long> versions = new ArrayList<>(commitActionsMap.keySet());
+ Collections.sort(versions);
+
+ for (long version : versions) {
+ CommitActionsInfo info = commitActionsMap.get(version);
+ if (info == null) {
+ throw new IllegalStateException("CommitActionsInfo was not found for
version " + version);
+ }
+ boolean hasCDC = !info.cdcInfo.isEmpty();
+
+ List<Row> rowsToProcess = hasCDC ? info.cdcInfo : info.insertInfo;
+
+ for (Row fileRow : rowsToProcess) {
+ String relPath;
+ long size;
+ Map<String, String> partitionValues;
+
+ if (hasCDC) {
+ relPath =
fileRow.getString(AddCDCFile.FULL_SCHEMA.indexOf("path"));
+ size = fileRow.getLong(AddCDCFile.FULL_SCHEMA.indexOf("size"));
+ partitionValues =
+ VectorUtils.toJavaMap(
+
fileRow.getMap(AddCDCFile.FULL_SCHEMA.indexOf("partitionValues")));
+ } else {
+ AddFile addFile = new AddFile(fileRow);
+ relPath = addFile.getPath();
+ size = addFile.getSize();
+ partitionValues =
VectorUtils.toJavaMap(addFile.getPartitionValues());
+ }
+
+ String fullPath = new org.apache.hadoop.fs.Path(tablePath,
relPath).toString();
+ List<Long> rowGroupSizes = getRowGroupSizes(fullPath, conf);
+
+ DeltaCDCReadTask task =
+ new DeltaCDCReadTask(
+ fullPath,
+ size,
+ partitionValues,
+ info.version,
+ info.timestamp,
+ hasCDC,
+ rowGroupSizes,
+ serializableScanState);
+
+ if (size >= MAX_TASK_SIZE_BYTES) {
+ if (!currentGroup.isEmpty()) {
+ emitGroup(currentGroup, out);
+ currentGroup = new ArrayList<>();
+ currentGroupSize = 0L;
+ }
+ out.output(task);
+ } else {
+ if (currentGroupSize + size > MAX_TASK_SIZE_BYTES) {
+ emitGroup(currentGroup, out);
+ currentGroup = new ArrayList<>();
+ currentGroup.add(task);
+ currentGroupSize = size;
+ } else {
+ currentGroup.add(task);
+ currentGroupSize += size;
+ }
+ }
+ }
+ }
+
+ if (!currentGroup.isEmpty()) {
+ emitGroup(currentGroup, out);
+ }
+ }
+ }
+
+ private void emitGroup(List<DeltaCDCReadTask> group,
OutputReceiver<DeltaCDCReadTask> out) {
+ for (DeltaCDCReadTask task : group) {
+ out.output(task);
+ }
+ }
+
+ private List<Long> getRowGroupSizes(String pathStr, Configuration conf) {
+ List<Long> sizes = new ArrayList<>();
+ try {
+ org.apache.hadoop.fs.Path hadoopPath = new
org.apache.hadoop.fs.Path(pathStr);
+ org.apache.parquet.hadoop.metadata.ParquetMetadata metadata =
+ org.apache.parquet.hadoop.ParquetFileReader.readFooter(
+ conf,
+ hadoopPath,
+
org.apache.parquet.format.converter.ParquetMetadataConverter.NO_FILTER);
+ for (org.apache.parquet.hadoop.metadata.BlockMetaData block :
metadata.getBlocks()) {
+ sizes.add(block.getTotalByteSize());
+ }
+ } catch (java.io.IOException e) {
+ throw new RuntimeException("Failed to read Parquet footer for " +
pathStr, e);
+ }
+ return sizes;
+ }
+
+ private static class CommitActionsInfo {
+ final long version;
+
+ final long timestamp;
+ final List<Row> cdcInfo = new ArrayList<>();
+ final List<Row> insertInfo = new ArrayList<>();
+
+ CommitActionsInfo(long version, long timestamp) {
+ this.version = version;
+ this.timestamp = timestamp;
+ }
+ }
+}
diff --git
a/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaCDCReadTask.java
b/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaCDCReadTask.java
new file mode 100644
index 00000000000..24c594343de
--- /dev/null
+++
b/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaCDCReadTask.java
@@ -0,0 +1,125 @@
+/*
+ * 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 java.io.Serializable;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import org.checkerframework.checker.nullness.qual.Nullable;
+
+/**
+ * A serializable task containing the necessary metadata to read a CDF (Change
Data Feed) file. This
+ * can be either a CDC parquet file or a regular data file (representing
inserts) from a commit.
+ */
+public class DeltaCDCReadTask implements Serializable {
+ private static final long serialVersionUID = 1L;
+
+ private final String path;
+ private final long size;
+ private final Map<String, String> partitionValues;
+ private final long version;
+ private final long timestamp;
+ private final boolean isCDC;
+ private final List<Long> rowGroupSizes;
+ private final SerializableRow scanStateRow;
+
+ public DeltaCDCReadTask(
+ String path,
+ long size,
+ Map<String, String> partitionValues,
+ long version,
+ long timestamp,
+ boolean isCDC,
+ List<Long> rowGroupSizes,
+ SerializableRow scanStateRow) {
+ this.path = path;
+ this.size = size;
+ this.partitionValues = partitionValues;
+ this.version = version;
+ this.timestamp = timestamp;
+ this.isCDC = isCDC;
+ this.rowGroupSizes = rowGroupSizes;
+ this.scanStateRow = scanStateRow;
+ }
+
+ public String getPath() {
+ return path;
+ }
+
+ public long getSize() {
+ return size;
+ }
+
+ public Map<String, String> getPartitionValues() {
+ return partitionValues;
+ }
+
+ public long getVersion() {
+ return version;
+ }
+
+ public long getTimestamp() {
+ return timestamp;
+ }
+
+ public boolean isCDC() {
+ return isCDC;
+ }
+
+ public List<Long> getRowGroupSizes() {
+ return rowGroupSizes;
+ }
+
+ public SerializableRow getScanStateRow() {
+ return scanStateRow;
+ }
+
+ @Override
+ public boolean equals(@Nullable Object o) {
+ if (this == o) {
+ return true;
+ }
+ if (!(o instanceof DeltaCDCReadTask)) {
+ return false;
+ }
+ DeltaCDCReadTask that = (DeltaCDCReadTask) o;
+ return size == that.size
+ && version == that.version
+ && timestamp == that.timestamp
+ && isCDC == that.isCDC
+ && Objects.equals(path, that.path)
+ && Objects.equals(partitionValues, that.partitionValues)
+ && Objects.equals(rowGroupSizes, that.rowGroupSizes)
+ && Objects.equals(scanStateRow, that.scanStateRow);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(
+ path, size, partitionValues, version, timestamp, isCDC, rowGroupSizes,
scanStateRow);
+ }
+
+ @Override
+ public String toString() {
+ return String.format(
+ "DeltaCDCReadTask{path='%s', size=%d, partitionValues=%s, version=%d,
timestamp=%d, isCDC=%b, "
+ + "rowGroupSizes=%s, scanStateRow=%s}",
+ path, size, partitionValues, version, timestamp, isCDC, rowGroupSizes,
scanStateRow);
+ }
+}
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
new file mode 100644
index 00000000000..cf10adc9865
--- /dev/null
+++
b/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaCDCSourceDoFn.java
@@ -0,0 +1,359 @@
+/*
+ * 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 io.delta.kernel.internal.DeltaErrors.wrapEngineException;
+
+import io.delta.kernel.Scan;
+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.engine.Engine;
+import io.delta.kernel.engine.FileReadResult;
+import io.delta.kernel.expressions.ExpressionEvaluator;
+import io.delta.kernel.expressions.Literal;
+import io.delta.kernel.internal.InternalScanFileUtils;
+import io.delta.kernel.internal.data.GenericRow;
+import io.delta.kernel.internal.data.ScanStateRow;
+import io.delta.kernel.internal.util.Utils;
+import io.delta.kernel.internal.util.VectorUtils;
+import io.delta.kernel.types.LongType;
+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.CloseableIterator;
+import io.delta.kernel.utils.FileStatus;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import org.apache.beam.sdk.io.range.OffsetRange;
+import org.apache.beam.sdk.schemas.Schema;
+import org.apache.beam.sdk.transforms.DoFn;
+import org.apache.beam.sdk.transforms.splittabledofn.RestrictionTracker;
+import org.apache.beam.sdk.values.Row;
+import org.apache.beam.sdk.values.ValueKind;
+import org.apache.hadoop.conf.Configuration;
+import org.checkerframework.checker.nullness.qual.Nullable;
+
+/**
+ * A Splittable DoFn that processes {@link DeltaCDCReadTask} elements and
reads Change Data Feed
+ * files, converting rows to Beam Rows.
+ */
[email protected]
+class DeltaCDCSourceDoFn extends DoFn<DeltaCDCReadTask, Row> {
+ @Nullable Map<String, String> hadoopConfig;
+ private transient @Nullable Engine engine;
+ private transient @Nullable Configuration conf;
+
+ public DeltaCDCSourceDoFn(@Nullable Map<String, String> hadoopConfig) {
+ this.hadoopConfig = hadoopConfig;
+ }
+
+ private synchronized Configuration getConfiguration() {
+ Configuration localConf = conf;
+ if (localConf == null) {
+ localConf = new Configuration();
+ if (hadoopConfig != null) {
+ for (Map.Entry<String, String> entry : hadoopConfig.entrySet()) {
+ localConf.set(entry.getKey(), entry.getValue());
+ }
+ }
+ conf = localConf;
+ }
+ return localConf;
+ }
+
+ private List<Long> getRowGroupSizes(DeltaCDCReadTask task) {
+ return task.getRowGroupSizes();
+ }
+
+ @GetInitialRestriction
+ public OffsetRange getInitialRestriction(@Element DeltaCDCReadTask task) {
+ List<Long> rowGroupSizes = getRowGroupSizes(task);
+ return new OffsetRange(0L, rowGroupSizes.size());
+ }
+
+ @NewTracker
+ public DeltaReadTaskTracker newTracker(
+ @Restriction OffsetRange restriction, @Element DeltaCDCReadTask task) {
+ return new DeltaReadTaskTracker(restriction, getRowGroupSizes(task));
+ }
+
+ @Setup
+ public void setUp() {
+ engine = DefaultEngine.create(getConfiguration());
+ }
+
+ @ProcessElement
+ public void processElement(
+ @Element DeltaCDCReadTask task,
+ RestrictionTracker<OffsetRange, Long> tracker,
+ OutputReceiver<Row> out)
+ throws Exception {
+
+ Engine currentEngine = engine;
+ if (currentEngine == null) {
+ throw new IllegalArgumentException("Expected the engine to not be null");
+ }
+
+ SerializableRow originalScanStateRow = task.getScanStateRow();
+ StructType logicalTableSchema =
ScanStateRow.getLogicalSchema(originalScanStateRow);
+ Schema publicBeamSchema =
DeltaIO.ReadRows.convertToBeamSchema(logicalTableSchema);
+ StructType physicalTableSchema =
ScanStateRow.getPhysicalDataReadSchema(originalScanStateRow);
+
+ StructType scanStateSchema = originalScanStateRow.getSchema();
+
+ // 1. Build modified scanState and scanFile rows depending on whether we
read a CDC file or ADD
+ // file.
+ io.delta.kernel.data.Row scanStateRow;
+ StructType readPhysicalSchema;
+ StructType readLogicalSchema;
+ Schema beamSchema;
+
+ if (task.isCDC()) {
+ readLogicalSchema = appendCDFColumns(logicalTableSchema);
+ readPhysicalSchema = appendCDFColumns(physicalTableSchema);
+ beamSchema = DeltaIO.ReadRows.convertToBeamSchema(readLogicalSchema);
+
+ HashMap<Integer, Object> valueMap = new HashMap<>();
+
+ // Tracking row level lineage is not needed.
+ Map<String, String> config =
+ new HashMap<>(ScanStateRow.getConfiguration(originalScanStateRow));
+ config.put("delta.enableRowTracking", "false");
+
+ valueMap.put(
+ scanStateSchema.indexOf("configuration"),
VectorUtils.stringStringMapValue(config));
+ valueMap.put(scanStateSchema.indexOf("logicalSchemaJson"),
readLogicalSchema.toJson());
+ valueMap.put(scanStateSchema.indexOf("physicalSchemaJson"),
readPhysicalSchema.toJson());
+ valueMap.put(
+ scanStateSchema.indexOf("partitionColumns"),
+
originalScanStateRow.getArray(scanStateSchema.indexOf("partitionColumns")));
+ valueMap.put(
+ scanStateSchema.indexOf("minReaderVersion"),
+
originalScanStateRow.getInt(scanStateSchema.indexOf("minReaderVersion")));
+ valueMap.put(
+ scanStateSchema.indexOf("minWriterVersion"),
+
originalScanStateRow.getInt(scanStateSchema.indexOf("minWriterVersion")));
+ valueMap.put(
+ scanStateSchema.indexOf("tablePath"),
+
originalScanStateRow.getString(scanStateSchema.indexOf("tablePath")));
+
+ scanStateRow = new ScanStateRow(valueMap);
+ } else {
+ // For ADD files, we read the table schema and append the CDF columns
manually afterwards.
+ // readLogicalSchema = logicalTableSchema;
+ readPhysicalSchema = physicalTableSchema;
+ beamSchema =
DeltaIO.ReadRows.convertToBeamSchema(appendCDFColumns(logicalTableSchema));
+ scanStateRow = originalScanStateRow;
+ }
+
+ io.delta.kernel.data.Row scanFileRow =
+ generateScanFileRow(task.getPath(), task.getPartitionValues());
+ FileStatus fileStatus = FileStatus.of(task.getPath(), task.getSize(),
task.getTimestamp());
+
+ BeamParquetHandler parquetHandler =
+ new BeamParquetHandler(getConfiguration(),
currentEngine.getParquetHandler(), tracker);
+ BeamEngine beamEngine = new BeamEngine(currentEngine, parquetHandler);
+
+ long currentStartRgIndex = 0L;
+
+ try (CloseableIterator<FileReadResult> fileReadResults =
+ parquetHandler.readParquetFiles(
+ Utils.singletonCloseableIterator(fileStatus),
+ readPhysicalSchema,
+ Optional.empty(),
+ currentStartRgIndex)) {
+
+ CloseableIterator<ColumnarBatch> physicalData =
+ new CloseableIterator<ColumnarBatch>() {
+ @Override
+ public void close() throws java.io.IOException {}
+
+ @Override
+ public boolean hasNext() {
+ return fileReadResults.hasNext();
+ }
+
+ @Override
+ public ColumnarBatch next() {
+ return fileReadResults.next().getData();
+ }
+ };
+
+ try (CloseableIterator<FilteredColumnarBatch> logicalBatches =
+ Scan.transformPhysicalData(beamEngine, scanStateRow, scanFileRow,
physicalData)) {
+
+ while (logicalBatches.hasNext()) {
+ FilteredColumnarBatch batch = logicalBatches.next();
+
+ if (!task.isCDC()) {
+ // For ADD files, we need to append the constant CDF columns:
+ // _change_type = "insert", _commit_version = task.version,
_commit_timestamp =
+ // task.timestamp
+ ColumnarBatch logicalBatch =
+ appendConstantCDFColumns(
+ currentEngine, batch.getData(), task.getVersion(),
task.getTimestamp());
+ // Make sure we use selection vector to considered filtered out or
deleted rows.
+ batch = new FilteredColumnarBatch(logicalBatch,
batch.getSelectionVector());
+ }
+
+ try (CloseableIterator<io.delta.kernel.data.Row> logicalRows =
batch.getRows()) {
+ while (logicalRows.hasNext()) {
+ io.delta.kernel.data.Row deltaRow = logicalRows.next();
+ Row beamRow = DeltaSourceDoFn.toBeamRow(deltaRow, beamSchema);
+ String changeType =
beamRow.getString(DeltaIO.CHANGE_TYPE_COLUMN);
+ if (changeType == null) {
+ throw new IllegalStateException(
+ "Field " + DeltaIO.CHANGE_TYPE_COLUMN + " must not be
null.");
+ }
+ ValueKind kind = getValueKind(changeType);
+ Row publicRow = projectRow(beamRow, publicBeamSchema);
+ out.builder(publicRow).setValueKind(kind).output();
+ }
+ }
+ }
+ }
+ }
+ }
+
+ 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;
+ }
+ Row.Builder builder = Row.withSchema(targetSchema);
+ for (Schema.Field field : targetSchema.getFields()) {
+ builder.addValue(row.getValue(field.getName()));
+ }
+ return builder.build();
+ }
+
+ private static ValueKind getValueKind(String changeType) {
+ // Maps Delta CDC change types to Beam's ValueKind enum.
+ //
https://docs.delta.io/delta-change-data-feed/#what-is-the-schema-for-the-change-data-feed
+ switch (changeType) {
+ case "insert":
+ return ValueKind.INSERT;
+ case "delete":
+ return ValueKind.DELETE;
+ case "update_preimage":
+ return ValueKind.UPDATE_BEFORE;
+ case "update_postimage":
+ return ValueKind.UPDATE_AFTER;
+ default:
+ throw new IllegalArgumentException("Unsupported change type: " +
changeType);
+ }
+ }
+
+ 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);
+ }
+
+ private ColumnarBatch appendConstantCDFColumns(
+ Engine engine, ColumnarBatch batch, long version, long timestamp) {
+ StructType schemaForEval = batch.getSchema();
+
+ ExpressionEvaluator changeTypeGenerator =
+ wrapEngineException(
+ () ->
+ engine
+ .getExpressionHandler()
+ .getEvaluator(schemaForEval, Literal.ofString("insert"),
StringType.STRING),
+ "Get the expression evaluator for change type");
+
+ ExpressionEvaluator commitVersionGenerator =
+ wrapEngineException(
+ () ->
+ engine
+ .getExpressionHandler()
+ .getEvaluator(schemaForEval, Literal.ofLong(version),
LongType.LONG),
+ "Get the expression evaluator for commit version");
+
+ ExpressionEvaluator commitTimestampGenerator =
+ wrapEngineException(
+ () ->
+ engine
+ .getExpressionHandler()
+ // Microseconds since epoch is expected for TimestampType
+ .getEvaluator(
+ schemaForEval,
+ Literal.ofTimestamp(timestamp * 1000L),
+ TimestampType.TIMESTAMP),
+ "Get the expression evaluator for commit timestamp");
+
+ ColumnVector changeTypeVector =
+ wrapEngineException(
+ () -> changeTypeGenerator.eval(batch), "Evaluating change type
expression");
+
+ ColumnVector commitVersionVector =
+ wrapEngineException(
+ () -> commitVersionGenerator.eval(batch), "Evaluating commit
version expression");
+
+ ColumnVector commitTimestampVector =
+ wrapEngineException(
+ () -> commitTimestampGenerator.eval(batch), "Evaluating commit
timestamp expression");
+
+ int numCols = batch.getSchema().length();
+ return batch
+ .withNewColumn(
+ numCols,
+ new StructField(DeltaIO.CHANGE_TYPE_COLUMN, StringType.STRING,
false),
+ changeTypeVector)
+ .withNewColumn(
+ numCols + 1,
+ new StructField(DeltaIO.COMMIT_VERSION_COLUMN, LongType.LONG,
false),
+ commitVersionVector)
+ .withNewColumn(
+ numCols + 2,
+ new StructField(DeltaIO.COMMIT_TIMESTAMP_COLUMN,
TimestampType.TIMESTAMP, false),
+ commitTimestampVector);
+ }
+
+ @SuppressWarnings("nullness")
+ private static io.delta.kernel.data.Row generateScanFileRow(
+ String path, Map<String, String> partitionValues) {
+ StructType addFileSchema =
+ (StructType)
InternalScanFileUtils.SCAN_FILE_SCHEMA.get("add").getDataType();
+ MapValue partMapValue = VectorUtils.stringStringMapValue(partitionValues);
+
+ Map<Integer, Object> addFileMap = new HashMap<>();
+ addFileMap.put(addFileSchema.indexOf("path"), path);
+ addFileMap.put(addFileSchema.indexOf("partitionValues"), partMapValue);
+ addFileMap.put(addFileSchema.indexOf("size"), 0L);
+ addFileMap.put(addFileSchema.indexOf("modificationTime"), 0L);
+ addFileMap.put(addFileSchema.indexOf("dataChange"), true);
+ addFileMap.put(addFileSchema.indexOf("deletionVector"), null);
+
+ io.delta.kernel.data.Row addFile = new GenericRow(addFileSchema,
addFileMap);
+
+ StructType scanFileSchema = InternalScanFileUtils.SCAN_FILE_SCHEMA;
+ Map<Integer, Object> scanFileMap = new HashMap<>();
+ scanFileMap.put(scanFileSchema.indexOf("add"), addFile);
+ scanFileMap.put(scanFileSchema.indexOf("tableRoot"), "/");
+
+ return new GenericRow(scanFileSchema, scanFileMap);
+ }
+}
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 d9c928ef659..3ac2c7a84a8 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
@@ -18,9 +18,11 @@
package org.apache.beam.sdk.io.delta;
import com.google.auto.value.AutoValue;
+import io.delta.kernel.Snapshot;
import io.delta.kernel.Table;
import io.delta.kernel.defaults.engine.DefaultEngine;
import io.delta.kernel.engine.Engine;
+import io.delta.kernel.internal.TableImpl;
import io.delta.kernel.types.ArrayType;
import io.delta.kernel.types.BinaryType;
import io.delta.kernel.types.BooleanType;
@@ -50,23 +52,6 @@ import org.checkerframework.checker.nullness.qual.Nullable;
/**
* A connector that reads from <a href="https://delta.io/">Delta Lake</a>
tables.
*
- * <p>{@link DeltaIO} is offered as a Managed transform. This class is subject
to change and should
- * not be used directly. Instead, use it like so:
- *
- * <pre>{@code
- * Map<String, Object> config = Map.of(
- * "table", "gs://my-bucket/delta-table",
- * "hadoop_config", Map.of(
- * "fs.gs.impl",
"com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem",
- * "fs.AbstractFileSystem.gs.impl",
"com.google.cloud.hadoop.fs.gcs.GoogleHadoopFS",
- * "fs.gs.project.id", "my-project-id"));
- *
- * pipeline
- * .apply(Managed.read(Managed.DELTA_LAKE).withConfig(config))
- * .getSinglePCollection()
- * .apply(ParDo.of(...));
- * }</pre>
- *
* <h2>Configuration Options</h2>
*
* Please check the <a
href="https://beam.apache.org/documentation/io/managed-io/">Managed IO
@@ -78,6 +63,10 @@ import org.checkerframework.checker.nullness.qual.Nullable;
@Internal
public class DeltaIO {
+ public static final String CHANGE_TYPE_COLUMN = "_change_type";
+ public static final String COMMIT_VERSION_COLUMN = "_commit_version";
+ public static final String COMMIT_TIMESTAMP_COLUMN = "_commit_timestamp";
+
/**
* Reads rows from a Delta Lake table.
*
@@ -88,6 +77,11 @@ public class DeltaIO {
return new AutoValue_DeltaIO_ReadRows.Builder().build();
}
+ /** Reads change data feed (CDC) from a Delta Lake table. */
+ public static ReadChanges readChanges() {
+ return new AutoValue_DeltaIO_ReadChanges.Builder().build();
+ }
+
@AutoValue
public abstract static class ReadRows extends PTransform<PBegin,
PCollection<Row>> {
@@ -211,4 +205,126 @@ public class DeltaIO {
}
}
}
+
+ @AutoValue
+ public abstract static class ReadChanges extends PTransform<PBegin,
PCollection<Row>> {
+ public abstract @Nullable String getTablePath();
+
+ public abstract @Nullable Long getStartVersion();
+
+ public abstract @Nullable String getStartTimestamp();
+
+ public abstract @Nullable Long getEndVersion();
+
+ public abstract @Nullable String getEndTimestamp();
+
+ public abstract @Nullable Map<String, String> getHadoopConfig();
+
+ abstract Builder toBuilder();
+
+ @AutoValue.Builder
+ abstract static class Builder {
+ abstract Builder setTablePath(String tablePath);
+
+ abstract Builder setStartVersion(@Nullable Long startVersion);
+
+ abstract Builder setStartTimestamp(@Nullable String startTimestamp);
+
+ abstract Builder setEndVersion(@Nullable Long endVersion);
+
+ abstract Builder setEndTimestamp(@Nullable String endTimestamp);
+
+ abstract Builder setHadoopConfig(@Nullable Map<String, String>
hadoopConfig);
+
+ abstract ReadChanges build();
+ }
+
+ public ReadChanges from(String tablePath) {
+ return toBuilder().setTablePath(tablePath).build();
+ }
+
+ public ReadChanges withStartVersion(long startVersion) {
+ return toBuilder().setStartVersion(startVersion).build();
+ }
+
+ public ReadChanges withStartTimestamp(String startTimestamp) {
+ return toBuilder().setStartTimestamp(startTimestamp).build();
+ }
+
+ public ReadChanges withEndVersion(long endVersion) {
+ return toBuilder().setEndVersion(endVersion).build();
+ }
+
+ public ReadChanges withEndTimestamp(String endTimestamp) {
+ return toBuilder().setEndTimestamp(endTimestamp).build();
+ }
+
+ public ReadChanges withConfig(Map<String, String> config) {
+ return toBuilder().setHadoopConfig(config).build();
+ }
+
+ @Override
+ public PCollection<Row> expand(PBegin input) {
+ String path = getTablePath();
+ if (path == null) {
+ throw new IllegalArgumentException("Table path must be set.");
+ }
+ if (getStartVersion() == null && getStartTimestamp() == null) {
+ // TODO: for unbounded reads, support using current HEAD or the latest
snapshot
+ // as the default starting point.
+ throw new IllegalArgumentException("Either startVersion or
startTimestamp must be set.");
+ }
+ if (getStartVersion() != null && getStartTimestamp() != null) {
+ throw new IllegalArgumentException("Cannot set both startVersion and
startTimestamp.");
+ }
+ if (getEndVersion() != null && getEndTimestamp() != null) {
+ throw new IllegalArgumentException("Cannot set both endVersion and
endTimestamp.");
+ }
+
+ Configuration conf = new Configuration();
+ Map<String, String> hadoopConfig = getHadoopConfig();
+ if (hadoopConfig != null) {
+ for (Map.Entry<String, String> entry : hadoopConfig.entrySet()) {
+ conf.set(entry.getKey(), entry.getValue());
+ }
+ }
+ Engine engine = DefaultEngine.create(conf);
+ Table table = Table.forPath(engine, path);
+
+ TableImpl tableImpl = (TableImpl) table;
+
+ long resolvedEndVersion;
+ Long endVersionVal = getEndVersion();
+ String endTimestampVal = getEndTimestamp();
+ if (endVersionVal != null) {
+ resolvedEndVersion = endVersionVal;
+ } else if (endTimestampVal != null) {
+ long endMillis =
java.time.Instant.parse(endTimestampVal).toEpochMilli();
+ resolvedEndVersion = tableImpl.getVersionBeforeOrAtTimestamp(engine,
endMillis);
+ } else {
+ resolvedEndVersion = table.getLatestSnapshot(engine).getVersion();
+ }
+
+ Snapshot endSnapshot = table.getSnapshotAsOfVersion(engine,
resolvedEndVersion);
+ StructType deltaSchema = endSnapshot.getSchema();
+ if (deltaSchema == null) {
+ throw new IllegalStateException("Table schema is null.");
+ }
+ Schema beamSchema = ReadRows.convertToBeamSchema(deltaSchema);
+
+ return input
+ .apply("Create Path", Create.of(path))
+ .apply(
+ "Plan CDF Files",
+ ParDo.of(
+ new CreateCDCReadTasksDoFn(
+ hadoopConfig,
+ getStartVersion(),
+ getStartTimestamp(),
+ getEndVersion(),
+ getEndTimestamp())))
+ .apply("Read CDF Data", ParDo.of(new
DeltaCDCSourceDoFn(hadoopConfig)))
+ .setRowSchema(beamSchema);
+ }
+ }
}
diff --git
a/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaReadSchemaTransformProvider.java
b/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaReadSchemaTransformProvider.java
index 42ca3f24def..48dc3a2c748 100644
---
a/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaReadSchemaTransformProvider.java
+++
b/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaReadSchemaTransformProvider.java
@@ -130,11 +130,11 @@ public class DeltaReadSchemaTransformProvider extends
TypedSchemaTransformProvid
abstract static class Builder {
abstract Builder setTable(String table);
- abstract Builder setVersion(@Nullable Long version);
+ abstract Builder setVersion(Long version);
- abstract Builder setTimestamp(@Nullable String timestamp);
+ abstract Builder setTimestamp(String timestamp);
- abstract Builder setHadoopConfig(@Nullable Map<String, String>
hadoopConfig);
+ abstract Builder setHadoopConfig(Map<String, String> hadoopConfig);
abstract Configuration build();
}
diff --git
a/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaSourceDoFn.java
b/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaSourceDoFn.java
index bd53c3c9d04..fca19e30cb8 100644
---
a/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaSourceDoFn.java
+++
b/sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaSourceDoFn.java
@@ -213,7 +213,7 @@ class DeltaSourceDoFn extends DoFn<DeltaReadTask, Row> {
}
// Convert Delta `Row` to Beam `Row`.
- private static Row toBeamRow(io.delta.kernel.data.Row deltaRow, Schema
beamSchema) {
+ static Row toBeamRow(io.delta.kernel.data.Row deltaRow, Schema beamSchema) {
Row.Builder builder = Row.withSchema(beamSchema);
StructType deltaSchema = deltaRow.getSchema();
List<StructField> fields = deltaSchema.fields();
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 bd8bf8b3c8c..f00b34be460 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,9 +17,23 @@
*/
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;
@@ -30,11 +44,19 @@ 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;
@@ -48,12 +70,16 @@ import org.apache.beam.sdk.testing.PAssert;
import org.apache.beam.sdk.testing.TestPipeline;
import org.apache.beam.sdk.transforms.Count;
import org.apache.beam.sdk.transforms.Create;
+import org.apache.beam.sdk.transforms.DoFn;
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.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;
import org.junit.Test;
@@ -122,9 +148,6 @@ public class DeltaIOTest {
io.delta.kernel.Snapshot snapshot = table.getLatestSnapshot(engine);
io.delta.kernel.Scan scan = snapshot.getScanBuilder().build();
- io.delta.kernel.data.Row scanState = scan.getScanState(engine);
- System.err.println("SCAN STATE SCHEMA: " +
scanState.getSchema().toString());
-
try
(io.delta.kernel.utils.CloseableIterator<io.delta.kernel.data.FilteredColumnarBatch>
scanFiles = scan.getScanFiles(engine)) {
while (scanFiles.hasNext()) {
@@ -300,53 +323,8 @@ public class DeltaIOTest {
writePipeline.run().waitUntilFinish();
- System.out.println("FILES IN TABLE DIR:");
- for (File f : tableDir.listFiles()) {
- System.out.println(
- " - " + f.getName() + " (size=" + f.length() + ", isDir=" +
f.isDirectory() + ")");
- if (f.isDirectory()) {
- for (File sub : f.listFiles()) {
- System.out.println(" - " + sub.getName() + " (size=" +
sub.length() + ")");
- }
- }
- }
-
File parquetFile = new File(tableDir, "part-00000.parquet");
byte[] fileBytes = Files.readAllBytes(parquetFile.toPath());
- System.out.println("PARQUET FILE LENGTH: " + fileBytes.length);
- if (fileBytes.length >= 8) {
- System.out.println(
- "PARQUET FIRST 4 BYTES: "
- + fileBytes[0]
- + ", "
- + fileBytes[1]
- + ", "
- + fileBytes[2]
- + ", "
- + fileBytes[3]
- + " ('"
- + (char) fileBytes[0]
- + (char) fileBytes[1]
- + (char) fileBytes[2]
- + (char) fileBytes[3]
- + "')");
- int len = fileBytes.length;
- System.out.println(
- "PARQUET LAST 4 BYTES: "
- + fileBytes[len - 4]
- + ", "
- + fileBytes[len - 3]
- + ", "
- + fileBytes[len - 2]
- + ", "
- + fileBytes[len - 1]
- + " ('"
- + (char) fileBytes[len - 4]
- + (char) fileBytes[len - 3]
- + (char) fileBytes[len - 2]
- + (char) fileBytes[len - 1]
- + "')");
- }
// 2. Create the Delta log
File logDir = new File(tableDir, "_delta_log");
@@ -392,28 +370,19 @@ public class DeltaIOTest {
@Test
public void testManagedDeltaRead() throws Exception {
File tableDir = tempFolder.newFolder("managed-delta-table");
+ Engine engine = DefaultEngine.create(new
org.apache.hadoop.conf.Configuration());
- // 1. Write a Parquet file to simulate a Delta table
Schema schema = Schema.builder().addField("name",
Schema.FieldType.STRING).build();
Row row = Row.withSchema(schema).addValues("test-name").build();
- writeParquetFile(new File(tableDir, "part-00000.parquet"), row);
+ StructType deltaSchema = new StructType().add("name", StringType.STRING);
- // 2. Create the Delta log
- File logDir = new File(tableDir, "_delta_log");
- logDir.mkdirs();
- File commitFile = new File(logDir, "00000000000000000000.json");
-
- File parquetFile = new File(tableDir, "part-00000.parquet");
- byte[] fileBytes = Files.readAllBytes(parquetFile.toPath());
-
- String commitContent =
- "{\"protocol\":{\"minReaderVersion\":1,\"minWriterVersion\":2}}\n"
- +
"{\"metaData\":{\"id\":\"test-id\",\"format\":{\"provider\":\"parquet\",\"options\":{}},\"schemaString\":\"{\\\"type\\\":\\\"struct\\\",\\\"fields\\\":[{\\\"name\\\":\\\"name\\\",\\\"type\\\":\\\"string\\\",\\\"nullable\\\":true,\\\"metadata\\\":{}}]}\",\"partitionColumns\":[],\"configuration\":{},\"createdAt\":123456789}}\n"
- +
"{\"add\":{\"path\":\"part-00000.parquet\",\"partitionValues\":{},\"size\":"
- + fileBytes.length
- + ",\"modificationTime\":123456789,\"dataChange\":true}}";
-
- Files.write(commitFile.toPath(),
commitContent.getBytes(StandardCharsets.UTF_8));
+ writeAppendCommit(
+ engine,
+ tableDir.getAbsolutePath(),
+ 0L,
+ 123456789L,
+ deltaSchema,
+ Collections.singletonList(row));
// 3. Read it using Managed
PCollection<Row> output =
@@ -810,4 +779,744 @@ public class DeltaIOTest {
}
}
}
+
+ @Test
+ public void testReadChanges() throws Exception {
+ File tableDir = tempFolder.newFolder("delta-table-changes");
+ 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);
+
+ 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();
+
+ writeCdcCommit(
+ engine,
+ tableDir.getAbsolutePath(),
+ 1L,
+ 200000000000L,
+ deltaSchema,
+ null,
+ null,
+ java.util.Arrays.asList(cdcRow1, cdcRow2, cdcRow3),
+ cdcWriteDeltaSchema);
+
+ // 3. Read CDF data from table using ReadChanges
+ PCollection<Row> output =
+ readPipeline.apply(
+
DeltaIO.readChanges().from(tableDir.getAbsolutePath()).withStartVersion(0L));
+
+ 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");
+ Engine engine = DefaultEngine.create(new
org.apache.hadoop.conf.Configuration());
+
+ Schema tableSchema = Schema.builder().addField("name",
Schema.FieldType.STRING).build();
+ StructType deltaSchema = new StructType().add("name", StringType.STRING);
+
+ // 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(
+ engine,
+ tableDir.getAbsolutePath(),
+ 0L,
+ 100000000000L,
+ deltaSchema,
+ java.util.Arrays.asList(tableRow1, tableRow2));
+
+ // 2. Write parquet files for Version 1 (commit with updates and deletes)
+ 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(200000000000L))
+ .build();
+ Row cdcRow2 =
+ Row.withSchema(cdcWriteSchema)
+ .addValues("row-1-updated", "update_postimage", 1L, new
Instant(200000000000L))
+ .build();
+ Row cdcRow3 =
+ Row.withSchema(cdcWriteSchema)
+ .addValues("row-2", "delete", 1L, new Instant(200000000000L))
+ .build();
+
+ writeCdcCommit(
+ engine,
+ tableDir.getAbsolutePath(),
+ 1L,
+ 200000000000L,
+ deltaSchema,
+ null,
+ null,
+ java.util.Arrays.asList(cdcRow1, cdcRow2, cdcRow3),
+ cdcWriteDeltaSchema);
+
+ // 3. Write parquet files for Version 2 (insert-only commit)
+ Row tableRow3 = Row.withSchema(tableSchema).addValues("row-3").build();
+ writeAppendCommit(
+ engine,
+ tableDir.getAbsolutePath(),
+ 2L,
+ 300000000000L,
+ deltaSchema,
+ java.util.Arrays.asList(tableRow3));
+
+ // Test 1: Read changes between start version 0 and end version 2
+ PCollection<Row> outputVersions =
+ readPipeline.apply(
+ "Read Changes Version Range",
+ DeltaIO.readChanges()
+ .from(tableDir.getAbsolutePath())
+ .withStartVersion(0L)
+ .withEndVersion(2L));
+
+ PCollection<String> formattedVersions =
+ outputVersions.apply("Format Version Output", ParDo.of(new
FormatValueKindAndRow()));
+
+ PAssert.that(formattedVersions)
+ .containsInAnyOrder(
+ "INSERT:row-1",
+ "INSERT:row-2",
+ "UPDATE_BEFORE:row-1",
+ "UPDATE_AFTER:row-1-updated",
+ "DELETE:row-2",
+ "INSERT:row-3");
+
+ // Test 2: Read changes between start timestamp (after version 0) and end
timestamp (after
+ // version 2)
+ String startTimestamp =
java.time.Instant.ofEpochMilli(150000000000L).toString();
+ String endTimestamp =
java.time.Instant.ofEpochMilli(350000000000L).toString();
+
+ PCollection<Row> outputTimestamps =
+ filteringPipeline.apply(
+ "Read Changes Timestamp Range",
+ DeltaIO.readChanges()
+ .from(tableDir.getAbsolutePath())
+ .withStartTimestamp(startTimestamp)
+ .withEndTimestamp(endTimestamp));
+
+ PCollection<String> formattedTimestamps =
+ outputTimestamps.apply("Format Timestamp Output", ParDo.of(new
FormatValueKindAndRow()));
+
+ PAssert.that(formattedTimestamps)
+ .containsInAnyOrder(
+ "UPDATE_BEFORE:row-1", "UPDATE_AFTER:row-1-updated",
"DELETE:row-2", "INSERT:row-3");
+
+ readPipeline.run().waitUntilFinish();
+ filteringPipeline.run().waitUntilFinish();
+ }
+
+ @Test
+ public void testReadChangesPartialRange() throws Exception {
+ File tableDir = tempFolder.newFolder("delta-table-changes-partial-range");
+ Engine engine = DefaultEngine.create(new
org.apache.hadoop.conf.Configuration());
+
+ Schema tableSchema = Schema.builder().addField("name",
Schema.FieldType.STRING).build();
+ StructType deltaSchema = new StructType().add("name", StringType.STRING);
+
+ // 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(
+ engine,
+ tableDir.getAbsolutePath(),
+ 0L,
+ 100000000000L,
+ deltaSchema,
+ java.util.Arrays.asList(tableRow1, tableRow2));
+
+ // 2. Write parquet files for Version 1 (commit with updates and deletes)
+ 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(200000000000L))
+ .build();
+ Row cdcRow2 =
+ Row.withSchema(cdcWriteSchema)
+ .addValues("row-1-updated", "update_postimage", 1L, new
Instant(200000000000L))
+ .build();
+ Row cdcRow3 =
+ Row.withSchema(cdcWriteSchema)
+ .addValues("row-2", "delete", 1L, new Instant(200000000000L))
+ .build();
+
+ writeCdcCommit(
+ engine,
+ tableDir.getAbsolutePath(),
+ 1L,
+ 200000000000L,
+ deltaSchema,
+ null,
+ null,
+ java.util.Arrays.asList(cdcRow1, cdcRow2, cdcRow3),
+ cdcWriteDeltaSchema);
+
+ // 3. Write parquet files for Version 2 (insert-only commit)
+ Row tableRow3 = Row.withSchema(tableSchema).addValues("row-3").build();
+ writeAppendCommit(
+ engine,
+ tableDir.getAbsolutePath(),
+ 2L,
+ 300000000000L,
+ deltaSchema,
+ java.util.Arrays.asList(tableRow3));
+
+ // 4. Write parquet files for Version 3 (commit with updates and deletes)
+ Row cdcRow4 =
+ Row.withSchema(cdcWriteSchema)
+ .addValues("row-3", "update_preimage", 3L, new
Instant(400000000000L))
+ .build();
+ Row cdcRow5 =
+ Row.withSchema(cdcWriteSchema)
+ .addValues("row-3-updated", "update_postimage", 3L, new
Instant(400000000000L))
+ .build();
+ Row cdcRow6 =
+ Row.withSchema(cdcWriteSchema)
+ .addValues("row-1-updated", "delete", 3L, new
Instant(400000000000L))
+ .build();
+
+ writeCdcCommit(
+ engine,
+ tableDir.getAbsolutePath(),
+ 3L,
+ 400000000000L,
+ deltaSchema,
+ null,
+ null,
+ java.util.Arrays.asList(cdcRow4, cdcRow5, cdcRow6),
+ cdcWriteDeltaSchema);
+
+ // 5. Write parquet files for Version 4 (insert-only commit)
+ Row tableRow4 = Row.withSchema(tableSchema).addValues("row-4").build();
+ writeAppendCommit(
+ engine,
+ tableDir.getAbsolutePath(),
+ 4L,
+ 500000000000L,
+ deltaSchema,
+ java.util.Arrays.asList(tableRow4));
+
+ // Read changes between start version 1 and end version 3
+ PCollection<Row> outputVersions =
+ readPipeline.apply(
+ "Read Changes Partial Version Range",
+ DeltaIO.readChanges()
+ .from(tableDir.getAbsolutePath())
+ .withStartVersion(1L)
+ .withEndVersion(3L));
+
+ PCollection<String> formattedVersions =
+ outputVersions.apply("Format Version Output", ParDo.of(new
FormatValueKindAndRow()));
+
+ PAssert.that(formattedVersions)
+ .containsInAnyOrder(
+ "UPDATE_BEFORE:row-1",
+ "UPDATE_AFTER:row-1-updated",
+ "DELETE:row-2",
+ "INSERT:row-3",
+ "UPDATE_BEFORE:row-3",
+ "UPDATE_AFTER:row-3-updated",
+ "DELETE:row-1-updated");
+
+ readPipeline.run().waitUntilFinish();
+ }
+
+ private static final class FormatValueKindAndRow extends DoFn<Row, String> {
+ @ProcessElement
+ public void process(
+ @Element Row row, ValueKind valueKind, OutputReceiver<String>
outputReceiver) {
+ outputReceiver.output(valueKind.name() + ":" + row.getString("name"));
+ }
+ }
+
+ 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());
+ }
+ 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));
+ }
+ }
+
+ 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);
+ }
+ }
+ 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);
+ }
+
+ @Override
+ public io.delta.kernel.data.MapValue getMap(int ord) {
+ return (io.delta.kernel.data.MapValue) getVal(ord);
+ }
+ }
}