Apache DataFusion Comet 1.1.0 Release
Posted on: Thu 01 October 2026 by pmc
The Apache DataFusion PMC is pleased to announce version 1.1.0 of the Comet subproject.
Comet is an accelerator for Apache Spark that translates Spark physical plans to DataFusion physical plans for improved performance and efficiency without requiring any code changes.
This release covers roughly seven weeks of development since 1.0.0: 398 commits from 40 contributors. See the change log for the full list of changes.
The highlights:
- Native Iceberg writes (experimental): Comet can write Iceberg data files with iceberg-rust instead of iceberg-java.
- Memory management: Comet now measures the native memory its pools don't track, logs it on every executor, and fixes several long-standing pool accounting bugs.
- Native shuffle over Apache Celeborn: 1.1.0 includes the native writer and reader, but released Celeborn clients currently fall back to ordinary Spark/Celeborn shuffle.
- Native Parquet writes on Spark 4.0+ (experimental), now built on Spark's own write path.
Native Iceberg Writes (Experimental)¶
Until now, Comet accelerated Iceberg reads but left writes to the JVM. In 1.1.0, Comet can write Iceberg data files natively using iceberg-rust. The feature is experimental and disabled by default, and it only engages under fairly strict conditions.
Splitting the write operator¶
Spark writes an Iceberg table with a single operator that writes the data files, writes the metadata, commits, and validates against the catalog. Because file writing is bundled with the metadata and commit steps, there was no separate piece for Comet to replace.
Setting spark.comet.write.iceberg.splitOperator.enabled=true splits eligible Iceberg writes into two operators:
IcebergWritewrites data files on the executors and returns each task's commit message.IcebergCommitcollects the commit messages on the driver and performs the normal Iceberg commit, once.
With only the split enabled, iceberg-java still writes the files; the native writer, described next, replaces
that step. The split covers INSERT INTO and DataFrame append, static and dynamic INSERT OVERWRITE, and
copy-on-write DELETE, UPDATE, and MERGE, on every supported Spark version. Merge-on-read writes are left
alone. When Comet can't split a write (an unrecognized write class, CTAS on Spark 3.4, or a write that needs
Spark's commit coordinator), it plans the write exactly as if Comet weren't there.
Writing Parquet natively¶
The native writer requires the split. Also setting spark.comet.iceberg.write.enabled=true hands each task's
Parquet writing to iceberg-rust. The driver passes the native writer everything it needs: schema, partition spec,
data location, Parquet settings, writer mode, and object-store configuration. Each task writes its files and
returns their metadata as an in-memory Iceberg manifest.
From there, iceberg-java takes over again. The JVM reads the manifest, recomputes each file's metrics from its Parquet footer, and builds the same commit message the JVM writer would have produced. Snapshots, manifest lists, commit validation, and retries all work as before.
The native writer consumes Arrow batches from a Comet operator, so the query feeding the write must run in Comet
too. Writes from a local relation, such as INSERT ... VALUES, also need
spark.comet.exec.localTableScan.enabled=true.
To see which path a write took, check the physical plan: a native write shows CometIcebergWrite under
IcebergCommit, and a write that fell back shows IcebergWrite.
Matching iceberg-java¶
The native writer aims to write the same table iceberg-java would, down to the metadata. Manifest metrics drive pruning for every future reader, so rather than trust the native writer's numbers, Comet recomputes them with iceberg-java's own code. Parity tests write the same rows through both writers and compare the committed counts and bounds. The cost is one small footer read per file.
Eligibility is an allowlist. A write goes native only if its whole configuration matches the documented set of
supported settings. Anything else, including properties added by future Iceberg versions, falls back to
iceberg-java and reports why in Comet's extended EXPLAIN output. Encryption, object-storage layout, custom
location providers, bloom filters, and unrecognized parquet.* properties all fall back. For a partitioned
table, Iceberg asks Spark to cluster and sort rows by partition value, such as bucket(16, id) or days(ts),
before writing. That step must also run natively, which it now can because Iceberg's system functions have
native implementations.
Comet decides all of this at planning time, including checking that every iceberg-java class it relies on is where it expects. An Iceberg release that moves one falls back instead of failing tasks mid-write. CI also runs Apache Iceberg's own Spark test suites, for Iceberg 1.8.1 through 1.11.0, with the native writer enabled.
Failure handling¶
Only successful tasks' commit messages are committed. A failed task deletes the files it wrote, as iceberg-java
does, and a failed job commits nothing and deletes the completed tasks' files. Retries can't collide, because
each attempt's ID is part of its file names. Anything cleanup misses is invisible to readers and is removed by
Iceberg's normal remove_orphan_files maintenance.
Known differences¶
Files written by parquet-rs, the Rust Parquet library iceberg-rust uses, differ from those written by
parquet-java (formerly parquet-mr), which iceberg-java uses. Enabling the feature accepts these differences.
Most are cosmetic, such as footer metadata, created_by, and page encoding labels. Two are worth knowing about:
- Files won't split at the same points. Both writers check file size every 1000 rows, but they estimate it differently, so they can roll over to a new file at different rows.
- High-cardinality columns keep a dictionary page. parquet-java drops dictionary encoding early when it isn't
saving space, while parquet-rs keeps it up to
write.parquet.dict-size-bytes. Results are the same, but selective reads of native-written files fetch more bytes (#6114).
The Iceberg Writes guide lists every supported setting and every difference. Please try it on a non-production table and tell us how it goes. Feedback from real workloads is what this feature needs before it loses the experimental label.
Thanks to @jordepic for designing and implementing the split-operator plan, write detection, and the native writer, and to @andygrove for the fidelity and failure-handling work, with contributions from @zhangfengcdt, @snmvaughan, @liupoyi-1031, and @0lai0, and reviews from @sunchao, @comphead, @unikdahal, and @mbutrovich. Related PRs: #4658, #5298, #5361, #5663, #5780.
Memory Management¶
A recurring problem for Comet users has been executors killed by the cluster manager (on Kubernetes,
ExecutorLostFailure with exit code 137) even though Comet stayed within its memory pool. 1.1.0 explains why
this happens, lets you measure it, and fixes pool bugs that made it worse.
Where the memory goes¶
Comet's native operators allocate from the Rust heap but charge their reservations against Spark's off-heap
pool, sized by spark.memory.offHeap.size. Operators only reserve memory for data they deliberately hold onto:
sort buffers, hash join build sides, aggregation state, and buffered shuffle partitions.
Plenty of memory is never reserved: working memory in expression kernels, decompression buffers, Parquet reader
state, object store buffers, JVM-side Arrow buffers, and allocator overhead. All of it has to fit in
spark.executor.memoryOverhead, and until now there was no way to see how much of it there was.
Measuring used memory¶
Comet now wraps its native allocator in a counter that tracks every byte allocated and not yet freed. It only observes and never rejects an allocation. Arrow buffers the JVM imports from native code are now tracked separately, so they aren't counted twice.
Each executor uses the counter to log its native memory usage every 10 seconds while Comet runs:
Comet native memory usage: allocated 5412.3 MiB, reserved 3890.0 MiB (16 native plans, 8 memory pools)
reserved is what the pools track, and the container already has room for it. allocated is everything
Comet's native code holds. The difference is what has to fit in spark.executor.memoryOverhead, next to the
JVM's own non-heap memory. To size the overhead, run a representative workload, take the largest difference from
the log, add it to the overhead you had before Comet, and leave a margin. Setting
spark.comet.memory.logInterval=1s for that run makes it less likely to miss a short peak.
The executor also warns when its native memory looks larger than the container allows, and the
tuning guide walks through sizing with examples. One gotcha: setting spark.executor.memoryOverhead
replaces the value Spark derives from spark.executor.memoryOverheadFactor instead of adding to it, so on a
large executor it can shrink the container. For large executors, raise the factor instead.
Memory pool fixes¶
Comet's off-heap memory comes from one of two pools. fair_unified, the default, caps each operator at an even
share of the task's memory. greedy_unified gives memory to operators first come, first served. When Spark
grants an operator less memory than it asked for (a partial grant), the operator spills.
fair_unifiedlimited a whole task to one operator's share. Since 0.15.0, the pool compared all of a task's reservations against a single operator's share, and each new operator shrank the limit for the rest. Each operator now gets its own share, so tasks with several operators can use more memory before spilling.- A partial grant from Spark no longer panics the task. DataFusion sometimes has to record memory that already exists, such as a spilled batch being read back, and both pools panicked when Spark granted less than they asked for. They now track the shortfall and repay it before returning memory to Spark.
- Leaks on failure paths. A plan that failed during setup or teardown could leak its task's memory pool, and a failed Arrow import leaked the vectors it had already imported.
- Configuration units. 1.0.0 misread three size settings, including
spark.memory.offHeap.sizewritten as a bare number. See Upgrading to 1.1.0 for what changes. - Less log noise when spilling. 1.0.0 logged a warning and a memory dump for every partial grant, so a spilling query could log hundreds. Partial grants now log at DEBUG, and the dump is gone because it could deadlock the task.
- Metrics. Spills from native sorts, sort-merge joins, and aggregates now count toward Spark's Spill (Disk) task metric in every stage. The native shuffle writer reports Spill (Disk) and Spill (Memory) separately, where 1.0.0 put the same number in both. Native aggregates with grouping keys also show spill counts, bytes, and rows in the SQL tab.
spark.comet.exec.memoryPool.fraction is deprecated. It was meant to leave room for untracked memory, but Spark
hands out the whole pool anyway, so it never did. Size spark.executor.memoryOverhead instead.
On-heap mode¶
On-heap mode exists so that Spark's and Iceberg's test suites can run against Comet; production runs off-heap.
Its memory accounting didn't protect anything, because native memory isn't on the JVM heap, so 1.1.0 removes that
accounting, along with six of the nine pool types and several testing-only settings. Outside tests, Comet now requires
off-heap memory however it's enabled, including when CometSparkSessionExtensions is registered directly, and
disables itself with a warning otherwise.
Contributors can find more detail in the new memory management guide.
Thanks to @andygrove for driving this work, @peterxcli for the memory pool lifecycle and shuffle spill accounting fixes, @ywskycn for reporting native memory usage to Spark, @1fanwang for the Arrow import leak fix, and @sunchao for the native aggregate spill metrics, with reviews from @sunchao, @comphead, and @mbutrovich. Related PRs: #5934, #6162, #6048, #6128, #6205, #6066, #5494.
More Iceberg Improvements¶
- V3 deletion vectors are applied on native scans.
- Iceberg system functions (
bucket,truncate,years,months,days, andhours) run natively. - Scan planning metrics and scan time appear in the Spark UI for native Iceberg scans.
- A wrong-results fix. A filter such as
bucket(4, id) = 2combined with another predicate was pushed to the native scan asid = 2, returning too few rows. - Tables partitioned by an unknown transform can be read natively, and
IS NULL/IS NOT NULLchecks on list and map columns no longer force a fallback.
Thanks to @mbutrovich for deletion vector support, @parthchandra for the scan metrics, @ErikBPF for the null-check fix, and @andygrove for the native system functions and residual fix, with reviews from @sunchao, @rich7420, @unikdahal, and @jordepic. Related PRs: #5853, #5638, #6027, #6154.
Remote Shuffle with Celeborn¶
1.1.0 adds Comet's side of native shuffle over Apache Celeborn: map tasks push Comet's Arrow data straight to Celeborn, and reducers read it back natively. In 1.0.0, Celeborn users always got ordinary Spark shuffle.
Native Celeborn shuffle is unavailable with released Celeborn 0.6.x and 0.7.x clients in Comet 1.1.0 because
these clients do not pass Comet's push-completion compatibility checks. Set spark.shuffle.manager to
org.apache.spark.sql.comet.execution.shuffle.CometCelebornShuffleManager to use Comet alongside
ordinary Spark/Celeborn shuffle; Comet can still accelerate other operators. These clients retain
ordinary Spark/Celeborn shuffle even when spark.comet.shuffle.mode=native is requested, and the
default auto mode also retains ordinary shuffle. Follow-up work to enable native shuffle with
released clients is tracked in #6523.
Native shuffle over Celeborn does not support spark.io.encryption.enabled=true, and Celeborn
is not bundled with Comet. The Celeborn tuning guide has the full requirements.
Thanks to @pingzh for this work, with reviews from @sunchao, @ziting-openai, and @andygrove. Related PRs: #5473, #5481, #5513, #5531, #5537.
Native Parquet Writes on Spark 4.0+ (Experimental)¶
Native Parquet writes remain experimental and disabled by default, but on Spark 4.0+ they now build on Spark's own write path.
In 1.0.0, a native write replaced Spark's entire write command, so Comet had to reimplement the commit protocol,
save modes, and job commit. On Spark 4.0+, Comet replaces only WriteFilesExec, the per-task write step that
Spark 4.0 made pluggable, and Spark handles the rest. That means:
- File names come from Spark's commit protocol, so committers that track individual files, like the S3A magic committer, work.
- Column names, nullability, and field IDs come from the target table, not the query.
- Bytes-written and rows-written metrics are correct on HDFS.
The writer handles non-partitioned, non-bucketed writes to local and HDFS paths, and skips writes that set
spark.sql.files.maxRecordsPerFile. To enable it, set spark.comet.parquet.write.enabled=true and
spark.comet.operator.WriteFilesExec.allowIncompatible=true. Spark 3.4 and 3.5 keep the old write path.
Thanks to @andygrove and @sunchao for this work, with reviews from @comphead, @peterxcli, @rich7420, and @parthchandra. Related PRs: #5763, #5369.
S3 Credentials¶
Comet's native S3 access now works with more of the credential setups Spark supports:
- Built-in credential provider adapters. In 1.0.0, a native Parquet scan failed with
Unsupported credential providerfor classes Spark accepts but Comet didn't reimplement, such asDefaultAWSCredentialsProviderChain. Two new adapters fix this by settingspark.hadoop.fs.s3a.comet.credential.provider.class.HadoopS3ACredentialProviderAdapter(recommended) uses Hadoop S3A's own provider chain, so it supports everything S3A does, including web identity, assumed roles, and per-bucket settings.AwsSdkCredentialProviderAdapterwraps a specific AWS SDK provider class. - STS throttling protection on EKS with IRSA. When many executors start at once, STS can throttle their
credential requests. In 1.0.0, the native Iceberg path then fell back to the EKS node role, which usually can't
read the bucket, and the job failed with
403errors. Native Iceberg reads and writes now fetch web-identity credentials themselves, retry throttled requests with backoff, never fall back to the node role, and share one credential per executor. This happens automatically. - Per-location credentials. A provider implementing
CometS3LocationScopedCredentialProvidercan return different credentials for different prefixes in the same bucket, such aswarehouse/salesandwarehouse/finance. - Custom S3-compatible schemes. Vendor filesystems that front an S3-compatible store with their own URL
scheme, such as
blob://, can be read natively by listing the scheme inspark.hadoop.fs.comet.s3Compliant.schemes.
See the S3 credential providers guide for details.
Thanks to @parthchandra for the credential adapters and STS throttling protection, @snmvaughan for per-location credentials, and @comphead for custom S3-compatible schemes, with reviews from @sunchao and @andygrove. Related PRs: #6023, #6025, #6031, #5314.
Performance¶
Shuffle¶
- In 1.0.0, the native shuffle writer ran with a 1-byte write buffer by default: its 1 MiB default was declared in MiB but passed to native code as bytes. It now gets the intended 1 MiB.
- Round-robin repartitioning assigns rows by position, as Spark does, instead of hashing every column of every row. That's much cheaper on wide nested schemas, and it spreads repeated values across reducers the way Spark does.
- A task spills all of its shuffle partitions into one file instead of one file each.
- Zstd and Arrow IPC compression contexts are reused across shuffle blocks.
- Shuffle blocks are decoded against a cached schema instead of re-parsing it for every block.
- Per-partition scratch buffers are reused, and
ArrowWriterbulk-copies fixed-width columns.
Thanks to the contributors who drove this work, especially @peterxcli, @dwsmith1983, @pingzh, and @andygrove, with reviews from @sunchao and @mbutrovich.
Planning and Execution¶
- Native dynamic filter pushdown: a hash join's build-side keys filter the probe-side Parquet scan at runtime, so the scan can skip data that can't match.
- Adaptive partial aggregation for eligible native shuffle plans: when grouping keys are nearly unique, the partial aggregate passes rows through instead of building a hash table that doesn't reduce them.
- Parsed plan data is cached across a stage's tasks.
- New native Parquet scan I/O and read-amplification metrics.
Expressions¶
Many expression kernels got faster. In kernel microbenchmarks, map_sort is up to 3x faster on multi-entry string
maps and 18x faster on single-entry maps. The map lookup behind element_at and GetMapValue is vectorized,
collect_list and collect_set have a native GroupsAccumulator, and regex patterns are compiled once per
planned expression. Smaller gains cover date and time field extraction, posexplode, unnesting, nested hashing,
decimal overflow checks, list_extract, and approximate percentile merges.
Expanded Coverage¶
New native support includes the mode, max_by / min_by, listagg / string_agg (Spark 4.0+), and
regr_* regression aggregates; WindowGroupLimitExec; explode_outer; make_interval; native
sequence for integral types and unbase64; _metadata constant columns in the native Parquet
scan; nested types as shuffle hash partitioning keys; BinaryType in sort-merge joins; and Spark 4's
EmptyRelationExec.
More expressions now use codegen dispatch by default, where Comet runs Spark's own generated code for an
expression inside the native pipeline, so the result matches Spark without falling back. They include
translate, to_csv, encode, lpad / rpad, round on floats, abs on intervals, timestampadd /
timestampdiff, next_day and levenshtein on collated input, and Spark expressions that call a Java or
Scala method directly. rlike runs natively by
default for literal patterns that behave the same as in Java.
Spark 4 Variant support also improved: native Parquet scans can project Variant columns directly. Thanks to @peterxcli for driving Variant support, with reviews from @sunchao. Related PRs: #5868, #5794.
There's also experimental native support for the in-memory cache (spark.comet.exec.inMemoryCache.enabled,
off by default).
Previewing Comet Plans¶
The new spark.comet.explain.planOnly.enabled setting logs the plan Comet would have run to the driver log,
then lets Spark run the query as normal. It shows how much of a production workload Comet would accelerate, and
why anything falls back, without running any of it through Comet.
Thanks to @andygrove for this feature, with reviews from @coderfender and @sunchao. Related PRs: #5394.
Correctness Fixes¶
1.1.0 fixes many cases where Comet returned different results from Spark, failed where Spark succeeds, or accepted input that Spark rejects. Besides the Iceberg and memory fixes above, these are the ones most likely to affect 1.0.0 users. The change log has the rest.
Wrong results¶
- Decimal
SUMreturned NULL, or raised an overflow error under ANSI, when an intermediate sum overflowed but the final result fit (#6041, @dwsmith1983). - After a late shuffle fallback,
avgcould return NULL andcollect_list/collect_setcould produce mismatched buffers (#5421, @sunchao). - Exchange reuse could share one shuffle between plans that differ, such as
COUNT(*) + 1andCOUNT(*) - 1, semi and anti joins, orexplodeandexplode_outer(#5470, @sunchao and #5828, @ErikBPF). - Two ABFS containers in the same storage account shared a cached object store, so a read could return the other container's data (#5053, @peterxcli).
- Dictionary-encoded values hashed differently from the same values decoded, and a null struct hashed its fields, which affected joins, aggregates and shuffle partitioning (#5757 and #5754, @viirya).
- Parquet field names containing non-ASCII characters that differ only in case read as NULL (#5602, @comphead).
IN,InSet, nested=,arrays_overlapandarray_positionnow treat-0.0and0.0, and every NaN encoding, as Spark does (#6073, @mizulun, #5235, @divyankshah and #5472, @sunchao).- Decimal to double and float casts were off by one unit in the last place for most
DECIMAL(38,18)values (#5684, @peterxcli). - String to timestamp casts now follow Spark's parsing rules for short fields, time zones and signed years (#5682 and #5858, @peterxcli).
Errors Spark raises that Comet did not¶
- Casts and expressions routed through codegen dispatch could skip ANSI errors raised inside a constant subexpression (#5623, @andygrove).
- Rejected
TIMESTAMP_NTZcasts returned NULL under ANSI instead of raisingCAST_INVALID_INPUT(#5752, @peterxcli). - Out-of-range Parquet
TIMESTAMP_MILLISvalues, top-level or nested, silently wrapped (#5177 and #5740, @peterxcli). - Nested Parquet struct, list and map fields now follow Spark's conversion rules instead of returning NULL on overflow or accepting values Spark rejects (#5681, @peterxcli).
Query failures and crashes¶
collect_listandcollect_setover nested arguments failed with "column types must match schema types" (#5159, @andygrove).- Native shuffle failed with a 2 GB task serialization error on jobs with very many partitions (#5392, @parthchandra).
- A Scala UDF from a user jar failed with a
ClassCastException(#5282, @andygrove). rpadandlpadpanicked on a NULL length (#5680, @peterxcli).- Structs with duplicate field names failed the task in native shuffle, and panicked in the native Parquet scan (#5866, @dwsmith1983 and #5786, @ErikBPF).
Hangs and resource use¶
- The JVM hung on exit when an application returned from
mainwithout callingspark.stop()(#5748, @zhangfengcdt). - The native scan busy-polled while waiting on S3 or HDFS reads, keeping one core per task at 100% (#6219, @mixermt and @andygrove).
- One task could force-spill another task's shuffle buffers, and a failed shuffle write leaked its memory reservation (#5493 and #5461, @peterxcli).
Upgrading to 1.1.0¶
A few changes in 1.1.0 can affect a deployment. The Comet Upgrade Guide has the details.
Platform¶
- JDK 17 or later is required. JDK 11 support was removed, as announced in 1.0.0.
- Spark 3.4 is still deprecated. Comet still publishes Spark 3.4 binaries, but Spark's SQL test suite runs against 3.4 only on demand, so 3.4-specific regressions are more likely. We recommend Spark 3.5 or later.
Size settings now read correctly¶
1.0.0 misread three size settings, so jobs that set them may behave differently:
spark.memory.offHeap.sizeas a bare number of bytes was read as MiB, which madefair_unified's per-operator shares about a million times too large. With correct shares, operators can spill sooner, and an operator that can't spill can fail when it exceeds its share. Sizes with a unit, like16g, aren't affected.spark.comet.exec.memoryPool=greedy_unifiedrestores the old behavior.spark.comet.maxTempDirectorySizewith a unit, like10g, was ignored in favor of the 100 GB default. It's now enforced, so a query that spills past it fails.spark.comet.shuffle.native.writeBufferSizewith a unit was treated as bytes, so64mmeant 64 bytes. It now means what it says. These buffers use untracked native memory, so check large values againstspark.executor.memoryOverhead.
Malformed values for spark.comet.maxTempDirectorySize and spark.comet.explain.native.enabled now fail the
query instead of falling back to the default.
More memory before spilling¶
With the fair_unified fix, tasks with several operators can reserve more memory before spilling than in
0.15.0 through 1.0.0, especially on executors running few tasks at once. If you sized executors for one of those
releases, check the memory usage log for headroom.
Conditions for enabling Comet¶
Comet requires Spark's off-heap memory. CometPlugin already enforced this, and 1.1.0 applies the same check
when CometSparkSessionExtensions is registered directly, disabling Comet with a warning. Both checks read
spark.memory.offHeap.enabled from the SparkContext, so set it when the application starts, not on a
SparkSession.builder after the context exists.
Comet also checks the shuffle manager that's actually running, not the session's spark.shuffle.manager. A
session that names CometShuffleManager after the context started with a different one now runs without Comet,
with a warning, instead of failing with a ClassCastException.
Deprecated and removed settings¶
spark.comet.exec.memoryPool.fractionis deprecated and will be removed in a future major release. It still works, and the driver warns when it's set.- Testing-only on-heap settings were removed:
spark.comet.memoryOverhead,spark.comet.exec.onHeap.memoryPool, andspark.comet.shuffle.jvm.memoryFactor(and its old name,spark.comet.columnar.shuffle.memory.factor). The versioning policy exempts testing settings, and Comet ignores them if they're still set.
Compatibility¶
Supported platforms:
- Spark 3.4.3 with Java 17 and Scala 2.12/2.13 (deprecated)
- Spark 3.5.9 with Java 17 and Scala 2.12/2.13
- Spark 4.0.4 with Java 17/21 and Scala 2.13
- Spark 4.1.3 with Java 17/21 and Scala 2.13
- Spark 4.2.0 with Java 17 and Scala 2.13 (experimental, for early evaluation only)
See the Spark Version Compatibility page for known limitations specific to each version.
This release builds on DataFusion 55.1 and Arrow 59.2.
Get Started with Comet 1.1.0¶
Follow the Comet 1.1.0 Installation Guide to get up and running, then point Comet at your existing Spark workloads.