-
Notifications
You must be signed in to change notification settings - Fork 4.6k
Adds support for reading at a given Delta Lake version or timestamp #39758
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from 1 commit
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| 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 |
|---|---|---|
|
|
@@ -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
|
||
| * with {@code Managed.DELTA_LAKE} instead of directly using this transform. | ||
| */ | ||
| public static ReadRows readRows() { | ||
|
|
@@ -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(); | ||
| } | ||
|
|
@@ -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(); | ||
|
|
@@ -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."); | ||
|
|
@@ -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); | ||
| } | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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"); | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. There are duplicated codes. Can we move the creation of hadoop Configuration into a helper function?
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Done. |
||
| 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); | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. sleep is generally unreliable especially on CI. Consider using a reliable way to detect initial setup is finished, e.g. poll the file system?
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Done. |
||
| 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); | ||
|
|
||
There was a problem hiding this comment.
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:buildnot:sdks:java:io:delta:integrationTestThere was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
It's included in https://github.com/apache/beam/blob/master/.github/workflows/beam_PostCommit_Java_Delta_IO_Dataflow.yml.
Executed for the current PR here: https://github.com/apache/beam/actions/runs/31854212817/job/94935775426?pr=39758