manuzhang commented on code in PR #6773:
URL: https://github.com/apache/datafusion-comet/pull/6773#discussion_r4220740181
##########
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:
this comment is stale?
--
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]