Skip to content
Merged
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": 3
}
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 @@ -132,14 +132,8 @@ public PCollection<Row> expand(PBegin input) {
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 +145,17 @@ public PCollection<Row> expand(PBegin input) {
}
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 +164,9 @@ public PCollection<Row> expand(PBegin input) {

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 @@ -49,6 +49,7 @@
import java.util.Optional;
import java.util.stream.Collectors;
import java.util.stream.IntStream;
import org.apache.beam.sdk.extensions.gcp.options.GcpOptions;
import org.apache.beam.sdk.managed.Managed;
import org.apache.beam.sdk.options.ExperimentalOptions;
import org.apache.beam.sdk.schemas.Schema;
Expand Down Expand Up @@ -114,19 +115,7 @@ public void setup() throws Exception {
LOG.info("Generating Delta Lake repository at {}", repoPath);

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

getHadoopConfig().forEach(configuration::set);
Engine engine = DefaultEngine.create(configuration);
Table table = Table.forPath(engine, repoPath);

Expand Down Expand Up @@ -278,18 +267,7 @@ public void testReadDeltaLakeTable() {
ExperimentalOptions options = readPipeline.getOptions().as(ExperimentalOptions.class);
ExperimentalOptions.addExperiment(options, "use_runner_v2");

Map<String, String> hadoopConfig = new HashMap<>();
hadoopConfig.put("fs.gs.impl", "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem");
hadoopConfig.put(
"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);
}
Map<String, String> hadoopConfig = getHadoopConfig();

PCollection<Row> output =
readPipeline
Expand All @@ -302,6 +280,85 @@ 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 = getHadoopConfig();
Configuration conf = new Configuration();
hadoopConfig.forEach(conf::set);
Engine engine = DefaultEngine.create(conf);

Table table = Table.forPath(engine, repoPath);
long commitTimestampV0 = table.getSnapshotAsOfVersion(engine, 0L).getTimestamp(engine);
String timestampV0 = java.time.Instant.ofEpochMilli(commitTimestampV0).toString();

// 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 = getHadoopConfig();
Configuration conf = new Configuration();
hadoopConfig.forEach(conf::set);
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 All @@ -314,24 +371,9 @@ public void testReadChangesDeltaLake() throws Exception {
options.setExperiments(modifiableExperiments);
}

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

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

StructType deltaSchema =
Expand Down Expand Up @@ -413,6 +455,19 @@ public void testReadChangesDeltaLake() throws Exception {
readPipeline.run().waitUntilFinish();
}

private Map<String, String> getHadoopConfig() {
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(GcpOptions.class).getProject();
if (project != null) {
hadoopConfig.put("fs.gs.project.id", project);
}
return hadoopConfig;
}

private static final class FormatITRowWithMetadata extends DoFn<Row, String> {
@ProcessElement
public void process(@Element Row row, OutputReceiver<String> out) {
Expand Down
Loading
Loading