What is the problem the feature request solves?
With #6773, a native fanout Iceberg write that the memory pool can't hold writes out and closes the partitions holding the most memory, and a closed partition's next rows open a new file. The write finishes where it used to fail its task, but with more, smaller files. When every partition keeps getting rows, the count depends on how far short the pool falls, not on how many partitions the task writes. That is the usual case for a fanout write, since Iceberg's default hash distribution sends a partition's rows to one task but doesn't sort them. Files per partition, with rows taking turns across partitions and the pool set to a fraction of what the same write reserves with room to spare (measured on #6773):
| pool |
16 partitions |
64 partitions |
256 partitions |
| 1/2 |
1.69 |
1.66 |
1.66 |
| 1/4 |
2.62 |
2.55 |
2.56 |
| 1/8 |
4.38 |
4.31 |
4.27 |
| 1/16 |
7.44 |
7.36 |
7.29 |
Each extra file costs a manifest entry, a file open in every scan that reads it, and eventually a compaction. iceberg-java's fanout writer keeps every file open until the task ends, so it writes one file per partition, rolling at the target size, as long as the JVM heap holds all of them.
How often this happens depends on the write. On Spark 3.5 and later, the default hash distribution gives each task only some of the partitions and asks AQE for tasks of an advisory size, so a task's open files usually fit its share of the pool. Early closes mostly come from write.distribution-mode=none, where a task's rows span every partition, from very large tasks, or from a small spark.memory.offHeap.size. The files closed early to free memory metric on CometIcebergWrite shows when it happens.
Describe the potential solution
When the pool refuses a fanout write's reservation, stop fanning out and sort the rest of the task's rows by partition. Spark's own file writer does this past spark.sql.maxConcurrentOutputFileWriters, an internal setting that is off by default: DynamicPartitionDataConcurrentWriter sorts the remaining rows by partition and writes them one partition at a time, continuing a partition's open writer if it has one and closing it when the partition's run ends.
For the native writer that would mean:
- Feeding the rest of the input through a spilling sort keyed on the partition values
PartitionSplitter computes, with the sort's memory reserved from the same pool, so that it spills when the pool refuses.
- Writing the sorted rows the way a clustered write does: each partition's rows in one run, into the partition's open file if it still has one, closing the file when the run ends.
Each partition then ends with about one file per target size, plus at most one from before the switch, at the cost of spilling the rest of the task's rows.
Open questions:
- When to switch. Sorting on the first refusal spills everything that follows, even for a task with little input left, which a few early closes would have handled more cheaply. The writer doesn't know how much input remains.
- Where the sort gets its memory. The open files hold what the pool just refused. Closing every open partition at the switch frees all of it for the sort, at the cost of one extra file per partition. Keeping them open leaves the sort to spill from the start.
- Where the sort lives: a DataFusion
SortExec over the remaining stream inside IcebergWriteExec, with the partition values appended as columns and dropped before the rows are written, or something the planner sets up.
Additional context
Raised in review of #6773 (comment).
Until then, a write that closes many partitions early has three workarounds. write.spark.fanout.enabled=false makes Spark sort every task's rows by partition up front, so the task keeps one file open at a time. A larger spark.memory.offHeap.size lets more partitions stay open. Iceberg's rewrite_data_files compacts the small files after the write.
Considered and set aside: flushing a partition's in-progress row group instead of closing its file would keep one file per partition, with smaller row groups. It needs iceberg-rust's ParquetWriter to expose a flush, and on S3 and GCS OpenDAL keeps a part of at least 5 MiB per open file until the next part completes or the file closes, so flushing alone doesn't bound memory over many partitions.
What is the problem the feature request solves?
With #6773, a native fanout Iceberg write that the memory pool can't hold writes out and closes the partitions holding the most memory, and a closed partition's next rows open a new file. The write finishes where it used to fail its task, but with more, smaller files. When every partition keeps getting rows, the count depends on how far short the pool falls, not on how many partitions the task writes. That is the usual case for a fanout write, since Iceberg's default hash distribution sends a partition's rows to one task but doesn't sort them. Files per partition, with rows taking turns across partitions and the pool set to a fraction of what the same write reserves with room to spare (measured on #6773):
Each extra file costs a manifest entry, a file open in every scan that reads it, and eventually a compaction. iceberg-java's fanout writer keeps every file open until the task ends, so it writes one file per partition, rolling at the target size, as long as the JVM heap holds all of them.
How often this happens depends on the write. On Spark 3.5 and later, the default hash distribution gives each task only some of the partitions and asks AQE for tasks of an advisory size, so a task's open files usually fit its share of the pool. Early closes mostly come from
write.distribution-mode=none, where a task's rows span every partition, from very large tasks, or from a smallspark.memory.offHeap.size. Thefiles closed early to free memorymetric onCometIcebergWriteshows when it happens.Describe the potential solution
When the pool refuses a fanout write's reservation, stop fanning out and sort the rest of the task's rows by partition. Spark's own file writer does this past
spark.sql.maxConcurrentOutputFileWriters, an internal setting that is off by default:DynamicPartitionDataConcurrentWritersorts the remaining rows by partition and writes them one partition at a time, continuing a partition's open writer if it has one and closing it when the partition's run ends.For the native writer that would mean:
PartitionSplittercomputes, with the sort's memory reserved from the same pool, so that it spills when the pool refuses.Each partition then ends with about one file per target size, plus at most one from before the switch, at the cost of spilling the rest of the task's rows.
Open questions:
SortExecover the remaining stream insideIcebergWriteExec, with the partition values appended as columns and dropped before the rows are written, or something the planner sets up.Additional context
Raised in review of #6773 (comment).
Until then, a write that closes many partitions early has three workarounds.
write.spark.fanout.enabled=falsemakes Spark sort every task's rows by partition up front, so the task keeps one file open at a time. A largerspark.memory.offHeap.sizelets more partitions stay open. Iceberg'srewrite_data_filescompacts the small files after the write.Considered and set aside: flushing a partition's in-progress row group instead of closing its file would keep one file per partition, with smaller row groups. It needs iceberg-rust's
ParquetWriterto expose a flush, and on S3 and GCS OpenDAL keeps a part of at least 5 MiB per open file until the next part completes or the file closes, so flushing alone doesn't bound memory over many partitions.