Skip to content

feat: enable native Iceberg writes by default - #6664

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

andygrove wants to merge 15 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, and so does a merge-on-read write: #6693 can plan one through the split plan on Spark 3.5+, but iceberg-java's DeltaWriter still writes its rows, so only the testing split flag does that.

Blockers

Blocker Fix
A fanout write whose open partitions outgrow the memory pool fails its task, where iceberg-java's on-heap writer succeeds (#6771) #6773

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, and plans a merge-on-read WriteDelta (Spark 3.5+) through the split plan only under the testing split flag.
  • 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, that a merge-on-read DELETE plans Spark's WriteDeltaExec, and that with transition reversion on the same write plans CometIcebergWriteExec under the default threshold and, with maxTransitions=0, 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, merge-on-read, eligibility, the cost of reading each natively written file's footer back, two GET requests per file on S3 or GCS, 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 FAQ, 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, 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?

  • 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, the merge-on-read test without the testing-flag check, and the transition-reversion test, which fails with an EOFException at commit without fix: restore Spark write execs when reverting transition-heavy stages #5957, and on its first run when the write is not eligible.
  • Every Iceberg suite that does not need MinIO, plus CometConfSuite, CometScanSchemeFallbackSuite, CometMergeRowsSuite, CometPluginsSuite, CometPlanEqualitySuite, RevertNativeForTransitionHeavyStagesSuite and 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 with INSERT, which goes through the split plan and, for a native input, the native writer.
  • CometIcebergWriteBenchmark at 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-write DELETE). The full table is in feat: enable native Iceberg writes by default #6664 (comment).
  • The Iceberg Spark test diffs set both flags, so those jobs run the same plan as before, merge-on-read included. run-iceberg-tests and run-all-spark-profiles are applied, so CI also runs the Iceberg Spark tests and the Comet suites against every Spark profile.

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.

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
@andygrove andygrove closed this Oct 7, 2026
@andygrove andygrove reopened this Oct 7, 2026
@andygrove andygrove added this to the 1.2.0 milestone Oct 7, 2026
@andygrove
andygrove requested a review from viirya October 7, 2026 21:30

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

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.

Comment on lines +102 to +106
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

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.

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.

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.

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.

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.

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.

Comment thread spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala Outdated

@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 writing required two experimental flags, including a combination where enabling the native writer alone did nothing.
  • Design approach: Make spark.comet.write.iceberg.enabled the 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 IcebergWriteExec and 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 withSessionConf helper appropriately preserves nested configuration values and prevents JVM baseline writes from accidentally becoming native.
  • Behavioral changes worth calling out: Compared with branch-1.1 at e9efd9f764ee0a59b7898ff028d6985d4a7a28e1, 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.
@andygrove

Copy link
Copy Markdown
Member Author

@viirya CometIcebergWriteBenchmark for the four arms, at 1cb149f on an Apple M3 Ultra (local filesystem, 4M rows, local[5], average of five runs):

Case Spark Comet scan, iceberg-java writer Comet scan, native writer Native vs iceberg-java writer
Unpartitioned INSERT 1426 ms 1250 ms 606 ms 2.1x
Partitioned INSERT, clustered writer 1948 ms 1496 ms 1040 ms 1.4x
Partitioned INSERT, fanout writer 1922 ms 1787 ms 951 ms 1.9x
Copy-on-write DELETE 2293 ms 2116 ms 1559 ms 1.4x

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 S3FileIO that is two ranged GETs per file, the 8-byte tail and then the footer, made one file after another. iceberg-java takes the same metrics from the footer it still holds in memory. A task that writes a few large files hardly notices, but a fanout task that writes many small files pays two request round trips for each of them. The user guide, under Native Parquet write eligibility, and the 1.2.0 upgrade entry now say so, and #6772 tracks removing the reads by handing each footer from the native writer to the JVM.

@viirya

viirya commented Oct 8, 2026

Copy link
Copy Markdown
Member

Thanks for the benchmark numbers and the footer-read docs. Both answer what I asked. 94915161e is a good catch: with #6693 on main, the new default would have sent every merge-on-read write through the split plan for no benefit. Since run-all-spark-profiles was removed on 10/07, the strategy change and the new merge-on-read test have only run on Spark 4.1. Could you add the label back for one more round before merge? #6740 is in now, so merging main should also clear the Spark 4.2 expressions failure.

The FAQ added on main called Comet's native Iceberg writes experimental. They are now
on by default for eligible writes.
@andygrove andygrove added 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 8, 2026
@andygrove

Copy link
Copy Markdown
Member Author

@viirya I merged main in (b8cadba), so #6740 is in and the Spark 4.2 expressions failure should clear. The only conflict was gluten_comparison.md, where main had just rewritten the Iceberg paragraph and still called the native writer experimental and off by default. I rewrote that sentence for the new default and fixed the same claim in the FAQ.

I added run-all-spark-profiles back, so the Comet suites run on Spark 3.4, 3.5, 4.0 and 4.2 for this head. That covers the strategy change and, on 3.5 and later, the merge-on-read test. The suites listed in the description also still pass locally on the merged head on 4.1.

@andygrove
andygrove requested review from comphead and viirya October 8, 2026 23:24
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:Iceberg enhancement New feature or request run-all-spark-profiles Run the Comet test suites against every Spark profile on this pull request, ahead of the merge queue run-iceberg-tests

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants