Skip to content

feat: enable native Iceberg writes by default - #6664

Open
andygrove wants to merge 9 commits into
apache:mainfrom
andygrove:iceberg-issue-5644
Open

andygrove wants to merge 9 commits into
apache:mainfrom
andygrove:iceberg-issue-5644

Conversation

@andygrove

@andygrove andygrove commented Oct 5, 2026 •

Copy link
Copy Markdown
Member

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:

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.

What changes are included in this PR?

  • spark.comet.write.iceberg.enabled defaults to true on 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; false plans Spark's own V2 write operator. The planner strategy also plans Spark's operator when spark.comet.enabled or spark.comet.exec.enabled is false, and in plan-only mode.
  • spark.comet.write.iceberg.splitOperator.enabled stays a testing setting, off by default. It plans the split operator with the native writer off, which CometIcebergWriteActionSuite uses so that the tables it writes outside withNativeEnabled are 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.
  • Tests that unset both flags check that an eligible write plans IcebergCommit over CometIcebergWriteExec, and that with transition reversion on and maxTransitions=0 the reverted stage keeps a JVM IcebergWriteExec and commits the rows, with and without AQE. The disabled-config test turns the write path off the way an application does, with spark.comet.write.iceberg.enabled=false alone, and the spark.comet.enabled=false test now also runs for spark.comet.exec.enabled=false. Both check that Spark's own operator is planned.
  • Suites and the write benchmark that set both flags now set only spark.comet.write.iceberg.enabled, except for a detection test that needs the split plan with the native writer off.
  • Docs: the Iceberg writes user guide (intro, configuration example, the split plan's description, eligibility, the fallback list, which now names spark.comet.enabled=false, spark.comet.exec.enabled=false and plan-only mode, and the accepted-divergences heading), iceberg.md, datasources.md, the operators table, understanding-comet-plans.md, the Gluten comparison, the roadmap, a 1.2.0 upgrade-guide entry naming the setting, its 1.1.0 name spark.comet.iceberg.write.enabled, which Comet now ignores, 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?

  • Each new or changed test fails without the change it covers: the default test with spark.comet.write.iceberg.enabled defaulting to false, the disabled-config test with the split flag defaulting to true, the spark.comet.exec.enabled=false test without the strategy's check of that setting, and the transition-reversion test, which fails with an EOFException at commit without fix: restore Spark write execs when reverting transition-heavy stages #5957.
  • Every Iceberg suite that does not need MinIO, plus CometConfSuite, CometMergeRowsSuite, CometPluginsSuite, CometPlanEqualitySuite and the Iceberg SQL file tests, ran locally on Spark 4.1 with the defaults on: 417 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 with INSERT, which goes through the split plan and, for a native input, the native writer. The full Iceberg write action and detection suites and RevertNativeForTransitionHeavyStagesSuite pass on Spark 4.1 at the current code, 219 tests.
  • The Iceberg Spark test diffs set both flags, so those jobs run the same plan as before. run-iceberg-tests and run-all-spark-profiles are applied for the other Spark versions.

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.
@andygrove andygrove added run-iceberg-tests run-all-spark-profiles Run the Comet test suites against every Spark profile on this pull request, ahead of the merge queue labels Oct 5, 2026
@github-actions github-actions Bot added enhancement New feature or request area:Iceberg labels Oct 5, 2026
@andygrove
andygrove marked this pull request as draft October 5, 2026 14:34
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.
@andygrove andygrove changed the title feat: enable the Iceberg split-operator write plan by default feat: enable the Iceberg split-operator plan and native writes by default Oct 5, 2026
….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.
@andygrove
andygrove marked this pull request as ready for review October 6, 2026 18:06
@andygrove

Copy link
Copy Markdown
Member Author

cc @jordepic

@comphead comphead left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

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?

@jordepic

jordepic commented Oct 6, 2026

Copy link
Copy Markdown
Contributor

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.
@andygrove

Copy link
Copy Markdown
Member Author

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 spark.comet.write.iceberg.enabled the only setting. On, which is the default, Comet plans IcebergCommit over IcebergWrite and writes eligible data files natively, with iceberg-java writing the rest inside the same plan. Off, Spark plans its own write operator, as in 1.1.0. spark.comet.write.iceberg.splitOperator.enabled is back to being a testing-only setting, which the suites use to get the split plan with the native writer off.

@sunchao sunchao left a comment

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.

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.enabled becomes 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)

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.

