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]