Apache DataFusion Comet 1.0.0 Release
Posted on: Fri 07 August 2026 by pmc
The Apache DataFusion PMC is pleased to announce version 1.0.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 six weeks of development since 0.17.0 and consists of 244 commits from 23 contributors. See the change log for the full list of changes.
The Road to 1.0¶
Comet was donated to the Apache DataFusion project in March 2024 and cut its first release, 0.1.0, five
months later with support for 13 operators and 106 expressions. Since then, the project has shipped 20
releases and drawn contributions from more than 120 developers, and the codebase now recognizes over 400
Spark expressions. Operator coverage has grown alongside it: 1.0 accelerates each of Spark's four join
operators, window functions, generators (explode, explode_outer, posexplode, and posexplode_outer
over arrays), sampling, in-memory table scans, and a fully native shuffle.
The 1.0 release marks the point at which Comet begins following semantic versioning. Users upgrading within the 1.x line can expect backward-compatible changes only; features slated for removal will be deprecated in a minor release before being dropped in the next major version. This is why the deprecations of JDK 11 and Spark 3.4 announced below are scheduled for 1.1 rather than landing in 1.0 itself.
Support for Spark 4.0+ with ANSI mode¶
Comet 1.0.0 supports Spark versions 3.4 through 4.1, with experimental support for 4.2. Comet fully supports Spark's ANSI mode, which is enabled by default starting with Spark 4.0.
Correctness Testing¶
It is important that queries accelerated by Comet produce the same results as Spark. Correctness checking has always been a large effort in Comet development, but the approach has evolved over time.
- Upstream Spark tests: Comet runs Spark's own test suite with Comet enabled, providing more than 24,000 unit tests effectively for free. These tests run in Comet's CI for all supported Spark versions.
- Scala tests: end-to-end queries that run with Comet enabled versus disabled, checking that results match.
- Fuzz testing: many of the Scala tests generate randomized data to catch regressions around edge cases such as nulls, NaN, Infinity, and timezone issues.
- Comet SQL tests: a sqllogictest-inspired approach that makes end-to-end tests easier to write.
- Generative AI audits: agentic skills sweep every expression, comparing Comet's implementation to Spark's source and ensuring tests cover important edge cases.
Performance¶
The early Comet releases provided a very modest speedup and the published benchmark results were based on running TPC workloads at small scale factors on a single node. There are now independent benchmark results published by AWS Labs that show significant speedups for TPC-DS @ 3TB running in EKS.
Codegen Dispatch¶
Comet 0.17.0 introduced a new approach to filling gaps in expression coverage. In earlier releases, whenever Comet's planner encountered an expression that lacked a native Rust implementation, it fell back to executing an entire subtree of the plan in Spark. That required converting Arrow columns back to Spark rows before the expression ran and back to Arrow after, and the cost was often enough to erase the speedup Comet had bought elsewhere in the plan.
Codegen dispatch narrows that fallback to the expression itself: the batch stays in the Comet pipeline and Comet invokes Spark's own generated code for just the missing expression, leaving the rest of the query running natively. Four consequences are worth calling out.
- Coverage. Expressions that would previously have blocked native execution of a whole subtree are now supported immediately, without a Rust port.
- Compatibility. For categories where a native reimplementation would inevitably diverge from Spark's semantics — regular expressions being the canonical case, given the gap between Java's regex engine and any Rust or C++ equivalent — codegen dispatch delivers bit-for-bit Spark parity because it is Spark's implementation.
- Expression fusion. A dispatched expression tree (i.e., nested expressions) is compiled into a single method, so the Arrow
input reads, the expression evaluation, and the Arrow output writes are fused together. The compiler
is free to optimize across the whole tree, and no intermediate Arrow
RecordBatchis materialized between one expression and the next. - Scala and Java UDFs. User-defined functions are compiled to the same codegen surface as built-in expressions, so they can flow through codegen dispatch without any change from the user. Queries that were previously disqualified from acceleration only because they contained a UDF can now benefit as long as the surrounding operators are supported. See the Scala and Java UDF guide for details.
Comet 1.0 widens the mechanism in three ways.
The first is the biggest. An expression that opts into codegen dispatch previously reached the dispatcher only when Comet reported it as incompatible for the given input; an unsupported report still sent the whole subtree back to Spark. In 1.0 both support levels route through the dispatcher, so an input that Spark handles and Comet's native code does not now stays inside the Comet pipeline.
Second, casts join the same path. Cast expressions that Comet declines to run natively — including legacy
configuration variants such as spark.sql.legacy.castComplexTypesToString.enabled — are now dispatched
rather than falling back, and more string, array, and interval expressions were opted in as well.
Third, the path is now visible. Comet's extended explain output reports native versus codegen-dispatch coverage for a plan, so you can see which path each expression actually took rather than inferring it from the absence of a fallback reason.
Improvements since 0.17.0¶
The rest of this post covers what is new since the 0.17.0 release.
Experimental PyArrow UDF Support¶
This release adds experimental support for accelerated PyArrow UDFs, allowing PyArrow-based user-defined functions to participate in native execution instead of forcing a fallback to Spark. When the feature is disabled, Comet now hints at the native PyArrow UDF path in its fallback reasons so users know the option exists. This is an early-stage feature and we welcome feedback from users experimenting with it.
New Expression and Aggregate Support¶
This release expands the set of Spark expressions and aggregates that are accelerated by Comet:
- Aggregates:
approx_percentile/percentile_approx, exactpercentile/median,approx_count_distinct, and nativecollect_list/array_agg. - Cast: Cast expressions where the native implementation is marked as incompatible or unsupported are now routed through codegen dispatch.
- Grouping:
grouping()andgrouping_id(). - Intervals: interval types via
make_ym_intervalandmake_dt_interval,CalendarIntervalType,multiply_dt_interval, and interval codegen dispatch for nested values and native shuffle. - String:
base64,split_partviaStringSplitSQL, nativelevenshtein, and nativerandstranduuid— both bit-for-bit compatible with Spark for a given seed. - Array / map:
array_prepend, theshuffle()array function,size()forMapType, andElementAtoverMapType. - Date/time: native
TimestampNTZinputs forhour/minute/secondandPreciseTimestampConversionfor native time-window grouping. - Windows: extended native window function support and Spark 4 decimal window average.
Faster Parquet Scans¶
Parquet reads pick up several improvements as well. Full Parquet metadata, including the page index, is now
cached via DataFusion's CachedParquetFileReaderFactory; identity casts are unwrapped in the schema adapter
so Parquet statistics pruning can engage; filter pushdown configuration has been revised; the native scan
passes a metadata size hint so a single read usually captures the footer; and the native Parquet scan seeds
its reader options from the session config so Parquet settings you already set take effect.
Iceberg Table Format V3¶
Comet now supports Iceberg 1.11 and its first Iceberg table format V3 feature: full table encryption. Other V3 features like deletion vectors and new data types (e.g., VARIANT) fall back gracefully. The native Iceberg scan supports the _pos, _spec, _file, and
_partition metadata columns, sizes delete files correctly to avoid dropped deletes, disambiguates scans that
share a metadata_location, and dedupes residuals and delete files in the native scan serde. A prior case
where Iceberg native scan exchange reuse with different pushed filters could produce wrong results is also
fixed.
Native Expression Performance¶
Many native expression implementations have been optimized to more efficiently leverage Arrow kernels or to avoid per-row builders.
- Casts between numeric, string, decimal, and date types, including a faster float-to-decimal cast, an
optimized integer-to-integer cast, shared no-overflow fast paths in
CheckOverflowandDecimalRescaleCheckOverflow, and acast_binary_to_stringthat is up to 27x faster on binary-format styles. - JSON, regex, and URL parsing:
get_json_object,regexp_extract, andparse_url. - Date/time and decimal kernels:
date_trunc,spark_ceil, and a vectorizedspark_unscaled_value. - String and array kernels:
lpad,unhex,size,arrays_overlap,escape_string, and thetry_*arithmetic kernel.
To make this kind of work repeatable, the release also adds a scalar expression optimization guide documenting how to benchmark a kernel, keep its output bit-identical to Spark, and gate changes on a no-regression check.
Deprecation Notice¶
With the move to a stable 1.0 release line, Comet begins deprecating older platforms under its versioning policy:
- JDK 11 is deprecated and scheduled for removal in Comet 1.1.0.
- Apache Spark 3.4 is deprecated and scheduled for removal in Comet 1.1.0.
Comet aligns its Spark support window with upstream Apache Spark maintenance. Spark 3.4 is no longer maintained upstream, so under the versioning policy it is deprecated in the first Comet minor release after that point and removed in the following one. Comet 1.0.0 still builds and publishes Spark 3.4 binaries.
Users on these platforms should plan to move to JDK 17+ and Spark 3.5 or later before upgrading to 1.1.0.
Compatibility¶
Supported platforms include:
- Spark 3.4.3 with Java 11/17 and Scala 2.12/2.13 (deprecated, removal in 1.1.0)
- Spark 3.5.9 with Java 11/17 and Scala 2.12/2.13
- Spark 4.0.4 with Java 17 and Scala 2.13
- Spark 4.1.3 with Java 17/21 and Scala 2.13
- Spark 4.2 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 upgrades to DataFusion 54.1 and Arrow 58.4.
Get Started with Comet 1.0.0¶
Ready to try it out? Follow the Comet 1.0.0 Installation Guide to get up and running, then point Comet at your existing Spark workloads and see the speedup for yourself.