Repository navigation
Spark 4.2: Apply session snapshot properties to maintenance actions - #18406
anuragmantri wants to merge 4 commits into
Conversation
Co-authored-by: Pucheng Yang <pyang@pinterest.com>
|
@mxm @dramaticlly - Could you please provide a review? |
mxm
left a comment
There was a problem hiding this comment.
Thanks @anuragmantri! Looks good to me. Just one question.
| summary.putAll( | ||
| PropertyUtil.propertiesWithPrefix( | ||
| JavaConverters.mapAsJavaMap(spark.conf().getAll()), | ||
| SparkSQLProperties.SNAPSHOT_PROPERTY_PREFIX)); |
There was a problem hiding this comment.
Is this the only place where we need to apply those?
There was a problem hiding this comment.
For Spark actions, yes. BaseSnapshotUpdateSparkAction is the base class of all four actions that commit snapshots. RewriteDataFilesSparkAction, RewritePositionDeleteFilesSparkAction, and RemoveDanglingDeletesSparkAction commit through commitSummary(), and RewriteManifestsSparkAction commits through commit(). The tests cover both paths.
A few other places commit snapshots without picking this up:
add_files,migrate, andsnapshotappend throughSparkTableUtil.cherrypick_snapshotandpublish_changesgo throughManageSnapshots, which has no way to set summary properties, so they would need an API change first.- Metadata-only
DELETEis a write-path gap, tracked in DELETE operations don't apply custom snapshot properties from session config #15060.
I'd like to keep this PR to the actions and handle SparkTableUtil in a follow-up.
There was a problem hiding this comment.
on a tangentially related note, I think we might drop the snapshotProperty on RevemoDanglingDeletesSparkAction before but not anymore after this change.
SparkActions.get(spark).rewriteDataFiles(table)
.option(RewriteDataFiles.REMOVE_DANGLING_DELETES, "true")
.snapshotProperty("audit-id", "A")
.execute();
There was a problem hiding this comment.
also I think this duplicates with SparkWriteConf to retrieve the session level snapshot properties, might consider a common method in SparkUtil?
public static Map<String, String> sessionSnapshotProperties(SparkSession spark) {
return PropertyUtil.propertiesWithPrefix(
JavaConverters.mapAsJavaMap(spark.conf().getAll()),
SparkSQLProperties.SNAPSHOT_PROPERTY_PREFIX);
}There was a problem hiding this comment.
Good catch on the dangling deletes commit. Explicit properties were still dropped there, because only session properties reached the nested action. That also meant an explicit value won on the rewrite commit while the session value won on the dangling deletes commit.
I also moved the session lookup into SparkUtil.sessionSnapshotProperties, which BaseSnapshotUpdateSparkAction and SparkWriteConf now share, as you suggested.
| assertThat(table.currentSnapshot().summary()).containsKeys(commitMetricsKeys); | ||
| } | ||
|
|
||
| @TestTemplate |
There was a problem hiding this comment.
can we add a test to cover if override session or explicit snapshot property collide with reserved added-data-files, just to pin down the behavior.
There was a problem hiding this comment.
Thanks for raising this. I tried it, and a colliding key fails the commit today with IllegalArgumentException: Multiple entries with same key: added-data-files=1 and added-data-files=999. The rewrite work is thrown away and the table is left unchanged.
That error comes from SnapshotSummary.Builder in core, and it happens the same way through snapshotProperty() and through writes with this session property, so this PR doesn't change it. #17009 tracks the reserved-key problem in core, and #17107 tried to fix it there before it went stale.
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?
| summary.putAll( | ||
| PropertyUtil.propertiesWithPrefix( | ||
| JavaConverters.mapAsJavaMap(spark.conf().getAll()), | ||
| SparkSQLProperties.SNAPSHOT_PROPERTY_PREFIX)); |
There was a problem hiding this comment.
on a tangentially related note, I think we might drop the snapshotProperty on RevemoDanglingDeletesSparkAction before but not anymore after this change.
SparkActions.get(spark).rewriteDataFiles(table)
.option(RewriteDataFiles.REMOVE_DANGLING_DELETES, "true")
.snapshotProperty("audit-id", "A")
.execute();
| summary.putAll( | ||
| PropertyUtil.propertiesWithPrefix( | ||
| JavaConverters.mapAsJavaMap(spark.conf().getAll()), | ||
| SparkSQLProperties.SNAPSHOT_PROPERTY_PREFIX)); |
There was a problem hiding this comment.
also I think this duplicates with SparkWriteConf to retrieve the session level snapshot properties, might consider a common method in SparkUtil?
public static Map<String, String> sessionSnapshotProperties(SparkSession spark) {
return PropertyUtil.propertiesWithPrefix(
JavaConverters.mapAsJavaMap(spark.conf().getAll()),
SparkSQLProperties.SNAPSHOT_PROPERTY_PREFIX);
}| RemoveDanglingDeletesSparkAction removeDanglingDeletesAction = | ||
| new RemoveDanglingDeletesSparkAction(spark(), table).toBranch(branch); | ||
| commitSummary().forEach(removeDanglingDeletesAction::snapshotProperty); |
There was a problem hiding this comment.
This looks error-prone. Could we pass commitSummary() through a RemoveDanglingDeletesSparkAction constructor?
There was a problem hiding this comment.
Agreed, that's cleaner. I made this change.
| | spark.sql.iceberg.executor-cache.max-entry-size | 67108864 (64MB) | Max size per cache entry (bytes) | | ||
| | spark.sql.iceberg.executor-cache.max-total-size | 134217728 (128MB) | Max total executor cache size (bytes) | | ||
| | spark.sql.iceberg.executor-cache.locality.enabled | false | Enables locality-aware executor cache usage | | ||
| | spark.sql.iceberg.snapshot-property._custom-key_ | null | Adds an entry with custom-key and corresponding value to the summary of snapshots committed by writes and by the `rewrite_data_files`, `rewrite_position_delete_files`, and `rewrite_manifests` procedures. Write options and properties set on an action take precedence | |
There was a problem hiding this comment.
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.
Thanks, good suggestion. I added a paragraph right after the CommitMetadata example that covers the session property, where it applies. Please take a look.
| protected RemoveDanglingDeletesSparkAction(SparkSession spark, Table table) { | ||
| super(spark); | ||
| this(spark, table, ImmutableMap.of()); | ||
| } |
There was a problem hiding this comment.
Should we remove this constructor?
| RuntimeException.class); | ||
| ``` | ||
|
|
||
| Snapshot properties can also be set for a whole session with `spark.sql.iceberg.snapshot-property._custom-key_`. They apply to writes and to the `rewrite_data_files`, `rewrite_position_delete_files`, and `rewrite_manifests` procedures. A `snapshot-property.` write option or `CommitMetadata` overrides a session property with the same key. When calling an action through the Java API, `snapshotProperty()` does too. |
There was a problem hiding this comment.
nit: maybe highlight that write option or explicit argument override the session property in general?
c4fc2be to
4588620
Compare
anuragmantri
left a comment
There was a problem hiding this comment.
Thanks @mxm and @dramaticlly! I pushed another change, which removes the extra constructor and adds the snapshot property docs next to CommitMetadata. @mxm, could you take another look, since this landed after your approval?
| protected RemoveDanglingDeletesSparkAction(SparkSession spark, Table table) { | ||
| super(spark); | ||
| this(spark, table, ImmutableMap.of()); | ||
| } |
| | spark.sql.iceberg.executor-cache.max-entry-size | 67108864 (64MB) | Max size per cache entry (bytes) | | ||
| | spark.sql.iceberg.executor-cache.max-total-size | 134217728 (128MB) | Max total executor cache size (bytes) | | ||
| | spark.sql.iceberg.executor-cache.locality.enabled | false | Enables locality-aware executor cache usage | | ||
| | spark.sql.iceberg.snapshot-property._custom-key_ | null | Adds an entry with custom-key and corresponding value to the summary of snapshots committed by writes and by the `rewrite_data_files`, `rewrite_position_delete_files`, and `rewrite_manifests` procedures. Write options and properties set on an action take precedence | |
There was a problem hiding this comment.
Thanks, good suggestion. I added a paragraph right after the CommitMetadata example that covers the session property, where it applies. Please take a look.
SQL users can tag write snapshots with
spark.sql.iceberg.snapshot-property.*(#14545), but snapshots fromrewrite_data_files,rewrite_position_delete_files, andrewrite_manifestsignore it, and procedures have no other way to set snapshot properties. This applies the session properties inBaseSnapshotUpdateSparkAction, through aSparkUtilhelper shared withSparkWriteConf. Properties set throughsnapshotProperty()take precedence, as write options do for writes, and now also reach the dangling deletes commit thatrewrite_data_filesmakes. Sessions that already set these properties for writes will now also tag maintenance snapshots.Takes over #15842 (credit @puchengy) and adds the commit()-path test @anoopj asked for, plus the missing docs entry.
Test plan:
CALL rewrite_manifestswith the session property set tags the snapshot.AI Disclosure