You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
Neither native writer charges its buffers to a memory pool. ParquetWriterExec (native/core/src/execution/operators/parquet_writer.rs) and the native Iceberg writer (native/core/src/execution/operators/iceberg_write.rs) both allocate through parquet-rs with no MemoryReservation. That memory is off-heap and invisible to Spark's TaskMemoryManager and to Comet's native pool. The JVM writers they replace allocate the same buffers on the JVM heap, where executors are already sized for them.
The unaccounted memory per task is roughly:
in-progress row group buffers (pages, dictionary pages, statistics) for each open file
one open file per partition for fanout and clustered writes, so everything above multiplies by the number of open partitions
Bloom filters, once feat: support bloom filters in native Iceberg writes #5724 lands: each enabled column allocates its full initial filter per open file up front. That is 1 MiB by default with no NDV set, and up to 128 MiB with write.parquet.bloom-filter-max-bytes. parquet-rs folds sparse filters at flush, but the fold truncates the length and keeps the vector's capacity, so it does not return memory.
As an example, a fanout Iceberg write with 10 Bloom-enabled columns over 200 open partitions holds about 2 GiB of filters alone, none of it tracked. When this exceeds what the executor has in reserve, the failure is a container OOM kill rather than a Comet or Spark out-of-memory error, and there is no spill or backpressure.
Describe the potential solution
Charge write buffers to the task's DataFusion MemoryPool. In off-heap mode that is the unified pool backed by Spark's off-heap memory through CometTaskMemoryManager, so writer buffers would share the same off-heap budget as native operators and Spark's own consumers, instead of growing outside it.
A sketch:
Give each open file writer a MemoryReservation from the task context and resize it to the writer's in-memory size after each batch and after each flush. DataFusion's own parquet sink uses this pattern with ArrowWriter::memory_size().
Check whether parquet-rs's reported size includes the Bloom filter allocation (via the column encoder's estimated_memory_size). If it does not, account for it separately from the configured sizes.
For the Iceberg writer, iceberg-rust's ParquetWriter wraps AsyncArrowWriter. It may need an upstream accessor for the in-progress size.
Decide what to do when a reservation cannot grow. Options include flushing the current row group early, closing the largest open file for fanout writes, or failing with a Comet out-of-memory error. Any of these is better than an OOM kill.
Add a test where a fanout write over many partitions with a small off-heap size fails with a Comet out-of-memory error, not a process-level one.
Additional context
#5648 is the audit for the Iceberg writer's buffers. This issue covers both native writers and proposes the shared-pool approach. Related: #5724 (Bloom filters in native Iceberg writes), #5649 (native Iceberg writes epic), #5212 (memory pool and accounting audit), #5997.
What is the problem the feature request solves?
Neither native writer charges its buffers to a memory pool.
ParquetWriterExec(native/core/src/execution/operators/parquet_writer.rs) and the native Iceberg writer (native/core/src/execution/operators/iceberg_write.rs) both allocate through parquet-rs with noMemoryReservation. That memory is off-heap and invisible to Spark'sTaskMemoryManagerand to Comet's native pool. The JVM writers they replace allocate the same buffers on the JVM heap, where executors are already sized for them.The unaccounted memory per task is roughly:
write.parquet.bloom-filter-max-bytes. parquet-rs folds sparse filters at flush, but the fold truncates the length and keeps the vector's capacity, so it does not return memory.As an example, a fanout Iceberg write with 10 Bloom-enabled columns over 200 open partitions holds about 2 GiB of filters alone, none of it tracked. When this exceeds what the executor has in reserve, the failure is a container OOM kill rather than a Comet or Spark out-of-memory error, and there is no spill or backpressure.
Describe the potential solution
Charge write buffers to the task's DataFusion
MemoryPool. In off-heap mode that is the unified pool backed by Spark's off-heap memory throughCometTaskMemoryManager, so writer buffers would share the same off-heap budget as native operators and Spark's own consumers, instead of growing outside it.A sketch:
MemoryReservationfrom the task context and resize it to the writer's in-memory size after each batch and after each flush. DataFusion's own parquet sink uses this pattern withArrowWriter::memory_size().estimated_memory_size). If it does not, account for it separately from the configured sizes.ParquetWriterwrapsAsyncArrowWriter. It may need an upstream accessor for the in-progress size.Additional context
#5648 is the audit for the Iceberg writer's buffers. This issue covers both native writers and proposes the shared-pool approach. Related: #5724 (Bloom filters in native Iceberg writes), #5649 (native Iceberg writes epic), #5212 (memory pool and accounting audit), #5997.