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
34 changes: 34 additions & 0 deletions docs/source/user-guide/latest/tuning/operators.md
Original file line number Diff line number Diff line change
Expand Up @@ -163,3 +163,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.
96 changes: 32 additions & 64 deletions native/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading
Loading