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 @@ -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 @@ -132,14 +132,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 +145,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 +164,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 @@ -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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is DeltaIOIT exercised by any GHA tests? Checking https://github.com/apache/beam/blob/master/.github/workflows/beam_PreCommit_Java_Delta_IO_Direct.yml it only runs :sdks:java:io:delta:build not :sdks:java:io:delta:integrationTest

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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