Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .github/workflows/pr_build_linux.yml
Original file line number Diff line number Diff line change
Expand Up @@ -575,6 +575,7 @@ jobs:
org.apache.spark.sql.CometCollationSuite
org.apache.spark.sql.CometSortCollationSuite
org.apache.comet.CometFuzzAggregateSuite
org.apache.spark.sql.comet.CometBroadcastInputLifecycleSuite
org.apache.spark.sql.comet.execution.arrow.CometArrowStreamSuite
org.apache.spark.sql.comet.execution.arrow.CachedBatchRowIteratorSuite
org.apache.spark.sql.comet.execution.arrow.CometArrowWriterSuite
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/pr_build_macos.yml
Original file line number Diff line number Diff line change
Expand Up @@ -279,6 +279,7 @@ jobs:
org.apache.spark.sql.CometCollationSuite
org.apache.spark.sql.CometSortCollationSuite
org.apache.comet.CometFuzzAggregateSuite
org.apache.spark.sql.comet.CometBroadcastInputLifecycleSuite
org.apache.spark.sql.comet.execution.arrow.CometArrowStreamSuite
org.apache.spark.sql.comet.execution.arrow.CachedBatchRowIteratorSuite
org.apache.spark.sql.comet.execution.arrow.CometArrowWriterSuite
Expand Down
2 changes: 2 additions & 0 deletions dev/ensure-jars-have-correct-contents.sh
Original file line number Diff line number Diff line change
Expand Up @@ -102,6 +102,8 @@ allowed_expr+="|^org/apache/spark/CometPlugin.class$"
allowed_expr+="|^org/apache/spark/CometDriverPlugin.*$"
allowed_expr+="|^org/apache/spark/CometExecutorPlugin.*$"
allowed_expr+="|^org/apache/spark/CometSource.*$"
# Native broadcast storage accounting needs Spark-private memory APIs.
allowed_expr+='|^org/apache/spark/CometBroadcastMemoryManager\.class$'
allowed_expr+="|^org/apache/spark/CometTaskMemoryManager.class$"
allowed_expr+="|^org/apache/spark/CometTaskMemoryManager.*$"
allowed_expr+="|^scala-collection-compat.properties$"
Expand Down
20 changes: 20 additions & 0 deletions docs/source/contributor-guide/ffi.md
Original file line number Diff line number Diff line change
Expand Up @@ -260,6 +260,26 @@ resident until the JVM closes it. See [Crossing the FFI boundary](memory_managem
| --------- | -------------------------------- | ---------------------------------------------------------- |
| All cases | Native allocates, JVM references | JVM must call `close()` to trigger native release callback |

## Lazy Union Inputs

An opt-in `CometUnionInput` is a task-scoped lazy iterator marker. Its native scan sets a
request flag, and the driving JNI thread opens the Spark iterator only after the join build
has published its completed filter. Native worker threads never open Spark Union iterators.

`openStream` borrows an address containing an `Arc<UnionFilterBundle>`. The JVM marker retains
its own owned handle before that callback returns. A bundle contains completed predicates,
leases on their accounted prepared builds, the task attempt, and the permitted native branch
root IDs. Each branch registers its handle before native planning; registration checks both
identities. The handle is released exactly once by the marker's close operation.

Branch EOF does not close its native execution while the enclosing Union can retain its Arrow
buffers. EOF remains sticky, and the enclosing owner holds the registered cleanup. After the
outer outputs and native context drop, native teardown releases the C stream and invokes the
marker's cleanup on the driving thread. That cleanup closes children in reverse order and
releases the filter leases. Only then does native teardown wait for its memory reservations.
A failed cleanup still allows other owners to close; pending Java exceptions are captured and
cleared between callbacks and the first failure is rethrown after cleanup.

## Further Reading

- [Arrow C Data Interface Specification](https://arrow.apache.org/docs/format/CDataInterface.html)
Expand Down
67 changes: 65 additions & 2 deletions docs/source/user-guide/latest/tuning/operators.md
Original file line number Diff line number Diff line change
Expand Up @@ -54,11 +54,40 @@ adaptations cannot fail. Nested column pruning still reads only the requested st
with supplied file statistics also skip reader attachment. These cases still use runtime filtering
on decoded batches.

Filters stay within the task's native plan and do not propagate across Spark exchanges or JVM/Arrow
boundaries. A shuffled hash join can still filter probe batches after shuffle, but it cannot send
By default, filters stay within the task's native plan and do not propagate across Spark exchanges
or JVM/Arrow boundaries. The separate Union setting below enables a scoped exception. A shuffled hash join can still filter probe batches after shuffle, but it cannot send
its filter back to an earlier scan stage. Compare the [runtime-filter and scan metrics](../metrics.md#hash-joins)
with the setting disabled to distinguish reduced hash-probe work from reader I/O savings.

### Runtime Filters Across UNION ALL

With both `spark.comet.exec.join.dynamicFilter.enabled=true` and
`spark.comet.exec.join.dynamicFilter.union.enabled=true`, an eligible broadcast inner join
can send its completed key filter into the native readers of a `UNION ALL` probe input.
Both settings default to false.

For example, when a small set of selected account IDs joins the union of current and archived
transactions, each transaction reader can discard row groups outside those IDs. The original
join still checks every match and preserves duplicates. Spark retains its original Union
partitions and its single broadcast exchange; Comet opens each branch lazily after the join's
build is ready.

This path accepts one direct signed integer join key and no join residual. Build columns must
be fixed-width scalars or plain UTF-8 strings; other build schemas keep ordinary execution.
Branches can rename or reorder columns and retain direct-column null checks. Computed projections, other residual
filters, limits, and intermediate joins stop reader propagation. Such branches can still be
filtered after their native output is produced. Existing per-file schema-conversion safeguards
remain in force. Missing or incomplete filters never delay opening a branch.

A filter holds a lease on the same prepared build that the original join probes. This keeps its
storage charged until every consumer finishes, including early termination. Task-attempt and
native-root checks prevent a filter from reaching unrelated executions. Nested Union owners
release child plans and Arrow streams before releasing the leases.

This draft depends on the prepared-build support in [PR #6037](https://github.com/apache/datafusion-comet/pull/6037)
and a compatible DataFusion release containing its prepared-build API. Runtime Union filtering
does not require enabling executor-wide broadcast reuse.

## Adaptive Partial Aggregation

Set `spark.comet.exec.aggregate.skipPartial.enabled=true` to let Comet bypass partial hash
Expand Down Expand Up @@ -163,3 +192,37 @@ native path with `spark.comet.expression.SortOrder.allowIncompatible=true`.

`sort_array` sorts array elements rather than ordering rows. It follows Spark's floating-point ordering as well, so it
also stays native with `spark.comet.exec.strictFloatingPoint=true`.

## Reusing Broadcast Hash Builds (Experimental)

`spark.comet.broadcast.reuse.enabled=true` allows tasks in one executor to share a prepared
native hash table for the same Spark broadcast. This avoids repeating Arrow decoding and hash-table
construction on a cache hit. It requires `org.apache.spark.CometPlugin` in `spark.plugins` and Spark
off-heap memory. The feature is disabled by default.

The implementation covers non-spilling inner broadcast joins with direct column keys,
matching key types, and fixed-width or plain UTF8 build columns. Residual join conditions remain
local to each consuming task. Dictionary, view, binary, and nested build columns remain unsupported.
Unsupported joins keep ordinary task-local execution. Tasks can share a build when they use the
same broadcast, build schema, and ordered build keys, even if their probe schemas differ.

`spark.comet.broadcast.reuse.maxMemory` defaults to `1g`. It caps native prepared-build allocations
through Spark off-heap storage memory while tasks are preparing or using a build. When the last
task releases that build, its storage charge is returned. A later task wave may prepare the same
broadcast again; a later query without a broadcast hash join inherits no idle build charge.
The first flat-schema broadcast input marked for reuse fixes the executor's cap, even if its join
later proves ineligible and creates no build. A later input with a different cap uses ordinary
execution. Cache admission never waits for another join to release its memory. When admission
fails, Comet discards the partial preparation and opens a fresh uncached input for that task.
Preparation must also admit the temporary overlap between decoded native batches and DataFusion's
compact build batch; the memory needed to prepare a build can exceed its final retained size.

This cap does not bound the existing driver broadcast representation or JVM decoder allocations.
The first task still decodes the existing broadcast format; bounded broadcast construction and
admission before JVM decoding are separate work.

Spark join metrics expose cache hits, successful preparations, admission fallbacks, preparation
rows/bytes, and preparation time. Probe metrics remain per task. A successful preparation should
normally be followed by cache hits from other tasks using that broadcast. Validate elapsed query
time and executor memory with the feature both enabled and disabled before enabling it broadly;
reuse counts alone do not establish a query speedup.
Loading
Loading