Skip to content
Merged
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
12 changes: 7 additions & 5 deletions docs/source/contributor-guide/iceberg-writes.md
Original file line number Diff line number Diff line change
Expand Up @@ -59,13 +59,15 @@ gets a failed query instead of a fallback.
## The Split-Operator Plan

Spark plans an Iceberg write as one physical operator (`AppendDataExec`, `ReplaceDataExec` and so
on) that runs the input query, writes the files, and commits, all outside AQE. The split plan
replaces it with two operators so that the write's input becomes an ordinary query stage.
on) that writes the files and commits. It is a `V2CommandExec`, so Spark's
`InsertAdaptiveSparkPlan` wraps its input query in AQE but leaves the operator itself outside. The
split plan replaces it with two operators so that data-file writing moves inside AQE, apart from
the commit.

```text
IcebergCommit driver: collect task commit messages, BatchWrite.commit
+- IcebergWrite executors: write data files, emit one commit message per task
+- <input query> scans, projects, exchanges, sorts; now visible to AQE and Comet
IcebergCommit driver: collect task commit messages, BatchWrite.commit; outside AQE
+- IcebergWrite executors: write data files, emit one commit message per task; inside AQE
+- <input query> scans, projects, exchanges, sorts; inside AQE with or without the split
```

| Component | Location | Role |
Expand Down
11 changes: 7 additions & 4 deletions docs/source/user-guide/latest/iceberg-writes.md
Original file line number Diff line number Diff line change
Expand Up @@ -25,9 +25,12 @@ against your own workloads.
## Overview

Spark writes an Iceberg table through a single physical operator that combines data-file
writing with metadata writing, committing, and catalog validation. Because that operator sits
outside Spark's Adaptive Query Execution (AQE), the sub-query feeding the write — the scans,
projects, sorts, and exchanges producing the rows — cannot be re-planned at runtime.
writing with metadata writing, committing, and catalog validation. Spark's Adaptive Query
Execution (AQE) already re-plans the sub-query feeding that operator — the scans, projects,
sorts, and exchanges producing the rows — but the operator itself sits outside AQE, so the
data-file writing cannot be re-planned in response to how its input ran. And because data-file
writing is bundled with the metadata and commit steps, there is no separate step for Comet to
replace.

When `spark.comet.write.iceberg.splitOperator.enabled=true`, Comet rewrites eligible Iceberg
writes into two operators:
Expand All @@ -39,7 +42,7 @@ writes into two operators:
Iceberg commit (including commit-time validation), outside AQE, exactly once.

With only the split plan enabled, data files are still written by iceberg-java; only the plan
shape changes. The split makes the write's input visible to AQE and to Comet's columnar rules,
shape changes. The split moves data-file writing inside AQE and separates it from the commit,
and it is the foundation for the second toggle: when
`spark.comet.iceberg.write.enabled=true` and the write passes the eligibility check below, the
`IcebergWrite` operator's per-task Parquet write is delegated to
Expand Down