Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion .ai/skills/review-comet-iceberg-write-pr/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -156,7 +156,9 @@ Read the ownership table in the contributor guide before reviewing any change ne
`MeteredParquetWriterBuilder`, and rows the writer holds outside iceberg-rust (the
`PartitionFeed`s, or anything a new feed holds back) must be added to what `run_write_task` reserves. A
builder that constructs `ParquetWriterBuilder` directly leaves its files out of the task's
memory reservation, so a wide fanout write grows past the pool instead of failing its task.
memory reservation, so a wide fanout write grows past the pool instead of closing
partitions early or failing its task. A fanout file must also report to its partition's
`OpenFileMemory`, or the partitions holding the most are not the ones closed.
A storage scheme newly supported for writes needs its entry in
`StorageWrites::for_location`: a file reports its flushed row groups until its storage has
written them out, which a local file does at once and an object store only part by part.
Expand Down
40 changes: 24 additions & 16 deletions docs/source/contributor-guide/iceberg-writes.md
Original file line number Diff line number Diff line change
Expand Up @@ -223,7 +223,7 @@ per task, and `PartitionWriterBuilder` builds the rest of the stack once per par
partition's Parquet writer properties:

```text
UnpartitionedWriter | FanoutWriter | ClusteredWriter
UnpartitionedWriter | FanoutPartitions | ClusteredWriter
-> PartitionWriterBuilder, per partition:
ParquetWriterBuilder -> RollingFileWriterBuilder -> DataFileWriterBuilder
```
Expand Down Expand Up @@ -270,8 +270,9 @@ Points where Comet adapts iceberg-rust to match iceberg-java:
- **Field ids and casting.** `decorate_batch_with_field_ids` casts each batch to the
field-id-annotated Arrow schema derived from the Iceberg schema, with `safe: false`, so a type
mismatch fails the task instead of writing NULLs.
- **Deterministic output order.** iceberg-rust's `FanoutWriter` closes its writers out of a
`HashMap`, so the fanout path sorts its `DataFile`s by path before returning them
- **Deterministic output order.** `FanoutPartitions`, Comet's version of iceberg-rust's
`FanoutWriter`, closes its files in an order that depends on how the rows arrived and on which
partitions it closed early, so the fanout path sorts its `DataFile`s by path before returning them
([#5776](https://github.com/apache/datafusion-comet/issues/5776)). Manifest order becomes the row
order of an unordered read, so any map iteration that reaches the output needs the same care.

Expand All @@ -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

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can the pool ever refuse write now?

task, as an unpartitioned or clustered write always does. What the reserved figure covers, and what it misses,
is described under [Native writers](memory_management.md#native-writers).

`FileIO` comes from `load_file_io` in `iceberg_common.rs`, shared with the native scan. It picks
the storage backend from the data location's scheme and wires in Comet's S3 credential bridge when
Expand Down Expand Up @@ -411,16 +419,16 @@ the writer: run the write suites and the Iceberg Spark tests. The pin policy is

## Testing

| Suite | What it covers |
| ---------------------------------------------------------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| `CometIcebergWriteActionSuite` | End-to-end writes through the split plan and the native writer: parity with iceberg-java, row-level DML, partition evolution, file order, cleanup on task and job failure, a fanout write outgrowing the memory pool, AQE re-planning. |
| `CometIcebergWriteDetectionSuite` | One case per eligibility rule, accepted and declined. |
| `CometIcebergSystemFunctionSuite` | Native `bucket`, `truncate`, `years`/`months`/`days`/`hours`, which keep a partitioned write's hash distribution and sort native end to end. |
| `CometIcebergRewriteActionSuite` | Iceberg's `rewrite_data_files` with the split plan and the native writer. |
| `IcebergWriteProtoTranslationSuite` | Translation of properties into `IcebergParquetWriteSettings` and the writer mode. |
| Rust tests in `iceberg_write.rs` and in `iceberg_partition_*.rs` | File rolling on the 1000-row grid, fanout order, clustered input checks, cleanup guard, manifest round trip, memory reservation, partition path rendering, partition values past `chrono`'s calendar. |
| Rust tests in `iceberg_dictionary.rs` | The per-column dictionary choice against answers recorded from parquet-mr, including columns either side of its cut-off. |
| `CometIcebergWriteBenchmark` | Native versus iceberg-java for unpartitioned, clustered, fanout and copy-on-write delete writes. It checks each arm's plan before timing it. |
| Suite | What it covers |
| ---------------------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ |
| `CometIcebergWriteActionSuite` | End-to-end writes through the split plan and the native writer: parity with iceberg-java, row-level DML, partition evolution, file order, cleanup on task and job failure, writes outgrowing the memory pool, AQE re-planning. |
| `CometIcebergWriteDetectionSuite` | One case per eligibility rule, accepted and declined. |
| `CometIcebergSystemFunctionSuite` | Native `bucket`, `truncate`, `years`/`months`/`days`/`hours`, which keep a partitioned write's hash distribution and sort native end to end. |
| `CometIcebergRewriteActionSuite` | Iceberg's `rewrite_data_files` with the split plan and the native writer. |
| `IcebergWriteProtoTranslationSuite` | Translation of properties into `IcebergParquetWriteSettings` and the writer mode. |
| Rust tests in `iceberg_write.rs` and in `iceberg_partition_*.rs` | File rolling on the 1000-row grid, fanout order, clustered input checks, cleanup guard, manifest round trip, memory reservation, partition path rendering, partition values past `chrono`'s calendar. |
| Rust tests in `iceberg_dictionary.rs` | The per-column dictionary choice against answers recorded from parquet-mr, including columns either side of its cut-off. |
| `CometIcebergWriteBenchmark` | Native versus iceberg-java for unpartitioned, clustered, fanout and copy-on-write delete writes. It checks each arm's plan before timing it. |

The Comet suites run against the Iceberg version each Spark profile pins in `spark/pom.xml`: 1.5.2
for Spark 3.4, 1.8.1 for 3.5, 1.10.0 for 4.0 and 4.2, and 1.11.0 for 4.1. Only the default profile
Expand Down
19 changes: 12 additions & 7 deletions docs/source/contributor-guide/memory_management.md
Original file line number Diff line number Diff line change
Expand Up @@ -378,9 +378,11 @@ An operator that never calls `try_grow` is invisible to the pool no matter how m
### Native writers

Both native writers reserve what they hold between batches through a single consumer per task,
`ParquetWriterExec[N]` or `IcebergWriteExec[N]`, resized after every batch. Neither can spill, so
when the pool refuses a resize the task fails with a `CometNativeException` whose message starts
`Additional allocation failed for` and names the consumer. That is a task failure Spark can retry.
`ParquetWriterExec[N]` or `IcebergWriteExec[N]`, resized after every batch. Neither can spill. A
fanout Iceberg write can give memory back by closing partitions early, and does so when the pool
refuses a resize, as described below. Otherwise, when the pool refuses a resize, the task fails
with a `CometNativeException` whose message starts `Additional allocation failed for` and names the
consumer. That is a task failure Spark can retry.
Unreserved, the same memory would count only toward the container limit, where exceeding it kills
the executor.

Expand All @@ -400,10 +402,13 @@ the executor.
it has flushed since its last part was uploaded. At the default row-group size
(`write.parquet.row-group-size-bytes`, 128 MiB) that is the last row group it flushed. The
reservation also covers the rows each partition holds back, first for its dictionary choice and
then until they fill the 1000-row unit the rolling writer is fed in. It does not cover what
parquet-rs holds beyond the encoded size: dictionary hash tables, unencoded dictionary indices
and buffer capacity. Native writes decline Bloom filters today, and neither figure would include
them.
then until they fill the 1000-row unit the rolling writer is fed in. When the pool refuses a
fanout write's resize, the writer writes out and closes the partitions holding the most, file
and held rows together, until the resize succeeds, so the write ends with more files rather than
failing. Each file also reports to its partition's total, which is how the writer finds them.
The reservation does not cover what parquet-rs holds beyond the encoded size: dictionary hash
tables, unencoded dictionary indices and buffer capacity. Native writes decline Bloom filters
today, and neither figure would include them.

The writers register one consumer per task rather than one per open file because every consumer
registered with `fair_unified` lowers the share of every other consumer in the task. A consumer per
Expand Down
34 changes: 23 additions & 11 deletions docs/source/user-guide/latest/iceberg-writes.md
Original file line number Diff line number Diff line change
Expand Up @@ -310,13 +310,21 @@ The native writer's buffers are charged to Comet's memory pool, the off-heap bud
native operators draw on, where iceberg-java's buffers sit on the JVM heap. A fanout write keeps a
data file open for every partition a task writes to. Each open file holds the row group it is
writing in memory, up to `write.parquet.row-group-size-bytes`, and on S3 or GCS also the last row
group it flushed, which is uploaded once the next one is complete or the file closes. So a task
writing to many partitions needs memory in proportion to them. When the pool cannot grant it, the
task fails with a `CometNativeException` reading `Additional allocation failed for IcebergWriteExec`
instead of exceeding the executor's memory, and Spark retries it like any other task failure. Such a
write fits in less memory with the fanout writer disabled (`write.spark.fanout.enabled=false`):
Spark then sorts each task's rows by partition, and the task keeps one file open at a time. A
smaller row-group size also helps. Otherwise the write needs a larger `spark.memory.offHeap.size`.
group it flushed, which is uploaded once the next one is complete or the file closes. Each
partition also holds the rows that have not reached its file yet. So a task writing to many
partitions needs memory in proportion to them. When the pool cannot grant it, the task writes out
and closes the partitions holding the most memory until what is left fits, and a closed
partition's next rows open a new file. The write then finishes with more, smaller files than
iceberg-java's would, and the `files closed early to free memory` metric of its
`CometIcebergWrite` operator counts the files it closed early. Disabling the fanout writer
(`write.spark.fanout.enabled=false`) avoids them: Spark then sorts each task's rows by partition,
and the task keeps one file open at a time. A larger `spark.memory.offHeap.size` also helps.

A write that keeps one file open, unpartitioned or clustered, has no partition to close. When that
file outgrows the pool, the task fails with a `CometNativeException` reading
`Additional allocation failed for IcebergWriteExec` instead of exceeding the executor's memory, and
Spark retries it like any other task failure. A smaller `write.parquet.row-group-size-bytes` or a
larger `spark.memory.offHeap.size` lets such a write fit.

Partial results are never committed. The commit set is exactly the commit messages returned by
successful tasks — a failed task contributes none — and if the job fails, the driver-side
Expand Down Expand Up @@ -411,17 +419,21 @@ a data file but not what any reader computes from it:
they may cross the target several grid steps apart, and the resulting files can differ in row
count by an arbitrary number of 1000-row blocks. Do not rely on file-layout parity between the
two writers; rely only on each file rolling on its own 1000-row boundary.
- A fanout write that the memory pool cannot hold closes partitions before the task ends (see
[Failure handling](#failure-handling)), where iceberg-java's fanout writer keeps every file open
until then. So a partition can have more, smaller files than iceberg-java writes, and a file
closed early ends off the 1000-row grid, as the last file of a task does. A partition closed
before its rows filled a first page makes its dictionary choice from the rows it has.
- A fanout write lists a task's data files in file-path order, where iceberg-java lists them in
its own `StructLikeMap` iteration order. Both are stable across runs, and neither is a
documented ordering, but the manifest entry order becomes the scan-task order and so the row
order of an unordered `SELECT *`. On a format-version 3 table it also decides which row ids the
commit gives each file's rows, so the same rows can get different `_row_id` values from the two
writers. The ids are unique either way, and across tasks iceberg-java's own assignment already
depends on the order in which the tasks finish. Only the sorted order is reproducible on the
native path: iceberg-rust's `FanoutWriter` closes its per-partition writers out of a `HashMap`,
which under Rust's per-process `RandomState` would otherwise give a different order on every
run. Clustered and unpartitioned writes append in creation order on both paths and are
unaffected.
native path: the order the native writer closes its files in depends on the order rows arrive in
and on which partitions it closes early to free memory. Clustered and unpartitioned writes append
in creation order on both paths and are unaffected.
- Compressed page bytes are implementation-defined: the codec and any explicit level are
translated, but parquet-rs and parquet-mr embed different encoder implementations and
defaults (zstd default levels, LZ4 framing), so byte-identical output is not achievable even
Expand Down
Loading
Loading