Skip to content
Open
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
15 changes: 15 additions & 0 deletions docs/docs/spark-configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -206,6 +206,7 @@ val spark = SparkSession.builder()
| 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 |

@dramaticlly dramaticlly Oct 8, 2026 •

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.

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.

Copy link
Copy Markdown
Collaborator Author

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 CommitMetadata example that covers the session property, where it applies. Please take a look.

| spark.sql.iceberg.merge-schema | false | Enables modifying the table schema to match the write schema. Only adds missing columns |
| spark.sql.iceberg.view.schema-binding-mode | BINDING | Coercion applied to view columns: `BINDING` (widening only), `COMPENSATION` (any ANSI cast) |
| spark.sql.iceberg.report-column-stats | true | Report Puffin Table Statistics if available to Spark's Cost Based Optimizer. CBO must be enabled for this to be effective |
Expand Down Expand Up @@ -286,3 +287,17 @@ CommitMetadata.withCommitProperties(properties,
},
RuntimeException.class);
```

Custom metadata can also be added to snapshot summaries for an entire Spark session by setting properties with the `spark.sql.iceberg.snapshot-property.` prefix. The prefix is removed from each property. These properties apply to writes and to the `rewrite_data_files`, `rewrite_position_delete_files`, and `rewrite_manifests` procedures. Session properties act as defaults. A property set explicitly, through a write option, `CommitMetadata`, or `snapshotProperty()` on a Spark action, overrides a session property with the same key. Here is an example:

```sql
SET spark.sql.iceberg.snapshot-property.created-by=maintenance-job;
CALL catalog.system.rewrite_data_files(table => 'db.sample');
```

```java
SparkActions.get(spark)
.rewriteDataFiles(table)
.snapshotProperty("created-by", "backfill-job")
.execute();
```
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@
import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList;
import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap;
import org.apache.iceberg.relocated.com.google.common.collect.Lists;
import org.apache.iceberg.spark.SparkSQLProperties;
import org.apache.spark.sql.AnalysisException;
import org.apache.spark.sql.RowFactory;
import org.apache.spark.sql.catalyst.analysis.NoSuchTableException;
Expand Down Expand Up @@ -106,6 +107,22 @@ public void testRewriteManifestsNoOp() {
.hasSize(1);
}

@TestTemplate
void rewriteManifestsWithSessionSnapshotProperty() {
sql(
"CREATE TABLE %s (id bigint NOT NULL, data string) USING iceberg PARTITIONED BY (data)",
tableName);
sql("INSERT INTO TABLE %s VALUES (1, 'a'), (2, 'b'), (3, 'c'), (4, 'd')", tableName);
sql("ALTER TABLE %s SET TBLPROPERTIES ('commit.manifest.target-size-bytes' '1')", tableName);

withSQLConf(
ImmutableMap.of(SparkSQLProperties.SNAPSHOT_PROPERTY_PREFIX + "key", "session-value"),
() -> sql("CALL %s.system.rewrite_manifests('%s')", catalogName, tableIdent));

Table table = validationCatalog.loadTable(tableIdent);
assertThat(table.currentSnapshot().summary()).containsEntry("key", "session-value");
}

@TestTemplate
public void testRewriteLargeManifestsOnDatePartitionedTableWithJava8APIEnabled() {
withSQLConf(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@
import org.apache.iceberg.types.Types;
import org.apache.iceberg.util.ByteBuffers;
import org.apache.iceberg.util.Pair;
import org.apache.iceberg.util.PropertyUtil;
import org.apache.spark.SparkEnv;
import org.apache.spark.scheduler.ExecutorCacheTaskLocation;
import org.apache.spark.sql.SparkSession;
Expand Down Expand Up @@ -293,6 +294,19 @@ public static boolean caseSensitive(SparkSession spark) {
return Boolean.parseBoolean(spark.conf().get("spark.sql.caseSensitive"));
}

/**
* Returns the snapshot summary properties set in the session through {@link
* SparkSQLProperties#SNAPSHOT_PROPERTY_PREFIX}, with the prefix removed.
*
* @param spark a Spark session
* @return a map of snapshot summary property names to values
*/
public static Map<String, String> sessionSnapshotProperties(SparkSession spark) {
return PropertyUtil.propertiesWithPrefix(
JavaConverters.mapAsJavaMap(spark.conf().getAll()),
SparkSQLProperties.SNAPSHOT_PROPERTY_PREFIX);
}

public static List<String> executorLocations() {
BlockManager driverBlockManager = SparkEnv.get().blockManager();
List<BlockManagerId> executorBlockManagerIds = fetchPeers(driverBlockManager);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,6 @@
import org.apache.spark.sql.util.CaseInsensitiveStringMap;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import scala.collection.JavaConverters;

/**
* A class for common Iceberg configs for Spark writes.
Expand Down Expand Up @@ -282,10 +281,7 @@ public Map<String, String> extraSnapshotMetadata() {
Map<String, String> extraSnapshotMetadata = Maps.newHashMap();

// Add session configuration properties with SNAPSHOT_PROPERTY_PREFIX if necessary
extraSnapshotMetadata.putAll(
PropertyUtil.propertiesWithPrefix(
JavaConverters.mapAsJavaMap(sessionConf.getAll()),
SparkSQLProperties.SNAPSHOT_PROPERTY_PREFIX));
extraSnapshotMetadata.putAll(SparkUtil.sessionSnapshotProperties(spark));

// Add write options, overriding session configuration if necessary
extraSnapshotMetadata.putAll(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,14 +21,22 @@
import java.util.Map;
import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap;
import org.apache.iceberg.relocated.com.google.common.collect.Maps;
import org.apache.iceberg.spark.SparkUtil;
import org.apache.spark.sql.SparkSession;

abstract class BaseSnapshotUpdateSparkAction<ThisT> extends BaseSparkAction<ThisT> {

private final Map<String, String> summary = Maps.newHashMap();

protected BaseSnapshotUpdateSparkAction(SparkSession spark) {
this(spark, ImmutableMap.of());
}

protected BaseSnapshotUpdateSparkAction(
SparkSession spark, Map<String, String> snapshotProperties) {
super(spark);
summary.putAll(SparkUtil.sessionSnapshotProperties(spark));
summary.putAll(snapshotProperties);
}

public ThisT snapshotProperty(String property, String value) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
import java.io.IOException;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.stream.StreamSupport;
import org.apache.iceberg.DeleteFile;
import org.apache.iceberg.FileScanTask;
Expand Down Expand Up @@ -53,8 +54,9 @@ class RemoveDanglingDeletesSparkAction
private final Table table;
private final RemoveDanglingDeleteFilesAction action;

protected RemoveDanglingDeletesSparkAction(SparkSession spark, Table table) {
super(spark);
RemoveDanglingDeletesSparkAction(
SparkSession spark, Table table, Map<String, String> snapshotProperties) {
super(spark, snapshotProperties);
this.table = table;
this.action = new RemoveDanglingDeleteFilesAction(table, this::findDanglingDeletes);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -487,7 +487,9 @@ private String jobDesc(
private RewriteDataFiles.Result executeRemoveDanglingDeletes(
ImmutableRewriteDataFiles.Result rewriteResult) {
RemoveDanglingDeletesSparkAction.Result result =
new RemoveDanglingDeletesSparkAction(spark(), table).toBranch(branch).execute();
new RemoveDanglingDeletesSparkAction(spark(), table, commitSummary())
.toBranch(branch)
.execute();
return rewriteResult.withRemovedDeleteFilesCount(
rewriteResult.removedDeleteFilesCount() + Iterables.size(result.removedDeleteFiles()));
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
import org.apache.iceberg.actions.ComputePartitionStats;
import org.apache.iceberg.actions.ComputeTableStats;
import org.apache.iceberg.actions.RemoveDanglingDeleteFiles;
import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap;
import org.apache.iceberg.spark.Spark3Util;
import org.apache.iceberg.spark.Spark3Util.CatalogAndIdentifier;
import org.apache.spark.sql.SparkSession;
Expand Down Expand Up @@ -112,7 +113,7 @@ public ComputePartitionStats computePartitionStats(Table table) {

@Override
public RemoveDanglingDeleteFiles removeDanglingDeleteFiles(Table table) {
return new RemoveDanglingDeletesSparkAction(spark, table);
return new RemoveDanglingDeletesSparkAction(spark, table, ImmutableMap.of());
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
import org.apache.iceberg.Table;
import org.apache.iceberg.actions.RemoveDanglingDeleteFiles;
import org.apache.iceberg.actions.TestRemoveDanglingDeleteFilesAction;
import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap;
import org.apache.iceberg.spark.TestBase;
import org.apache.spark.sql.SparkSession;
import org.junit.jupiter.api.AfterAll;
Expand Down Expand Up @@ -55,6 +56,6 @@ public static void stopSpark() {

@Override
protected RemoveDanglingDeleteFiles removeDanglingDeleteFiles(Table actionTable) {
return new RemoveDanglingDeletesSparkAction(spark, actionTable);
return new RemoveDanglingDeletesSparkAction(spark, actionTable, ImmutableMap.of());
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -2028,6 +2029,52 @@ public void testSnapshotProperty() {
assertThat(table.currentSnapshot().summary()).containsKeys(commitMetricsKeys);
}

@TestTemplate

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.

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.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The 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 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?

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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,7 @@
import org.apache.iceberg.relocated.com.google.common.collect.Iterables;
import org.apache.iceberg.relocated.com.google.common.collect.Lists;
import org.apache.iceberg.relocated.com.google.common.collect.Maps;
import org.apache.iceberg.spark.SparkSQLProperties;
import org.apache.iceberg.spark.SparkTableUtil;
import org.apache.iceberg.spark.SparkWriteOptions;
import org.apache.iceberg.spark.TestBase;
Expand Down Expand Up @@ -389,6 +390,31 @@ public void testRewriteSmallManifestsNonPartitionedTable() {
assertThat(actualRecords).as("Rows must match").isEqualTo(expectedRecords);
}

@TestTemplate
void snapshotPropertyFromSessionConf() {
Table table =
TABLES.create(
SCHEMA,
PartitionSpec.unpartitioned(),
ImmutableMap.of(TableProperties.FORMAT_VERSION, String.valueOf(formatVersion)),
tableLocation);
writeRecords(Lists.newArrayList(new ThreeColumnRecord(1, null, "AAAA")));
writeRecords(Lists.newArrayList(new ThreeColumnRecord(2, "CCCC", "CCCC")));
table.refresh();

withSQLConf(
ImmutableMap.of(SparkSQLProperties.SNAPSHOT_PROPERTY_PREFIX + "key", "session-value"),
() -> {
SparkActions.get()
.rewriteManifests(table)
.rewriteIf(manifest -> true)
.option(RewriteManifestsSparkAction.USE_CACHING, useCaching)
.execute();
table.refresh();
assertThat(table.currentSnapshot().summary()).containsEntry("key", "session-value");
});
}

@TestTemplate
public void testRewriteManifestsWithCommitStateUnknownException() {
PartitionSpec spec = PartitionSpec.unpartitioned();
Expand Down
Loading