andygrove opened a new pull request, #6247:
URL: https://github.com/apache/datafusion-comet/pull/6247

   ## Which issue does this PR close?
   
   Closes #5648.
   Closes #6115.
   
   ## Rationale for this change
   
   Neither native writer reserved the memory it buffers, so none of it was 
visible to Comet's memory pool or to Spark.
   
   - `ParquetWriterExec` holds its in-progress row group. It uses parquet-rs's 
default limit of 1Mi rows and no byte limit.
   - The native Iceberg writer holds an in-progress row group for every open 
file. A fanout write keeps a file open for every partition a task writes to. 
Each partition also holds up to 999 rows back, waiting to fill the 1000-row 
unit the rolling writer is fed in.
   
   A wide fanout write therefore grew in memory that no budget accounted for, 
until the container was killed. The JVM writers hold the same buffers on the 
JVM heap, which executors are sized for.
   
   #5648 asked for an audit of the Iceberg writer's buffers, and a test that a 
fanout write over many partitions with a small memory budget fails with a Comet 
out-of-memory error rather than a process-level one. #6115 asked the same of 
both writers, with the buffers charged to the task's pool. This PR does both.
   
   ## What changes are included in this PR?
   
   - Both writers register one `MemoryConsumer` per task, 
`ParquetWriterExec[N]` and `IcebergWriteExec[N]`, and `try_resize` it after 
every batch. It is one per task rather than one per file because every consumer 
registered with `fair_unified` lowers the share of every other consumer in the 
task.
   - `ParquetWriterExec` reserves `ArrowWriter::memory_size()`, which counts 
encoded pages, encoder buffers, dictionaries and Bloom filters. On the remote 
(HDFS) path it adds the staging buffer's capacity.
   - The Iceberg writer wraps iceberg-rust's `ParquetWriterBuilder` in a 
`MeteredParquetWriterBuilder`. That is the only hook, because the rolling and 
partitioning writers keep their file writers private. iceberg-rust's 
`ParquetWriter` exposes `current_written_size()`, which is the bytes flushed 
plus the in-progress row group's encoded size. It does not expose parquet-rs's 
`memory_size()`. So each open file reports `current_written_size()` capped at 
`write.parquet.row-group-size-bytes`. That is exact until the file flushes its 
first row group, and at most one row group high after it. The reservation also 
includes the rows each `RowPacer` holds back.
   - Neither writer can spill, so a resize the pool refuses fails the task, the 
way DataFusion's own parquet sink does. The JVM sees `CometNativeException: 
Additional allocation failed for IcebergWriteExec[0] ...`, a task failure Spark 
retries. The Iceberg writer's abort guard deletes the files the task had 
opened, as on any other task failure.
   - Docs:
     - `memory_management.md` has a new "Native writers" section.
     - The contributor `iceberg-writes.md` describes the metered builder.
     - The user guide's failure-handling section explains what a fanout write 
needs and how to make it fit.
     - The Iceberg-write review skill gets a checklist item.
   
   What this does not do:
   
   - Nothing sheds memory when the pool refuses. Flushing the row group early 
in `ParquetWriterExec`, or closing the largest open file in a fanout write, 
would let the write continue with a different file layout. Those can be 
follow-ups.
   - The Iceberg reservation misses what parquet-rs holds beyond the encoded 
estimate: dictionary hash tables, unencoded dictionary indices and buffer 
capacity. An exact figure needs an upstream accessor on iceberg-rust's 
`ParquetWriter`. Native Iceberg writes decline Bloom filters today, and #5724 
would need to reserve them itself.
   - There is no cap on open fanout writers. Accounting for them keeps the 
shape of iceberg-java's `FanoutDataWriter`, which also keeps every partition's 
file open.
   - `ParquetWriterExec` still writes 1Mi-row row groups with no byte limit 
(#5304), so a wide write needs a correspondingly large share of the pool.
   
   Behavior change: both native writers are off by default. A write whose 
buffers exceed the task's share of the pool now fails the task. Before, it 
succeeded in untracked memory, or got the executor killed.
   
   For reviewers: #6241 and #6238 also edit `iceberg_write.rs`. #6241's 
`PartitionWriterBuilder` builds a `ParquetWriterBuilder` per partition, which 
has to go through `MeteredParquetWriterBuilder` once the two meet.
   
   ## How are these changes tested?
   
   - Rust, `iceberg_write.rs`:
     - `a_fanout_write_reserves_the_rows_its_partitions_hold_back`: 900 rows 
held back in each of 1 and 8 partitions reserve 9.5 KB and 74 KB.
     - `the_reservation_follows_the_files_that_are_open`: eight one-unit 
partitions reserve 44 KB fanout, 5.5 KB clustered, and 5.9 KB when an 
unpartitioned write rolls through eight files.
     - `a_file_reserves_no_more_than_one_row_group`: a 295 KB file written with 
16 KiB row groups reserves exactly 16 KiB.
     - `a_write_the_pool_cannot_hold_fails_and_deletes_its_files`: fails with 
`ResourcesExhausted` naming the consumer, leaves no files, and gives back its 
reservation.
     - `an_open_file_gives_back_its_share_when_it_closes_or_is_dropped` and 
`pacer_counts_the_memory_of_the_rows_it_holds_back`.
   - Rust, `parquet_writer.rs`: the row group is reserved until the file 
closes, and a pool too small for it fails the write.
   - Mutations, each caught by at least one test: no resize; a resize without 
the held-back rows; a resize without the open files; no row-group cap; a file's 
share leaked on close; a share not given back on drop; the pacer's count not 
reset after a hand-over; only one fanout partition counted.
   - JVM. Both tests shrink the pool to about 4 MiB with 
`spark.comet.exec.memoryPool.fraction=0.002`.
     - `CometIcebergWriteActionSuite` "native acceleration: a fanout write that 
outgrows the memory pool fails its task" writes 64 partitions from one task. It 
checks the error, that nothing was committed, and that no data files were left 
(the writer deleted the 12 it had opened). Then the same write with the whole 
pool succeeds natively.
     - `CometParquetWriterSuite` "a row group the pool cannot hold fails its 
task with a native out-of-memory error" does the same for a single 10 MB row 
group.
     - Both fail with the resizes removed, which is main's behaviour.
   - Runs:
     - `cargo test -p datafusion-comet --lib` passes (503 tests), and clippy 
with `-D warnings` and `cargo fmt` are clean.
     - On Spark 4.1, `CometIcebergWriteActionSuite`, 
`CometIcebergRewriteActionSuite` and `CometIcebergSystemFunctionSuite` pass (96 
tests), and so does `CometParquetWriterSuite` (50 tests).
     - The two new JVM tests also pass on Spark 3.5.
     - I have not run Spark 3.4, 4.0 or the Iceberg Spark tests locally.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to