hsiang-c commented on code in PR #201:
URL: https://github.com/apache/datafusion-site/pull/201#discussion_r3739085648
##########
content/blog/2026-08-07-datafusion-comet-1.0.0.md:
##########
@@ -0,0 +1,221 @@
+---
+layout: post
+title: Apache DataFusion Comet 1.0.0 Release
+date: 2026-08-07
+author: pmc
+categories: [subprojects]
+---
+
+<!--
+{% comment %}
+Licensed to the Apache Software Foundation (ASF) under one or more
+contributor license agreements. See the NOTICE file distributed with
+this work for additional information regarding copyright ownership.
+The ASF licenses this file to you under the Apache License, Version 2.0
+(the "License"); you may not use this file except in compliance with
+the License. You may obtain a copy of the License at
+
+http://www.apache.org/licenses/LICENSE-2.0
+
+Unless required by applicable law or agreed to in writing, software
+distributed under the License is distributed on an "AS IS" BASIS,
+WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+See the License for the specific language governing permissions and
+limitations under the License.
+{% endcomment %}
+-->
+
+[TOC]
+
+The Apache DataFusion PMC is pleased to announce version 1.0.0 of the
[Comet](https://datafusion.apache.org/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.
+
+[change log]:
https://github.com/apache/datafusion-comet/blob/main/docs/source/changelog/1.0.0.md
+
+## 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.
+
+[donated]: https://datafusion.apache.org/blog/2024/03/06/comet-donation/
+
+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.
+
+[semantic versioning]:
https://datafusion.apache.org/comet/about/versioning_policy.html
+
+### 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](https://awslabs.github.io/data-on-eks/docs/benchmarks/spark-datafusion-comet-benchmark).
+
+### 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
`RecordBatch` is 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.
+
+[Scala and Java UDF guide]:
https://datafusion.apache.org/comet/user-guide/latest/scala_java_udfs.html
+
+## 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`, exact
`percentile` / `median`,
+ `approx_count_distinct`, and native `collect_list` / `array_agg`.
+- **Cast**: Cast expressions where the native implementation is marked as
incompatible or unsupported
+ are now routed through codegen dispatch.
+- **Grouping**: `grouping()` and `grouping_id()`.
+- **Intervals**: interval types via `make_ym_interval` and `make_dt_interval`,
`CalendarIntervalType`,
Review Comment:
(nit) Maybe we can call out interval support b/c adding new types is
infrequent.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]