Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run.",
"modification": 4
"modification": 2
}
Original file line number Diff line number Diff line change
Expand Up @@ -38,9 +38,20 @@
class CreateReadTasksDoFn extends DoFn<String, DeltaReadTask> {
private static final long MAX_TASK_SIZE_BYTES = 1024L * 1024L * 1024L; // 1 GB
private final @Nullable Map<String, String> hadoopConfig;
private final @Nullable Long version;
private final @Nullable String timestamp;

public CreateReadTasksDoFn(@Nullable Map<String, String> hadoopConfig) {
this(hadoopConfig, null, null);
}

public CreateReadTasksDoFn(
@Nullable Map<String, String> hadoopConfig,
@Nullable Long version,
@Nullable String timestamp) {
this.hadoopConfig = hadoopConfig;
this.version = version;
this.timestamp = timestamp;
}

@ProcessElement
Expand All @@ -54,7 +65,17 @@ public void processElement(@Element String tablePath, OutputReceiver<DeltaReadTa
}
Engine engine = DefaultEngine.create(conf);
Table table = Table.forPath(engine, tablePath);
Snapshot snapshot = table.getLatestSnapshot(engine);
Snapshot snapshot;
Long versionVal = version;
String timestampVal = timestamp;
if (versionVal != null) {
snapshot = table.getSnapshotAsOfVersion(engine, versionVal);
} else if (timestampVal != null) {
long timestampMillis = java.time.Instant.parse(timestampVal).toEpochMilli();
snapshot = table.getSnapshotAsOfTimestamp(engine, timestampMillis);
} else {
snapshot = table.getLatestSnapshot(engine);
}
Scan scan = snapshot.getScanBuilder().build();
Row scanState = scan.getScanState(engine);
SerializableRow serializableScanState = new SerializableRow(scanState);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -72,7 +72,7 @@
/**
* Reads rows from a Delta Lake table.
*
* <p>Normally, it is recommended to use {@link org.apache.beam.sdk.managed.Managed#read(String)}

Check warning on line 75 in sdks/java/io/delta/src/main/java/org/apache/beam/sdk/io/delta/DeltaIO.java

View workflow job for this annotation

GitHub Actions / beam_PreCommit_Java_Delta_IO_Direct (Run Java_Delta_IO_Direct PreCommit)

Tag @link: reference not found: org.apache.beam.sdk.managed.Managed#read(String)
* with {@code Managed.DELTA_LAKE} instead of directly using this transform.
*/
public static ReadRows readRows() {
Expand Down Expand Up @@ -114,10 +114,23 @@
return toBuilder().setTablePath(tablePath).build();
}

/**
* Specifies the version of the Delta Lake table to read.
*
* <p>Only one of version or timestamp should be provided. If neither is provided, the latest
* version (HEAD) is read.
*/
public ReadRows withVersion(@Nullable Long version) {
return toBuilder().setVersion(version).build();
}

/**
* Specifies the timestamp of the Delta Lake table to read as an ISO 8601 string (e.g.
* "2026-05-20T15:43:26Z").
*
* <p>Only one of version or timestamp should be provided. If neither is provided, the latest
* version (HEAD) is read.
*/
public ReadRows withTimestamp(@Nullable String timestamp) {
return toBuilder().setTimestamp(timestamp).build();
}
Expand All @@ -132,14 +145,8 @@
if (path == null) {
throw new IllegalArgumentException("Table path must be set.");
}
if (getTimestamp() != null) {
throw new UnsupportedOperationException(
"Reading from a specific timestamp is not supported yet");
}

if (getVersion() != null) {
throw new UnsupportedOperationException(
"Reading from a specific version is not supported yet");
if (getVersion() != null && getTimestamp() != null) {
throw new IllegalArgumentException("Cannot set both version and timestamp.");
}

Configuration conf = new Configuration();
Expand All @@ -151,7 +158,17 @@
}
Engine engine = DefaultEngine.create(conf);
Table table = Table.forPath(engine, path);
io.delta.kernel.Snapshot snapshot = table.getLatestSnapshot(engine);
Snapshot snapshot;
Long versionVal = getVersion();
String timestampVal = getTimestamp();
if (versionVal != null) {
snapshot = table.getSnapshotAsOfVersion(engine, versionVal);
} else if (timestampVal != null) {
long timestampMillis = java.time.Instant.parse(timestampVal).toEpochMilli();
snapshot = table.getSnapshotAsOfTimestamp(engine, timestampMillis);
} else {
snapshot = table.getLatestSnapshot(engine);
}
StructType deltaSchema = snapshot.getSchema();
if (deltaSchema == null) {
throw new IllegalStateException("Table schema is null.");
Expand All @@ -160,7 +177,9 @@

return input
.apply("Create Path", Create.of(path))
.apply("Plan Files", ParDo.of(new CreateReadTasksDoFn(hadoopConfig)))
.apply(
"Plan Files",
ParDo.of(new CreateReadTasksDoFn(hadoopConfig, getVersion(), getTimestamp())))
.apply("Read Logical Data", ParDo.of(new DeltaSourceDoFn(hadoopConfig)))
.setRowSchema(beamSchema);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -114,11 +114,13 @@ static Builder builder() {
@SchemaFieldDescription("Identifier of the Delta Lake table.")
abstract String getTable();

@SchemaFieldDescription("Version of the Delta Lake table to read.")
@SchemaFieldDescription(
"Version of the Delta Lake table to read. Cannot be set if timestamp is set.")
@Nullable
abstract Long getVersion();

@SchemaFieldDescription("Timestamp of the Delta Lake table to read.")
@SchemaFieldDescription(
"Timestamp of the Delta Lake table to read (in UTC ISO 8601 format, e.g. 2026-05-20T15:43:26Z). Cannot be set if version is set.")
@Nullable
abstract String getTimestamp();

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -302,6 +302,116 @@ public void testReadDeltaLakeTable() {
readPipeline.run().waitUntilFinish();
}

@Test
public void testReadDeltaLakeTableAtTimestamp() throws Exception {
ExperimentalOptions options = readPipeline.getOptions().as(ExperimentalOptions.class);
ExperimentalOptions.addExperiment(options, "use_runner_v2");

Map<String, String> hadoopConfig = new HashMap<>();
hadoopConfig.put("fs.gs.impl", "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem");
hadoopConfig.put(
"fs.AbstractFileSystem.gs.impl", "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFS");
hadoopConfig.put("fs.gs.auth.type", "APPLICATION_DEFAULT");
String project =
readPipeline
.getOptions()
.as(org.apache.beam.sdk.extensions.gcp.options.GcpOptions.class)
.getProject();
if (project != null) {
hadoopConfig.put("fs.gs.project.id", project);
}

org.apache.hadoop.conf.Configuration conf = new org.apache.hadoop.conf.Configuration();
for (Map.Entry<String, String> entry : hadoopConfig.entrySet()) {
conf.set(entry.getKey(), entry.getValue());
}
Engine engine = DefaultEngine.create(conf);

// Wait briefly to ensure timestamp is after version 0 commit
Thread.sleep(1000);
String timestampV0 = java.time.Instant.ofEpochMilli(System.currentTimeMillis()).toString();
Thread.sleep(1000);

// Write version 1 with additional rows
List<Row> additionalRows =
IntStream.range(100, 150)
.mapToObj(i -> Row.withSchema(ROW_SCHEMA).addValues(i, "name_" + i).build())
.collect(Collectors.toList());

StructType deltaSchema =
new StructType().add("id", IntegerType.INTEGER).add("name", StringType.STRING);

DeltaWriteTestUtils.writeAppendCommit(
engine, repoPath, 1L, System.currentTimeMillis(), deltaSchema, additionalRows);

PCollection<Row> output =
readPipeline
.apply(
Managed.read(Managed.DELTA_LAKE)
.withConfig(
ImmutableMap.of(
"table",
repoPath,
"timestamp",
timestampV0,
"hadoop_config",
hadoopConfig)))
.getSinglePCollection();

PAssert.that(output).containsInAnyOrder(TEST_ROWS);
readPipeline.run().waitUntilFinish();
}

@Test
public void testReadDeltaLakeTableAtVersion() throws Exception {
ExperimentalOptions options = readPipeline.getOptions().as(ExperimentalOptions.class);
ExperimentalOptions.addExperiment(options, "use_runner_v2");

Map<String, String> hadoopConfig = new HashMap<>();
hadoopConfig.put("fs.gs.impl", "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem");
hadoopConfig.put(
"fs.AbstractFileSystem.gs.impl", "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFS");
hadoopConfig.put("fs.gs.auth.type", "APPLICATION_DEFAULT");
String project =
readPipeline
.getOptions()
.as(org.apache.beam.sdk.extensions.gcp.options.GcpOptions.class)
.getProject();
if (project != null) {
hadoopConfig.put("fs.gs.project.id", project);
}

org.apache.hadoop.conf.Configuration conf = new org.apache.hadoop.conf.Configuration();
for (Map.Entry<String, String> entry : hadoopConfig.entrySet()) {
conf.set(entry.getKey(), entry.getValue());
}
Engine engine = DefaultEngine.create(conf);

// Write version 1 with additional rows
List<Row> additionalRows =
IntStream.range(100, 150)
.mapToObj(i -> Row.withSchema(ROW_SCHEMA).addValues(i, "name_" + i).build())
.collect(Collectors.toList());

StructType deltaSchema =
new StructType().add("id", IntegerType.INTEGER).add("name", StringType.STRING);

DeltaWriteTestUtils.writeAppendCommit(
engine, repoPath, 1L, System.currentTimeMillis(), deltaSchema, additionalRows);

PCollection<Row> output =
readPipeline
.apply(
Managed.read(Managed.DELTA_LAKE)
.withConfig(
ImmutableMap.of(
"table", repoPath, "version", 0L, "hadoop_config", hadoopConfig)))
.getSinglePCollection();

PAssert.that(output).containsInAnyOrder(TEST_ROWS);
readPipeline.run().waitUntilFinish();
}

@Test
public void testReadChangesDeltaLake() throws Exception {
ExperimentalOptions options = readPipeline.getOptions().as(ExperimentalOptions.class);
Expand Down
Loading
Loading