Repository navigation
Spark 4.2: Apply session snapshot properties to maintenance actions #18406
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
base: main
Are you sure you want to change the base?
Changes from all commits
0b1d5db
e6977a5
8b8b849
4588620
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 |
|---|---|---|
|
|
@@ -113,6 +113,7 @@ | |
| import org.apache.iceberg.spark.ScanTaskSetManager; | ||
| import org.apache.iceberg.spark.SparkReadConf; | ||
| import org.apache.iceberg.spark.SparkReadOptions; | ||
| import org.apache.iceberg.spark.SparkSQLProperties; | ||
| import org.apache.iceberg.spark.SparkSchemaUtil; | ||
| import org.apache.iceberg.spark.SparkTableUtil; | ||
| import org.apache.iceberg.spark.SparkWriteOptions; | ||
|
|
@@ -2028,6 +2029,52 @@ public void testSnapshotProperty() { | |
| assertThat(table.currentSnapshot().summary()).containsKeys(commitMetricsKeys); | ||
| } | ||
|
|
||
| @TestTemplate | ||
|
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. can we add a test to cover if override session or explicit snapshot property collide with reserved
Collaborator
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. Thanks for raising this. I tried it, and a colliding key fails the commit today with That error comes from I'd rather not pin the current Guava error in a Spark test, because it would break as soon as core fixes #17009. Would you suggest a test that only checks that the commit fails and the table is unchanged? |
||
| void snapshotPropertyFromSessionConf() { | ||
| Table table = createTable(4); | ||
| withSQLConf( | ||
| ImmutableMap.of(SparkSQLProperties.SNAPSHOT_PROPERTY_PREFIX + "key", "session-value"), | ||
| () -> { | ||
| Result ignored = basicRewrite(table).execute(); | ||
| assertThat(table.currentSnapshot().summary()).containsEntry("key", "session-value"); | ||
| }); | ||
| } | ||
|
|
||
| @TestTemplate | ||
| void explicitSnapshotPropertyOverridesSessionConf() { | ||
| Table table = createTable(4); | ||
| withSQLConf( | ||
| ImmutableMap.of(SparkSQLProperties.SNAPSHOT_PROPERTY_PREFIX + "key", "session-value"), | ||
| () -> { | ||
| Result ignored = basicRewrite(table).snapshotProperty("key", "explicit-value").execute(); | ||
| assertThat(table.currentSnapshot().summary()).containsEntry("key", "explicit-value"); | ||
| }); | ||
| } | ||
|
|
||
| @TestTemplate | ||
| void snapshotPropertyAppliedToRemoveDanglingDeletesCommit() { | ||
| Table table = | ||
| TABLES.create( | ||
| SCHEMA, | ||
| SPEC, | ||
| ImmutableMap.of(TableProperties.FORMAT_VERSION, String.valueOf(formatVersion)), | ||
| tableLocation); | ||
| writeRecords(Lists.newArrayList(new ThreeColumnRecord(1, null, "AAAA"))); | ||
| writeEqDeleteRecord(table, "c1", 2, "c3", "CCCC"); | ||
| table.refresh(); | ||
|
|
||
| Result ignored = | ||
| basicRewrite(table) | ||
| .option(RewriteDataFiles.REMOVE_DANGLING_DELETES, "true") | ||
| .snapshotProperty("key", "value") | ||
| .execute(); | ||
|
|
||
| table.refresh(); | ||
| assertThat(table.currentSnapshot().summary()) | ||
| .containsEntry(SnapshotSummary.REMOVED_EQ_DELETE_FILES_PROP, "1") | ||
| .containsEntry("key", "value"); | ||
| } | ||
|
|
||
| @TestTemplate | ||
| public void testBinPackRewriterWithSpecificUnparitionedOutputSpec() { | ||
| Table table = createTable(10); | ||
|
|
||
Uh oh!
There was an error while loading. Please reload this page.
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.
noticed there's later section which also populate the custom metadata to a snapshot summary during a SQL execution
CommitMetadata.withCommitProperties, but list to per SQL writes. I think this complements with spark rewrite actions, might consider move later to it sits together? Some example on how explicit override session conf can also be helpful IMO.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.
Thanks, good suggestion. I added a paragraph right after the
CommitMetadataexample that covers the session property, where it applies. Please take a look.