Repository navigation
Conversation
spark.comet.write.iceberg.splitOperator.enabled now defaults to true, so Comet plans an Iceberg append, overwrite, or copy-on-write DELETE, UPDATE or MERGE as IcebergWrite under IcebergCommit instead of Spark's single V2 write operator. iceberg-java still writes and commits the data files, and the native writer (spark.comet.write.iceberg.enabled) stays off. The setting moves from the testing config category to query execution, since it is now the switch that restores Spark's own operator. The user guide, the upgrade guide, the operators table, the Iceberg contributor guide and the Iceberg write review skill describe the new default. This is step 1 of apache#5644.
spark.comet.write.iceberg.enabled now defaults to true, so an eligible Iceberg write is written by iceberg-rust and every other write still falls back to iceberg-java. The setting moves to the query execution category, and the docs and the 1.2.0 upgrade-guide entry describe the new default. CometIcebergWriteActionSuite pins the flag off, so the tables it writes outside withNativeEnabled stay an iceberg-java baseline, and withNativeEnabled now restores the previous value instead of unsetting it, which would turn the native writer back on.
….0 upgrade guide With spark.comet.write.iceberg.enabled on by default, every eligible Iceberg write draws its buffers from Comet's off-heap memory pool (apache#6247) instead of the JVM heap. The 1.2.0 upgrade-guide entry now says so, names the error a task fails with when a fanout write to many partitions outgrows the pool, and lists what lets such a write fit.
|
cc @jordepic |
comphead
left a comment
There was a problem hiding this comment.
Thanks @andygrove it looks fine for me, would be that useful to document how Iceberg writes work for true/false permutations that set to split config and native writes config?
|
Exciting! |
…witch spark.comet.write.iceberg.enabled now plans the split operator by itself, so one setting turns Comet's Iceberg write path on or off: on plans IcebergCommit over IcebergWrite and writes eligible data files natively, off plans Spark's own V2 write operator. Separate flags for the two layers gave four combinations but only three behaviours: the native writer flag did nothing without the split flag, and the split plan with the native writer off gives a user nothing over Spark's own operator. spark.comet.write.iceberg.splitOperator.enabled goes back to the testing category, off by default. It plans the split operator with the native writer off, which CometIcebergWriteActionSuite uses for its iceberg-java baselines. The Iceberg Spark test diffs set both flags to true, so they run the same plan as before.
|
Thanks @comphead. Writing out the four combinations showed we didn't need two settings: the native flag did nothing without the split flag, and the split plan without the native writer gives users nothing over Spark's own operator. So instead of documenting the matrix, I've made |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Iceberg writes used Spark’s combined operator by default. Native writes required two experimental flags.
- Design approach:
spark.comet.write.iceberg.enabledbecomes the supported, default-on switch for both split planning and eligible native writes. - Correctness / compatibility analysis: Reviewed the entire base-relative diff and surrounding code against Spark 3.4–4.2 and the pinned iceberg-java implementations. The rollout materially exposes the unresolved writer-removal defect tracked in #5719 to existing transition-reversion users.
- Key design decisions: Reusing the eligibility gate and JVM committer keeps the configuration simpler without introducing another execution abstraction. Ineligible writes retain iceberg-java’s writer.
- Implementation sketch: Change the default and planner condition, retain the testing-only split flag, restore test configuration correctly, and update tests, benchmark configuration and documentation.
- Behavioral changes worth calling out: Compared with
branch-1.1, eligible writes now use native file production and Comet’s memory pool by default. Native writes also perform a JVM footer read per output file. The upgrade guide documents the rollout and opt-out. No performance benchmark was run locally. - Suggested improvements: Preserve a valid write operator during transition reversion before enabling this default, and cover the configuration below with a regression test. The existing fix work in #5957 addresses the underlying #5719 concern.
Reviewed full SHA: 1dc0d67df2ee8b15dc5b8704598dfb71e2d22831, against 6940c2e584b96ae9e6666228ab14d5d8bcc693f6. Routed skills: review-comet-pr and review-comet-iceberg-write-pr. Existing PR discussion contained no duplicate finding.
Exact-head CI: 72 successful checks, 11 skipped and 2 failed. All four Iceberg matrix families passed. The Spark 4.2 expressions job failed two trunc_timestamp.sql cases with long overflow during Spark reference evaluation, causing Required Checks to fail. Upstream Spark SQL suites were skipped.
Validation: A bounded Spark 4.1.3 JVM reproduction confirmed writer removal through the full post-columnar rule pipeline. The stock Spark control committed its input successfully. Changed production classes and relevant rule/writer classes were freshly compiled against cached unchanged dependencies. No full native rebuild, end-to-end Iceberg storage test or object-store test was run locally. The checkout remains unchanged.
| "commits every write. Set this to false to plan Spark's own V2 write operator.") | ||
| .booleanConf | ||
| .createWithDefault(false) | ||
| .createWithDefault(true) |
There was a problem hiding this comment.
[P2] Preserve the write operator during transition reversion before enabling this default. With spark.sql.adaptive.enabled=false, spark.comet.exec.transitionRevert.enabled=true and spark.comet.exec.transitionRevert.maxTransitions=0, an eligible Iceberg INSERT ... SELECT now selects CometIcebergWriteExec without either write flag being set. Spark inserts a ColumnarToRowExec beneath it, which triggers reversion before redundant transitions are removed. Because the native writer’s originalPlan is its child, reversion removes the writer entirely. IcebergCommitExec then attempts to deserialize table data as commit messages and the write fails. Previously, these settings retained Spark’s working V2 writer. The underlying defect is tracked in #5719, but this default materially expands its exposure to users who never enabled experimental Iceberg writes. Restore an IcebergWriteExec when reverting, or safely exclude this combination from native conversion, and add a regression covering the default flags.
Evidence: Bounded reproduction: /tmp/6664-current-review-075caeqj/WriteRevertProbe.scala, with results in probe-full-pipeline.log. Using exact-head classes, Spark’s transition insertion and CometRule.postColumnarRules transformed IcebergCommit -> CometIcebergWrite -> ColumnarToRow -> CometProject -> CometLocalTableScan into IcebergCommit -> Project -> Project -> LocalTableScan. Output: COUNT=1, WRITERS_LEFT=0, followed by java.io.EOFException for one ordinary binary row. The equivalent stock AppendDataExec control reported COMMIT_ROWS=1 and STOCK_SUCCEEDED=true. Disabling reversion retained the native writer. This was a JVM plan/commit harness with a test BatchWrite, without native storage I/O.
There was a problem hiding this comment.
This is #5719, and #5957 fixed it by making the reverted stage restore an IcebergWriteExec instead of dropping the write. This branch now includes #5957. I added transition-heavy fallback keeps the write when no flag is set, which writes from a Parquet source with transition reversion on, maxTransitions=0 and no write flag set, under both AQE settings. Without #5957 it fails with the same EOFException at commit, and with it the write goes through iceberg-java's writer and commits the expected rows.
With the new default, an application that turns off spark.comet.exec.enabled and uses Comet only for scans or shuffle got IcebergCommit over a JVM IcebergWrite. The strategy now plans Spark's own operator in that case too. Also brings the remaining docs in line with the single, default-on setting, and notes the 1.1.0 name of the setting in the migration guide.
…th default flags With no write flag set, an eligible write becomes a CometIcebergWriteExec, so transition reversion has to restore a JVM writer (apache#5719, fixed by apache#5957).
|
Thanks @comphead, these were all still describing the old behavior. |
viirya
left a comment
There was a problem hiding this comment.
Thanks for pulling all the prerequisites together, and for folding the two flags into one. Writing out the four combinations made a good case that only three behaviours were ever reachable, and IcebergWriteStrategy reads cleanly with the single switch. I checked the earlier threads at 5664442ef, and they all look addressed. #5957 is in, the spark.comet.exec.enabled check and its test are there, and the docs comphead listed now describe the new default.
I have two questions about the default itself, and one about a test, inline.
Could you attach CometIcebergWriteBenchmark results for the four arms (unpartitioned, clustered, fanout, copy-on-write delete) to this PR? #5644 lists them under "Before 1.2.0 ships", but this is the PR that turns the native writer on for every eligible write. It seems like the right place to show there is no regression against iceberg-java. The benchmark writes to the local filesystem, so it can't show the footer read the JVM does for each written file to rebuild metrics. On S3 that is an extra GET per file. Could the PR or the user guide say what that costs on an object store, especially for writes that produce many small files?
Once #6740 lands, the Required Checks failure from the Spark 4.2 expressions job should clear.
| The native writer's buffers count against Comet's off-heap memory pool, 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, and | ||
| each open file holds the row group it is writing, so a task that writes to many partitions can need | ||
| more memory than the pool grants it. The task then fails with | ||
| `Additional allocation failed for IcebergWriteExec`. Disabling the fanout writer |
There was a problem hiding this comment.
This paragraph is clear about the failure mode, but it leaves the user to discover it after a job fails. On Iceberg 1.5+, a partitioned table with write.distribution-mode=none and no sort order uses the fanout writer by default. That is a fairly common setup for people avoiding the shuffle. Under this PR those writes go native, their buffers come out of the off-heap pool, and the writer can't spill. A write that fit in the executor heap on 1.1.0 can now fail with Additional allocation failed for IcebergWriteExec. Spark's retry takes the same native path, so the job fails.
Would it be reasonable to keep fanout writes on iceberg-java by default for now? Another option is for the native writer to close the largest open file when the pool refuses a reservation, then retry. Roll points are already an accepted divergence, so that wouldn't add a new kind of difference. If neither fits, could you share some numbers on how many partitions a fanout task can hold with a typical spark.memory.offHeap.size? Then we can judge whether the default is safe for that shape.
There was a problem hiding this comment.
I'd rather not keep fanout writes on iceberg-java. On Iceberg 1.5+ an unsorted partitioned table gets the fanout writer by default whatever its distribution mode, since useFanoutWriter is true whenever the write has no required ordering, so that would take most partitioned writes off the native path.
#6773 takes your second suggestion instead. When the pool refuses the writer's reservation, it writes out and closes the partitions holding the most memory, counting both a partition's open file and the rows it has not handed to the file yet, until what is left fits. A closed partition's next rows open a new file, so the write finishes with more, smaller files rather than failing, and a files closed early to free memory metric on CometIcebergWrite counts them. More files is the same kind of difference as the roll points already listed under accepted divergences, and the PR adds it there.
The suite's 64-partition fanout test now interleaves its partitions and passes under the same 4 MiB pool, writing 219 files where the full pool writes 64. A write that keeps one file open, unpartitioned or clustered, can still fail, but only when that one file's row group outgrows the pool. I've added #6773 to the blockers in the description.
There was a problem hiding this comment.
You're right about useFanoutWriter, and I had the scope too narrow. Since apache/iceberg#8621, writeOrdering drops the local sort whenever the table is unsorted and fanout is on by default, so any unsorted partitioned table gets the fanout writer whatever its distribution mode. Keeping fanout on iceberg-java would take most partitioned writes off the native path, so closing partitions early is the better fix. I'll review #6773 separately, and I'd like to keep this thread open until it lands, since it's listed as a blocker.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Iceberg writes used Spark’s combined operator by default. Native writing required two experimental flags, including a combination where enabling the native writer alone did nothing.
- Design approach: Make
spark.comet.write.iceberg.enabledthe supported, default-on switch for split planning and eligible native writes. Keep iceberg-java responsible for commits and declined writes. - Correctness: Checked writer selection, commit memoization, AQE replanning, fallback and configuration restoration. The earlier writer-removal concern is addressed by #5957. A freshly compiled Spark 4.1.3 JVM probe restored exactly one
IcebergWriteExecand committed successfully. The existing fanout-memory concern remains substantiated below. - Compatibility analysis: Compared Spark’s V2 write, commit and AQE contracts across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0, plus writer selection in the pinned Iceberg versions. The existing shims preserve version-specific MERGE summaries and Spark 4.2 transactions. Disabling Comet, native execution or the supported write flag retains Spark’s operator.
- Key design decisions: Retain the eligibility allowlist, native-child requirement, JVM metadata reconstruction and JVM committer. Preserve the split-only testing flag so parity tests still have an iceberg-java baseline.
- Implementation sketch: Change the configuration default and category, broaden the planner condition, add the native-execution guard, restore previous test settings, and update plan assertions, benchmark configuration and documentation.
- Performance: Native writing avoids row-based file production but adds a JVM footer read per output file and charges writer buffers to Comet’s pool. No throughput or object-store benchmark was run locally, so this review makes no net-speed claim. Fanout exhaustion is demonstrated, not hypothetical.
- Design: One supported switch makes the public behavior easier to understand. Reusing the existing planner and eligibility gate keeps the rollout small. The unresolved memory behavior limits the safety of making fanout writes native by default.
- Abstraction & complexity: No new execution abstraction is introduced. The scoped
withSessionConfhelper appropriately preserves nested configuration values and prevents JVM baseline writes from accidentally becoming native. - Behavioral changes worth calling out: Compared with
branch-1.1ate9efd9f764ee0a59b7898ff028d6985d4a7a28e1, eligible writes now use native file production by default, with different physical file layout and memory budgeting. The upgrade guide explains the opt-out and the ignored old key,spark.comet.iceberg.write.enabled. - Suggested improvements: Address the existing fanout-memory concern by retaining JVM fanout writing by default or providing bounded flushing when reservations fail. No additional P1/P2-level issue or recommendation was identified.
Reviewed full SHA: 5664442efb6fc043519503256f085ed008d2b0ae, against supplied base c8e6553dad376537c144e36be9dab7e31df2df0e. Reviewed the complete 20-file endpoint diff and existing non-Copilot discussion. The two datetime-file differences come from #6348 landing on the base after the branch diverged at 46d91ef02d80f2a42cd89034699550f45fa13394. They are unchanged from that common ancestor, not a PR-introduced reversion.
Routed skills: review-comet-pr, review-comet-iceberg-write-pr, review-comet-expression-pr and review-comet-memory-pr.
Existing blocker: The fanout-memory thread remains unresolved. Exact-head CI passes the test that intentionally confirms native failure for 64,000 rows across 64 partitions, with 500-character payloads and approximately 4 MiB of configured native budget. run_write_task propagates a refused reservation without flushing or switching writers. A local Spark 4.1.3/iceberg-java 1.11.0 control wrote the same data successfully into 64 files. The default switch exposes previously working JVM writes to this failure. This is existing feedback, so no duplicate inline finding is returned.
Exact-head CI: Latest workflow 37598531752 has 54 successful and 11 skipped jobs, including successful Required Checks and all four Iceberg matrix families. Its Spark 4.1 scans job passed 681 tests, including the new default, disable-switch and AQE-on/off reversion cases. The latest result for Spark 4.2 expressions remains failed from the earlier same-head run: two trunc_timestamp.sql cases raise long overflow during Spark reference evaluation. That profile was not rerun in the newest workflow. The earlier Iceberg 1.11 shard failure reported lost runner communication and subsequently passed. Upstream Spark SQL jobs were skipped.
Validation limits: The local transition probe freshly compiled the changed configuration/planner and relevant writer/rule classes against cached dependencies, using a test BatchWrite. The local storage control used iceberg-java. No full native rebuild, local end-to-end native storage test, object-store test or performance benchmark was run. Native execution evidence comes from exact-head CI. Project files remain unchanged.
No additional introduced P1/P2 issues found within this review beyond the substantiated existing concern above.
# Conflicts: # docs/source/user-guide/latest/iceberg-writes.md
…is reverted Run the same write under the default threshold first and require a CometIcebergWriteExec, so the reverted run cannot pass because the write was never eligible.
…stores The JVM reads each natively written file's footer back to compute its Iceberg metrics. Say that this is two reads per file, made one after another, which on S3 or GCS are GET requests, so a write of many small files pays two round trips for each.
…e testing split flag is on apache#6693 plans Spark 3.5+ WriteDelta through the split plan, with Iceberg's JVM DeltaWriter still writing the rows. With spark.comet.write.iceberg.enabled on by default, every merge-on-read write would take that path although the plan gives it nothing over Spark's own operator. Only the testing split flag plans it now, so users keep Spark's operator until a native delta writer exists (apache#6240), while the Iceberg Spark tests, which set both flags, still cover it.
|
@viirya
The native writer is faster than iceberg-java's in all four. On the footer read: for every file it writes natively, the task reads the file's footer back before it returns its commit message, so that iceberg-java's own code computes the file's metrics. Through |
|
Thanks for the benchmark numbers and the footer-read docs. Both answer what I asked. |
The FAQ added on main called Comet's native Iceberg writes experimental. They are now on by default for eligible writes.
|
@viirya I merged I added |
Which issue does this PR close?
Part of #5644, which makes Comet's Iceberg write path the default in 1.2.0: the split-operator plan and the native writer.
Rationale for this change
#5644 makes Comet's Iceberg write path the default. This PR does it for 1.2.0, targeted for late October or early November (#6550), so the path gets nightly runs before the release:
IcebergWriteunderIcebergCommit. IcebergWriteStrategy ignores spark.comet.enabled, so the split-operator write plan survives the kill switch #6142 (the planner strategy ignoredspark.comet.enabled) and Preserve catalog transactions for split Iceberg writes on Spark 4.2 #6586 (Spark 4.2 catalog transactions, fixed by fix(iceberg): preserve Spark 4.2 write transactions #6587) are fixed, and the Iceberg Spark test jobs have run with the split plan on since it was added in feat: Optionally split the Iceberg V2 write operator into distinct writer and committer operations #4658.CometNativeException, fixed by fix: rethrow the exception a JVM input throws instead of a CometNativeException #6243), Account the native Iceberg writer's buffers in Comet's memory pool #5648 and Charge native write buffers to the shared off-heap memory pool #6115 (the writer's buffers were not charged to Comet's memory pool, fixed by feat: charge native write buffers to the task's memory pool #6247) and revertToSpark erases CometIcebergWriteExec / CometNativeWriteExec because originalPlan is the node's own child #5719 (reverting a transition-heavy stage droppedCometIcebergWriteExec, so the commit failed, fixed by fix: restore Spark write execs when reverting transition-heavy stages #5957) are fixed, and the Iceberg Spark test jobs have run with the native writer on since test: run the Iceberg Spark tests with the native Iceberg writer enabled #5677.One setting,
spark.comet.write.iceberg.enabled, switches both. Separate flags for the two layers gave four combinations but only three behaviours: the native writer flag did nothing without the split flag, and the split plan with the native writer off gives a user nothing over Spark's own operator. Turning the one setting off plans Spark's own operator, as in Comet 1.1.0. For the same reason, an application that turns off Comet's native execution and uses Comet only for scans or shuffle keeps Spark's operator too, and so does a merge-on-read write: #6693 can plan one through the split plan on Spark 3.5+, but iceberg-java'sDeltaWriterstill writes its rows, so only the testing split flag does that.Blockers
What changes are included in this PR?
spark.comet.write.iceberg.enableddefaults totrueon every Spark version and moves from the testing config category to query execution. It plans the split operator by itself and writes eligible data files natively;falseplans Spark's own V2 write operator. The planner strategy also plans Spark's operator whenspark.comet.enabledorspark.comet.exec.enabledisfalse, and in plan-only mode, and plans a merge-on-readWriteDelta(Spark 3.5+) through the split plan only under the testing split flag.spark.comet.write.iceberg.splitOperator.enabledstays a testing setting, off by default. It plans the split operator with the native writer off, whichCometIcebergWriteActionSuiteuses so that the tables it writes outsidewithNativeEnabledare an iceberg-java baseline for the native writer's output. That suite's session-conf helper restores previous values instead of unsetting them, which would turn the native writer back on.IcebergCommitoverCometIcebergWriteExec, that a merge-on-readDELETEplans Spark'sWriteDeltaExec, and that with transition reversion on the same write plansCometIcebergWriteExecunder the default threshold and, withmaxTransitions=0, keeps a JVMIcebergWriteExecand commits the rows, with and without AQE. The disabled-config test turns the write path off the way an application does, withspark.comet.write.iceberg.enabled=falsealone, and thespark.comet.enabled=falsetest now also runs forspark.comet.exec.enabled=false. Both check that Spark's own operator is planned.spark.comet.write.iceberg.enabled, except for a detection test that needs the split plan with the native writer off.spark.comet.enabled=false,spark.comet.exec.enabled=falseand plan-only mode, and the accepted-divergences heading),iceberg.md,datasources.md, the operators table,understanding-comet-plans.md, the Gluten comparison, the FAQ, the roadmap, a 1.2.0 upgrade-guide entry naming the setting, its 1.1.0 namespark.comet.iceberg.write.enabled, which Comet now ignores, the footer reads, and that the native writer's buffers draw on Comet's off-heap memory pool, the Iceberg writes and Iceberg Spark tests contributor guides, and the Iceberg write review skill.How are these changes tested?
spark.comet.write.iceberg.enableddefaulting tofalse, the disabled-config test with the split flag defaulting totrue, thespark.comet.exec.enabled=falsetest without the strategy's check of that setting, the merge-on-read test without the testing-flag check, and the transition-reversion test, which fails with anEOFExceptionat commit without fix: restore Spark write execs when reverting transition-heavy stages #5957, and on its first run when the write is not eligible.CometConfSuite,CometScanSchemeFallbackSuite,CometMergeRowsSuite,CometPluginsSuite,CometPlanEqualitySuite,RevertNativeForTransitionHeavyStagesSuiteand the Iceberg SQL file tests, ran locally on Spark 4.1 at b8cadba with the defaults on: 497 passed, none failed, and 3 were canceled by version assumptions (the existing SPARK-55626 one and two Spark 4.2-only tests). Most of these suites create their tables withINSERT, which goes through the split plan and, for a native input, the native writer.CometIcebergWriteBenchmarkat 1cb149f, on an Apple M3 Ultra with a local filesystem: the native writer is faster than iceberg-java's in all four cases (2.1x unpartitioned, 1.4x clustered, 1.9x fanout, 1.4x for a copy-on-writeDELETE). The full table is in feat: enable native Iceberg writes by default #6664 (comment).run-iceberg-testsandrun-all-spark-profilesare applied, so CI also runs the Iceberg Spark tests and the Comet suites against every Spark profile.