pingzh commented on code in PR #205:
URL: https://github.com/apache/datafusion-site/pull/205#discussion_r4197798158


##########
content/blog/2026-10-01-datafusion-comet-1.1.0.md:
##########
@@ -0,0 +1,602 @@
+---
+layout: post
+title: Apache DataFusion Comet 1.1.0 Release
+date: 2026-10-01
+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.1.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 seven weeks of development since 1.0.0: 398 
commits from 40 contributors. See the
+[change log] for the full list of changes.
+
+[change log]: 
https://github.com/apache/datafusion-comet/blob/branch-1.1/docs/source/changelog/1.1.0.md
+
+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**: Comet's side is in place, but it 
needs a Celeborn client API that no
+  released Celeborn version provides yet.

Review Comment:
   This makes both the native implementation and the current fallback with 
released clients explicit.
   
   ```suggestion
   - **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.
   ```



##########
content/blog/2026-10-01-datafusion-comet-1.1.0.md:
##########
@@ -0,0 +1,602 @@
+---
+layout: post
+title: Apache DataFusion Comet 1.1.0 Release
+date: 2026-10-01
+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.1.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 seven weeks of development since 1.0.0: 398 
commits from 40 contributors. See the
+[change log] for the full list of changes.
+
+[change log]: 
https://github.com/apache/datafusion-comet/blob/branch-1.1/docs/source/changelog/1.1.0.md
+
+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**: Comet's side is in place, but it 
needs a Celeborn client API that no
+  released Celeborn version provides yet.
+- **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](https://github.com/apache/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:
+
+1. **`IcebergWrite`** writes data files on the executors and returns each 
task's commit message.
+2. **`IcebergCommit`** collects 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](https://github.com/apache/datafusion-comet/issues/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].
+
+[Iceberg Writes guide]: 
https://datafusion.apache.org/comet/user-guide/latest/iceberg-writes.html
+
+## 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.
+
+<img
+src="/blog/images/comet-1.1.0/comet-executor-memory.svg"
+width="100%"
+class="img-fluid"
+alt="The executor container holds the JVM heap, the off-heap memory pool, and 
the memory overhead. Spark and Comet share the off-heap pool, where Comet's 
sorts, joins, aggregations, and shuffles reserve memory. The memory overhead 
holds the JVM's own overhead plus Comet's native memory that the pool does not 
track."
+/>
+
+### 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.
+
+[tuning guide]: 
https://datafusion.apache.org/comet/user-guide/latest/tuning/memory.html
+
+### 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_unified` limited 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.size` written as
+  a bare number. See [Upgrading to 1.1.0](#upgrading-to-110) 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].
+
+[memory management guide]: 
https://datafusion.apache.org/comet/contributor-guide/memory_management.html
+
+## More Iceberg Improvements
+
+- **V3 deletion vectors** are applied on native scans.
+- **Iceberg system functions** (`bucket`, `truncate`, `years`, `months`, 
`days`, and `hours`) 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) = 2` combined with 
another predicate was pushed to
+  the native scan as `id = 2`, returning too few rows.
+- Tables partitioned by an unknown transform can be read natively, and `IS 
NULL` / `IS NOT NULL` checks 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.
+
+It does not run with a released Celeborn yet. Native shuffle needs Celeborn to 
report reliably when an in-flight
+push has completed, and no released client does, including 0.7.0, the latest 
release. Those clients keep ordinary
+shuffle even when native shuffle is requested. Once a Celeborn release 
provides that API, set
+`spark.shuffle.manager` to 
`org.apache.spark.sql.comet.execution.shuffle.CometCelebornShuffleManager` and
+`spark.comet.shuffle.mode=native`; the default `auto` mode keeps ordinary 
shuffle. Native shuffle over Celeborn
+doesn't support `spark.io.encryption.enabled=true`, and Celeborn isn't bundled 
with Comet. The
+[Celeborn tuning guide] has the full requirements.

Review Comment:
   The current blocker is Comet's push-completion compatibility checks. The 
follow-up discussion in 
[#6523](https://github.com/apache/datafusion-comet/issues/6523#issuecomment-6019598622)
 considers changes to those checks as well as a supported completion API, so I 
would avoid implying that a future Celeborn release alone will enable the 
native path. This also makes clear that other Comet operators can still be 
accelerated.
   
   ```suggestion
   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](https://github.com/apache/datafusion-comet/issues/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.
   ```



-- 
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]

Reply via email to