[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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

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.

@comphead comphead left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Thanks @andygrove. Collapsing this into one setting answers my question about the permutations, and the upgrade entry is clear about the opt-out and the memory impact. I read the PR the way someone upgrading from 1.1.0 would and found a few places where the docs and the behavior don't line up.

Docs that still describe the old behavior

These pages are outside the diff, so I can't comment on them inline. Each one still says the native writer is experimental, opt-in or two settings, or that Iceberg writes go through Spark.

  • datasources.md line 32 says Comet "has an experimental, opt-in native Iceberg writer". This is the data sources overview, so it is probably the first place a user reads about Iceberg writes.
  • gluten_comparison.md line 116 says "Iceberg writes still go through Spark."
  • roadmap.md lines 140 to 147 say the feature is "experimental, disabled by default, and controlled by two settings". They also list "enabling both settings by default" as remaining work.
  • understanding-comet-plans.md still labels CometIcebergWrite (line 332), IcebergWrite (line 346) and IcebergCommit (line 347) as experimental. This PR already edits line 225 of the same file, and every Iceberg write plan will now show these operators.

Could these be updated in this PR? I'd rather not leave them for a follow-up, since the user guide and the roadmap are where people check the current state.

The other three points are inline comments.

// Planner strategies run whether or not Comet is enabled, so check it here too: with Comet
// off, Spark must plan its own V2 write operator.
if (!isCometLoaded(conf) || !CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.get(conf)) {
if (!isCometLoaded(conf) || !splitEnabled) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Should spark.comet.exec.enabled=false also keep Spark's own write operator? This strategy is injected unconditionally in CometSparkSessionExtensions, and it only checks isCometLoaded, the two write flags and plan-only mode. isCometLoaded covers spark.comet.enabled, off-heap memory, the shuffle manager and a few other startup conditions, but not spark.comet.exec.enabled. So with the new default, someone who runs Comet for scans or shuffle only would get IcebergCommit over a JVM IcebergWrite after upgrading. The operators page tells them spark.comet.exec.enabled=false turns off native execution, and the new fallback list in iceberg-writes.md names only spark.comet.enabled=false. CometIcebergWriteActionSuite covers that case (line 122) but not this one.

I read this from the code and haven't run it. The description says the split plan alone gives users nothing over Spark's operator, so I'd lean toward adding COMET_EXEC_ENABLED to this check, with a test next to the spark.comet.enabled=false one. If you'd rather keep the split plan in that mode, could the fallback list say so? It could also mention plan-only mode, which this method checks too.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Agreed, the split plan with iceberg-java's writer gives a scan- or shuffle-only deployment nothing, so it shouldn't change their plans. The strategy now also checks spark.comet.exec.enabled, and the spark.comet.enabled=false test runs for both settings. The fallback list in iceberg-writes.md names spark.comet.exec.enabled=false and plan-only mode now too.

# Native Parquet writer (experimental, off by default; requires the split plan)
# Split-operator plan and native Parquet writer (on by default since Comet 1.2.0; false plans
# Spark's own write operator)
spark.comet.write.iceberg.enabled=true

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

This block still reads as an enable step. Line 89 introduces it as "plus the write-side toggle", and this line sets spark.comet.write.iceberg.enabled=true, which is now the default. Copying the block changes nothing for this setting. Of the write-related lines, only spark.comet.exec.localTableScan.enabled=true changes behavior.

Could the example show spark.comet.write.iceberg.enabled=false as the opt-out instead? Or the line could go, with the default stated in the text above.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Good point. I dropped the line from the block, and the text above it now says the setting is on by default and that false plans Spark's own operator. The local table scan setting is the only write-related line left in the example.


### Iceberg Writes

`spark.comet.write.iceberg.enabled` now defaults to `true`. Comet plans an Iceberg `INSERT INTO`,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

The 1.1.0 name for this setting was spark.comet.iceberg.write.enabled, and the 1.1.0 user guide told people to set that name. #6494 renamed it to spark.comet.write.iceberg.enabled after the 1.1 branch was cut. Comet ignores the old name now, and nothing in this entry mentions it.

That matters here because the entry says the setting "now defaults to true" and offers spark.comet.write.iceberg.enabled=false as the way back to 1.1.0 behavior. A deployment that set spark.comet.iceberg.write.enabled=false in 1.1.0 would get native writes after upgrading. This PR is also the one that moves the key out of testing, so it seems like the right place to say it.

Could the entry name the old key and say Comet ignores it, the way the 1.1.0 section does for the testing keys it removed?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Thanks, I'd missed that the rename landed after the 1.1 branch was cut. The entry now says 1.1.0 named the setting spark.comet.iceberg.write.enabled, that it was a testing setting Comet now ignores, and that a deployment that set it to false needs spark.comet.write.iceberg.enabled=false to keep iceberg-java's writer.

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).
@andygrove

Copy link
Copy Markdown
Member Author

Thanks @comphead, these were all still describing the old behavior. datasources.md, gluten_comparison.md, the roadmap's Iceberg section and the operator tables in understanding-comet-plans.md now describe the single, default-on setting, and none of them call the writer or the split-plan operators experimental any more. I've replied on each inline thread with what changed there.

@andygrove

Copy link
Copy Markdown
Member Author

@sunchao #5957 is now merged but we're blocked waiting for review on #6740

@andygrove andygrove changed the title feat: enable the Iceberg split-operator plan and native writes by default feat: enable native Iceberg writes by default Oct 7, 2026
@andygrove andygrove removed the run-all-spark-profiles Run the Comet test suites against every Spark profile on this pull request, ahead of the merge queue label Oct 7, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants