andygrove commented on code in PR #6773:
URL: https://github.com/apache/datafusion-comet/pull/6773#discussion_r4231475541


##########
native/core/src/execution/operators/iceberg_write.rs:
##########
@@ -1057,7 +1105,204 @@ async fn run_write_task(
     }
 }
 
-/// Enum-based dispatch over the three iceberg-rust partitioning writers, each 
paired with the
+/// A fanout write: like iceberg-rust's `FanoutWriter`, a data file writer 
open for every
+/// partition the task has written to, except that a partition can be closed 
before the task ends,
+/// to give back the memory it holds. Its next rows then open a new file, with 
the same properties.
+/// `FanoutWriter` keeps its writers private and closes them only all at once, 
so the fanout path
+/// keeps its own.
+///
+/// iceberg-java's fanout writer keeps every file open until the task ends, 
its buffers growing on
+/// the JVM heap. The native writer's buffers count against the task's memory 
pool instead, so when
+/// the pool refuses them, closing partitions early lets the write finish with 
more, smaller files
+/// rather than fail (see [`InnerWriter::reserve`]).
+struct FanoutPartitions {
+    builder: PartitionWriterBuilder,
+    /// Every partition the task has seen, in the order it first saw them.
+    partitions: Vec<FanoutPartition>,
+    /// Where each partition value sits in `partitions`.
+    index: HashMap<IcebergStruct, usize>,
+    /// Bytes the feeds hold back for their dictionary choice, across all 
partitions. Every
+    /// partition stays open, so the rows they hold back share one limit: once 
they reach it, every
+    /// partition still holding makes its choice from the rows it has. A task 
fanning out to many
+    /// partitions that each get less than a page would otherwise hold all of 
its rows back,
+    /// uncompressed, until it closed.
+    held_bytes: usize,
+    /// The data files of the partitions closed early.
+    closed: Vec<DataFile>,
+}
+
+/// One partition of a fanout write.
+struct FanoutPartition {
+    /// Kept so held rows can still be written out at close.

Review Comment:
   Yes, it was left over from before the refactor, when the feeds kept the key 
only so held rows could go out at close. The key now opens every file the 
partition writes, its first and any after an early close, and the comment says 
so in a7eb6cdce3.
   



##########
docs/source/contributor-guide/iceberg-writes.md:
##########
@@ -281,9 +282,16 @@ hold in memory. It hands `ParquetWriter` each file's 
`OutputFile` behind a `Coun
 writer counts the bytes that leave memory on their way to storage, and a file 
reports what it has
 written less those. When they leave depends on the storage (`StorageWrites`). 
After every batch
 `run_write_task` resizes the task's reservation to what the open files report 
plus the rows each
-`PartitionFeed` holds, for the dictionary choice or for pacing, and a resize 
the pool refuses fails
-the task. What that figure covers, and what it misses, is described under
-[Native writers](memory_management.md#native-writers).
+`PartitionFeed` holds, for the dictionary choice or for pacing 
(`InnerWriter::reserve`). When the
+pool refuses a fanout write, `reserve` writes out and closes the partitions 
holding the most, in
+their open file and their feed, until the resize succeeds, and 
`run_write_task` counts them in the
+`files_closed_early` metric. Each partition's files report to a child of the 
task's
+`OpenFileMemory`, which is how `reserve` finds them. That is why the fanout 
path uses
+`FanoutPartitions` rather than iceberg-rust's `FanoutWriter`, which cannot 
close one partition's
+writer. A closed partition's next rows open a new file with the properties its 
first file used,
+which the partition keeps. A write the pool still refuses, with nothing left 
to close, fails the

Review Comment:
   Not a fanout write. Everything it reserves belongs to one of its partitions, 
so once every partition holding memory is closed the resize is down to nothing, 
and shrinking a reservation can't be refused. The sentence was wrong to suggest 
a fanout write could still fail with nothing left to close. In a7eb6cdce3 it 
says the pool can't fail a fanout write, and that only an unpartitioned or 
clustered write, with no file it can close early, still fails its task, which 
`a_write_the_pool_cannot_hold_fails_and_deletes_its_files` and its counterpart 
in `CometIcebergWriteActionSuite` cover. The same commit fixes the matching 
sentence in `memory_management.md` and puts a comment on the `Err(refused)` at 
the end of `reserve`, which isn't reached.
   



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