Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
72 commits
Select commit Hold shift + click to select a range
65be8b7
fix: retain Comet execution for warehouse workaround regressions
notimesea Sep 25, 2026
45fad9e
fix: preserve zero partitions when native scan files are pruned
notimesea Sep 25, 2026
d153b3a
revert: drop the sort merge reservation cap
msafonov Sep 27, 2026
e90fcec
feat: spill all non-streaming native window shapes
msafonov Sep 27, 2026
e435466
test: compare spilled window shapes with Spark
msafonov Sep 27, 2026
943936f
test: reproduce native sort failing to spill after the Spark share sh…
msafonov Sep 27, 2026
6eaf12e
build: vendor datafusion-physical-plan 55.1.0 unchanged
msafonov Sep 27, 2026
880575b
fix: keep native sort spill memory when the Spark share shrinks
msafonov Sep 27, 2026
b3cd74a
fix: keep the spill merge's headroom across all merge passes
msafonov Sep 27, 2026
d7d47f3
fix: reserve the encoded rows a merge keeps for reuse
msafonov Sep 27, 2026
d82e3a3
fix: count a single-column merge's sort key once
msafonov Sep 27, 2026
69e6389
fix: leave as much memory again for the final spill merge's consumers
msafonov Sep 27, 2026
a4dc994
fix: let a sort's spill merge split a run batch wider than its budget
msafonov Sep 27, 2026
426ca8e
fix: keep a sort spill's batches reserved until they are written
msafonov Sep 27, 2026
1b8d480
fix: count buffers a sort's chunks share once
msafonov Sep 27, 2026
6ae0323
fix: reserve a buffer a sort's buffered batches share once
msafonov Sep 27, 2026
5b75a38
fix: keep the merge headroom for a sort's final in-memory merge
msafonov Sep 27, 2026
399e6c9
fix: reserve what a sort's spill merge holds for each run
msafonov Sep 28, 2026
4ff3708
fix: restore sort merge reservation cap and native test contract
msafonov Sep 28, 2026
ea38aad
fix: allow small Comet batches and cap JVM shuffle batches at use time
msafonov Sep 28, 2026
ca6b963
debug: measure wide batches at shuffle decode and FFI output
msafonov Sep 28, 2026
5f6cbb8
test: bound unreserved sort output at the JVM handoff
msafonov Sep 28, 2026
39f49a0
test: measure wide binary payload copies in sort kernels
msafonov Sep 28, 2026
be99cc1
fix: avoid repeated copies of wide binary payloads in native sort
msafonov Sep 28, 2026
5da7791
test: measure dictionary expansion before wide binary sort
msafonov Sep 28, 2026
4956d83
fix: select direct shuffle reads during initial planning
msafonov Sep 28, 2026
c0b0448
test: benchmark binary columnar-to-row allocation overhead
msafonov Sep 28, 2026
3d99265
perf: avoid intermediate binary copies in row materialization
msafonov Sep 28, 2026
a52c168
test: cover default Kryo broadcast payloads
msafonov Sep 28, 2026
d9ea392
perf: reuse unique dictionary values instead of casting in JVM shuffl…
msafonov Sep 28, 2026
2eda535
feat: revert isolated native operators feeding Spark operators
msafonov Sep 28, 2026
edaf6ca
feat: run each query stage in one engine and pick shuffle formats per…
msafonov Sep 29, 2026
95dc9ae
perf: evaluate small window partitions together and coalesce window o…
msafonov Sep 29, 2026
85bef53
fix: do not fail an aggregate whose first reservation was refused whe…
msafonov Sep 29, 2026
b65fdbe
fix: let the task's other memory consumers make the JVM shuffle write…
msafonov Sep 29, 2026
0a92f71
revert: remove UnifyStageEngines
msafonov Sep 29, 2026
4f92985
revert: remove RevertIsolatedNativeOperators
msafonov Sep 29, 2026
01df400
feat: pick shuffle and broadcast formats from the engines on both sides
msafonov Sep 29, 2026
fd190a6
feat: choose each operator's engine by a cost over the whole plan
msafonov Sep 29, 2026
c69bd9c
feat: native Delta scan (port of apache/datafusion-comet#5365)
msafonov Sep 29, 2026
a435236
fix(delta): claim GCS scans whose Hadoop auth keys only select ADC
msafonov Sep 29, 2026
88cfb29
fix: do not hash decimals above precision 18 in native shuffle (port …
msafonov Sep 29, 2026
2d3cc89
feat: run sorts of wide rows with narrow keys in Spark when Spark rea…
msafonov Sep 29, 2026
909335e
feat: move sorts with variable-width payloads to Spark on the initial…
msafonov Sep 29, 2026
3b37b4a
test: time native sorts of wide binary payloads end to end
msafonov Sep 29, 2026
d796f09
perf: sort the keys of wide rows and gather each payload once in nati…
msafonov Sep 29, 2026
f66011d
fix: gather a native sort's output from the batches it references
msafonov Sep 30, 2026
9ff3736
refactor: keep wide-row sorts in Spark across AQE re-plans without pe…
msafonov Sep 30, 2026
b2771dd
fix: count the encoded rows a merge keeps instead of summing them per…
msafonov Sep 30, 2026
f7a13eb
perf: coalesce small input batches of a native sort before buffering …
msafonov Sep 30, 2026
0e17ce3
fix(delta): do not resolve subqueries when planning reads the scan's …
msafonov Oct 1, 2026
d913126
fix: rebase the offsets of sliced build batches when coalescing a bro…
msafonov Oct 1, 2026
f761b31
fix: keep the streamed order of a sort-merge join with a join filter
msafonov Oct 1, 2026
2d30872
perf: coalesce small blocks when reading a Comet shuffle
msafonov Oct 1, 2026
c2562fb
feat: keep shuffles of wide rows in Spark
msafonov Oct 1, 2026
9560d4a
perf: reuse generated columnar-to-row projections across tasks
msafonov Oct 1, 2026
b83e917
feat: decide wide-row sorts and shuffles by leaf columns
msafonov Oct 1, 2026
568955d
fix: spill a large in-memory sort before producing its output
msafonov Oct 2, 2026
94abcc0
perf: gather wide shuffle rows from the batches each partition refere…
msafonov Oct 2, 2026
d1148b2
feat: price the cost-based engine choice by rows and leaf columns
msafonov Oct 2, 2026
cea0ae5
feat: price the cost-based engine choice per row and decide it on eve…
msafonov Oct 2, 2026
f9793a7
feat: enable the cost-based engine choice by default
msafonov Oct 2, 2026
1f7b89b
feat: price the cost-based engine choice with the calib29 measurements
msafonov Oct 3, 2026
6b11f2e
fix: clone the shared sort source with Arc::clone in the handoff test
msafonov Oct 3, 2026
71a7c07
feat: keep the cost-based engine choice disabled by default
msafonov Oct 3, 2026
2b28fed
test: give the child sort spill metrics test less memory so the sort …
msafonov Oct 3, 2026
bc64dfc
feat: keep PartitionAggregateWindowExec disabled by default
msafonov Oct 3, 2026
8e06527
fix: pool C2R projections only after every reader of their rows is done
msafonov Oct 3, 2026
bb53574
ci: list the missing suites so check-suites.py passes
msafonov Oct 3, 2026
c8b47b4
style: drop a needless string interpolator flagged by scalafix
msafonov Oct 3, 2026
1a52c8c
build: compile with Scala 2.13 and Rust 1.99
msafonov Oct 3, 2026
4f150e1
style: drop an unused pattern binding flagged by scalafix on Scala 2.13
msafonov Oct 3, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
The table of contents is too big for display.
Diff view
Diff view
  •  
  •  
  •  
10 changes: 10 additions & 0 deletions .github/workflows/pr_build_linux.yml
Original file line number Diff line number Diff line change
Expand Up @@ -515,11 +515,13 @@ jobs:
org.apache.spark.sql.comet.execution.shuffle.CometDiskBlockWriterSuite
org.apache.comet.exec.CometShuffleEncryptionSuite
org.apache.comet.exec.CometShuffleManagerSuite
org.apache.comet.exec.CometShuffleReadCoalesceSuite
org.apache.comet.exec.CometAsyncShuffleSuite
org.apache.comet.exec.DisableAQECometShuffleSuite
org.apache.comet.exec.DisableAQECometAsyncShuffleSuite
org.apache.spark.shuffle.comet.CometUnboundedShuffleMemoryAllocatorSuite
org.apache.spark.shuffle.sort.SpillSorterSuite
org.apache.spark.shuffle.sort.CometShuffleExternalSorterSpillSuite
- name: "exec"
value: |
org.apache.comet.exec.CometAggregateSuite
Expand All @@ -528,8 +530,11 @@ jobs:
org.apache.comet.exec.CometEmptyRelationExecSuite
org.apache.comet.exec.CometInMemoryCacheSuite
org.apache.comet.exec.CometInMemoryCacheKryoSuite
org.apache.spark.sql.comet.CometBroadcastKryoPayloadSuite
org.apache.spark.sql.comet.CometBroadcastDefaultKryoSuite
org.apache.comet.exec.CometGenerateExecSuite
org.apache.comet.exec.CometWindowExecSuite
org.apache.comet.exec.CometPartitionAggregateWindowSuite
org.apache.comet.exec.CometJoinSuite
org.apache.spark.sql.comet.CometMapInBatchSuite
org.apache.spark.sql.execution.python.CometArrowPythonRunnerSuite
Expand Down Expand Up @@ -557,12 +562,17 @@ jobs:
org.apache.spark.sql.comet.CometPlanEqualitySuite
org.apache.comet.rules.CometEmptyRelationExecRuleSuite
org.apache.comet.rules.RevertNativeForTransitionHeavyStagesSuite
org.apache.comet.rules.ChooseBoundaryFormatsSuite
org.apache.comet.rules.CostBasedEngineChoiceSuite
org.apache.comet.rules.WideRowSortFallbackSuite
org.apache.comet.rules.WideRowShuffleFallbackSuite
org.apache.spark.sql.CometTPCDSQuerySuite
org.apache.spark.sql.CometTPCDSQueryTestSuite
org.apache.spark.sql.CometTPCHQuerySuite
org.apache.spark.sql.comet.CometTPCDSV1_4_PlanStabilitySuite
org.apache.spark.sql.comet.CometTPCDSV2_7_PlanStabilitySuite
org.apache.spark.sql.comet.CometTaskMetricsSuite
org.apache.spark.sql.comet.CometBatchRowProjectionSuite
org.apache.spark.sql.comet.CometDppFallbackRepro3949Suite
org.apache.spark.sql.comet.CometShuffleFallbackStickinessSuite
org.apache.spark.sql.comet.PlanDataInjectorSuite
Expand Down
10 changes: 10 additions & 0 deletions .github/workflows/pr_build_macos.yml
Original file line number Diff line number Diff line change
Expand Up @@ -163,11 +163,13 @@ jobs:
org.apache.spark.sql.comet.execution.shuffle.CometDiskBlockWriterSuite
org.apache.comet.exec.CometShuffleEncryptionSuite
org.apache.comet.exec.CometShuffleManagerSuite
org.apache.comet.exec.CometShuffleReadCoalesceSuite
org.apache.comet.exec.CometAsyncShuffleSuite
org.apache.comet.exec.DisableAQECometShuffleSuite
org.apache.comet.exec.DisableAQECometAsyncShuffleSuite
org.apache.spark.shuffle.comet.CometUnboundedShuffleMemoryAllocatorSuite
org.apache.spark.shuffle.sort.SpillSorterSuite
org.apache.spark.shuffle.sort.CometShuffleExternalSorterSpillSuite
- name: "exec"
value: |
org.apache.comet.exec.CometAggregateSuite
Expand All @@ -176,8 +178,11 @@ jobs:
org.apache.comet.exec.CometEmptyRelationExecSuite
org.apache.comet.exec.CometInMemoryCacheSuite
org.apache.comet.exec.CometInMemoryCacheKryoSuite
org.apache.spark.sql.comet.CometBroadcastKryoPayloadSuite
org.apache.spark.sql.comet.CometBroadcastDefaultKryoSuite
org.apache.comet.exec.CometGenerateExecSuite
org.apache.comet.exec.CometWindowExecSuite
org.apache.comet.exec.CometPartitionAggregateWindowSuite
org.apache.comet.exec.CometJoinSuite
org.apache.spark.sql.comet.CometMapInBatchSuite
org.apache.spark.sql.execution.python.CometArrowPythonRunnerSuite
Expand Down Expand Up @@ -205,12 +210,17 @@ jobs:
org.apache.spark.sql.comet.CometPlanEqualitySuite
org.apache.comet.rules.CometEmptyRelationExecRuleSuite
org.apache.comet.rules.RevertNativeForTransitionHeavyStagesSuite
org.apache.comet.rules.ChooseBoundaryFormatsSuite
org.apache.comet.rules.CostBasedEngineChoiceSuite
org.apache.comet.rules.WideRowSortFallbackSuite
org.apache.comet.rules.WideRowShuffleFallbackSuite
org.apache.spark.sql.CometTPCDSQuerySuite
org.apache.spark.sql.CometTPCDSQueryTestSuite
org.apache.spark.sql.CometTPCHQuerySuite
org.apache.spark.sql.comet.CometTPCDSV1_4_PlanStabilitySuite
org.apache.spark.sql.comet.CometTPCDSV2_7_PlanStabilitySuite
org.apache.spark.sql.comet.CometTaskMetricsSuite
org.apache.spark.sql.comet.CometBatchRowProjectionSuite
org.apache.spark.sql.comet.CometDppFallbackRepro3949Suite
org.apache.spark.sql.comet.CometShuffleFallbackStickinessSuite
org.apache.spark.sql.comet.PlanDataInjectorSuite
Expand Down
70 changes: 70 additions & 0 deletions contrib/delta-spark/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,70 @@
<!--
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.
-->

# Comet Delta Lake Contrib (experimental)

Native Delta Lake reads for Comet. Delta tables are scanned through Comet's
existing native Parquet reader, so they get row-group pruning, page-index
pruning, and filter pushdown, with deletion vectors applied inside the scan.

Support is experimental and explicitly opt-in. Two things are required:

1. This module's jar (`comet-contrib-delta-spark`) on the classpath, alongside
`delta-spark`. It is never bundled into `comet-spark`; without it, Comet
has no Delta surface at all.
2. `spark.comet.scan.delta.enabled=true`. The default is `false`, so the jar
alone does nothing.

Unsupported tables and features fall back to Spark's reader. See the
[user guide](https://datafusion.apache.org/comet/user-guide/latest/delta.html)
for configuration details.

## Supported versions

| Spark | Delta | Status |
| ----- | -------------- | --------------------------------------------- |
| 3.5 | 3.3.x | supported |
| 4.0 | 4.0.x | supported |
| 4.1 | 4.3.x | supported |
| 3.4 | delta-core 2.4 | not supported (older Delta, would need shims) |
| 4.2 | none released | inert until Delta ships a Spark 4.2 release |

## Building and testing

The module builds under the `delta` Maven profile. It resolves `comet-spark`
from the local Maven repository, so install `common` and `spark` from the same
checkout immediately before, as CI does, with the `delta` profile active so that
install produces the spark test-jar the contrib suites depend on; a stale sibling
install is the trap the contributor guide warns about:

```shell
./mvnw -Pspark-3.5,delta install -pl common,spark -DskipTests
./mvnw -Pspark-3.5,delta install -pl contrib/delta-spark
```

Run the test suites the same way (`test` instead of `install` on the second
line). CI runs them via `.github/workflows/delta_contrib_test.yml`: on Spark
3.5 in the merge queue, and on 4.0 and 4.1 nightly or on a pull request
carrying the `run-delta-tests` label.
`CometDeltaS3Suite` starts MinIO through Testcontainers and cancels itself when
no Docker daemon is reachable; setting `COMET_DELTA_S3_REQUIRED=1`, as the
`contrib-delta-s3` CI job does, turns that cancel into a suite failure.

`dev/` contains a benchmark script (`bench_delta_comet.py`) and a harness for
running Delta's own test suites against Comet (`run-delta-regression.sh`).
Loading
Loading