Skip to content

Charge native write buffers to the shared off-heap memory pool #6115

Description

@andygrove

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 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.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions