Skip to content

perf: reuse prepared broadcast builds across executor tasks - #6037

Draft
sunchao wants to merge 5 commits into
apache:mainfrom
sunchao:codex/oss-broadcast-build-reuse-6013
Draft

sunchao wants to merge 5 commits into
apache:mainfrom
sunchao:codex/oss-broadcast-build-reuse-6013

Conversation

@sunchao

@sunchao sunchao commented Sep 19, 2026 •

Copy link
Copy Markdown
Member

Which issue does this PR close?

Part of #6013. This adds reuse for eligible inner joins while tasks on an executor are actively probing the same broadcast.

Rationale for this change

Spark broadcasts the build relation once, but Comet still decodes that relation and constructs a native hash table separately for each probe task. Sharing the broadcast bytes therefore does not share the native preparation work.

For example, consider eight concurrent tasks joining different partitions of events to the same broadcast users table. Today, the executor can decode users eight times and build eight copies of its hash table. Those tasks could instead probe one immutable build. Switching to DataFusion's CollectLeft alone does not solve this: it shares preparation within a single join plan, whereas each Comet task has its own plan.

What changes are included in this PR?

This PR lets compatible tasks on an executor prepare a broadcast once and share it while they run. The first task opens and decodes the broadcast; other tasks using the same broadcast and build keys borrow the prepared rows and hash table. Each task still runs its own probe input and join condition. A cache hit can skip both decoding and hash-table construction.

The shared build belongs to the executor's storage memory pool, so finishing the task that created it does not invalidate another task's join. Active tasks keep it alive, and the last task releases its memory. The lookup retains only weak references: a later wave of tasks may build the broadcast again. If preparation cannot obtain memory, affected tasks use the ordinary join path with fresh broadcast streams.

Reuse is disabled by default and enabled with spark.comet.broadcast.reuse.enabled=true. It requires CometPlugin and Spark off-heap memory; spark.comet.broadcast.reuse.maxMemory defaults to 1g per executor and caps the native memory used to prepare and retain shared builds. Existing Spark broadcast and JVM decoder buffers are outside this cap.

The initial scope is inner joins with matching direct-column keys and fixed-width or plain UTF-8 build columns. Other joins continue through the existing path. The tuning guide explains the memory scope and metrics for evaluating reuse.

This draft now pins the immutable Arrow 59-compatible prepared-build companion, so ordinary native builds and CI can exercise the feature. It is two commits ahead of official DataFusion 55.1.0, changing only eight prepared-join implementation/test and serialization-test files. All 32 DataFusion crates use that same source at version 55.1.0 to avoid mixing incompatible Rust types; Arrow and all other dependency versions remain unchanged. Remove this temporary pin when Comet adopts an official compatible release containing the API. The PR remains a draft.

How are these changes tested?

Current head 2b94d426d107187d4ade50c3fd88d557396b3bf9 is rebased on Apache main 0ac4dadae70838d8675bfd846c2d1c7617c25c78.

  • All 12 selected native broadcast/cache/filter tests pass using the committed dependency manifest and lockfile, without validation-only dependency overrides.
  • All five selected Spark 3.5 runtime tests pass with the native library built from this head: four prepared-broadcast reuse cases covering duplicate/null inputs, concurrent tasks, composite string keys/residuals and admission fallback, plus the lazy-input lifecycle test.
  • Full root-reactor Spark 3.5 compilation with strict Scala warnings, Spotless, Scalastyle, Rust formatting, and whitespace checks pass.
  • Lockfile inspection confirms one source for every DataFusion crate, no registry DataFusion duplicates, and no non-DataFusion dependency changes from the temporary pin.

Broader current-head CI is pending. No current-head benchmark was run. The PR remains a draft while using the temporary companion dependency.

@github-actions github-actions Bot added enhancement New feature or request performance area:scan Parquet scan / data reading area:memory Memory pools, reservations, OOM handling area:joins Join operators and dynamic filter pushdown labels Sep 19, 2026
@sunchao sunchao added the run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue label Sep 19, 2026
adriangb pushed a commit to pydantic/datafusion that referenced this pull request Sep 26, 2026
## Which issue does this PR close?

Part of [Comet
apache#6013](apache/datafusion-comet#6013). [Comet
apache#6037](apache/datafusion-comet#6037) uses this
API. The shared row-index accounting fix in
[apache#25508](apache#25508) is included in
the base branch.

## Rationale for this change

Independent joins against the same build snapshot repeat preparation and
retain separate hash tables. `CollectLeft` shares a build among
partitions of one join, but cannot share it across independently created
plans.

For example, eight concurrent requests joining events to the same users
snapshot can share one users lookup while probing their own event
batches. Preparing that lookup once avoids seven redundant builds.
Applications can replace the snapshot while existing consumers finish
with the old one. This supports both Comet broadcast tasks and other
DataFusion applications.

Comet creates independent native joins for its executor tasks, so their
existing builds can overlap across task slots. The intended benefit is
sharing one retained build batch and hash table and avoiding repeated
materialization and hashing work. Removing repeated builds does not
imply a proportional latency improvement: cold consumers wait for
preparation, then probe independently. Latency depends on scheduling,
contention, build size, coordination, and whether a prepared build is
already retained. This PR does not establish end-to-end Comet latency,
process CPU-time, or RSS savings.

## What changes are included in this PR?

`HashJoinExec::prepare_build` returns an immutable
`PreparedHashJoinBuild` that compatible joins attach explicitly:

```rust
let prepared = first_join
    .prepare_build(users_snapshot, durable_pool, config)
    .await?;
let first = first_join.builder()
    .with_prepared_build(Arc::clone(&prepared))
    .build()?;
let second = second_join.builder()
    .with_prepared_build(prepared)
    .build()?;
```

Consumers share build rows, lookup storage, and the retained memory
reservation. Probe progress, residual predicates, and dynamic filters
remain independent. Preparation charges retained data and temporary
working memory; errors and cancellation release unfinished work.
Concatenation alignment is charged once per output buffer, so small
input batches do not multiply that reservation.

Supported builds use non-spilling `CollectLeft` inner joins,
direct-column keys with matching types, and fixed-width or UTF-8 build
columns. Attachment and execution check schema, keys, and null equality,
replaces the unused build subtree with a schema placeholder, and
preserves the handle through supported resets and projections. Prepared
state is process-local and cannot be serialized.

The caller owns snapshot identity, concurrent preparation, and eviction.
Input buffers must survive producer cleanup, and the supplied pool must
cover every consumer. Attach after physical optimization, supply fresh
probe plans and dynamic-filter expressions, and retain a build handle
while any associated filter outlives its plan. Compatibility checks
cannot verify snapshot contents.

## What is the testing strategy for this PR?

Fifteen focused tests cover equivalence with ordinary joins, both
null-equality modes, empty and duplicate keys, batch boundaries,
concurrent consumers with independent filters and predicates, plan
rewrites, public key changes after attachment, admission failures,
small-batch concatenation under a fixed memory limit, cancellation, and
buffer/reservation lifetimes. The latest regressions exercise
under-admission without shrinking the reservation, checked copy/scratch
arithmetic, and IN-list buffer sharing across one/two keys and one/two
input batches. Existing serialization coverage rejects prepared state;
the API example is a doctest. The benchmark checks outputs across 24
cases.

Prepared copy and scratch sizing retain checked arithmetic. A concat
output that exceeds its admitted bound returns an error in every build
profile before shrinking the reservation. Numeric IN-list membership
shares the retained batch's key buffers; charging its full array size
again would duplicate that payload charge. General array/schema metadata
and consumer-local filter expressions are outside the buffer
reservation.

Validation of
[8a964a8](apache@8a964a8)
passed: all 15 focused prepared-build tests in the normal `ci` profile;
`cargo fmt --all`; full workspace Clippy with all targets/features and
`-D warnings`; the complete repository lint suite, including Rust and
HTML documentation builds; and the extended workspace suite (12,159 Rust
tests and all 524 SQL logic files, with 8 Rust tests ignored). All 15
prepared-build tests also passed in the optimized `release-nonlto`
profile.

The public-key-mutation regression was verified to fail before restoring
execution-time validation: the ordinary join returned its two expected
rows while the stale prepared lookup returned none. With the guard
restored, the mutated descriptor is rejected with a planning error.

Local validation used Rust 1.98.1 and temporary source overrides for the
official Arrow 60.0.0 (`ef1fa157`), object_store 0.14.2 (`279572ea`),
and sqlparser 0.63.0 (`85b1a6f2`) releases. Dependency overrides and
generated lockfile changes are not included.

The rebased head has not been benchmarked.

```bash
cargo bench -p datafusion-physical-plan --features test_utils --profile release-nonlto --bench prepared_hash_join
```

## Are there any user-facing changes?

Opt-in Rust APIs and an `EXPLAIN` marker for prepared builds. Ordinary
SQL planning and join selection are unchanged.
@sunchao
sunchao force-pushed the codex/oss-broadcast-build-reuse-6013 branch from 40bdc29 to 7a930fa Compare October 4, 2026 17:42
@sunchao

sunchao commented Oct 4, 2026

Copy link
Copy Markdown
Member Author

Rebased in 7a930fa205fe0076de859a1fd950da991c4d94ee. Adapted to main's ArrowArrayStreamReader, kept the upstream BatchProducer shutdown and plugin lifecycle cleanup, and registered the broadcast lifecycle suite in both CI matrices.

Native core/test compilation passes with the same validation-only public DataFusion companion (c9a142d834b6b47c2ac15524171bbf84f3022f73), and full Spark 4.1 reactor test compilation, Rust formatting, Spotless, suite registration and whitespace checks pass. Dependency pins and the lockfile remain unchanged.

The existing draft blocker remains: Comet's released DataFusion55.1 dependency does not contain the prepared-build APIs, even though upstream #25491 is merged. Runtime tests/benchmarks were not rerun on this rebased head, and ordinary native CI cannot pass until that dependency prerequisite is available.

@sunchao

sunchao commented Oct 4, 2026

Copy link
Copy Markdown
Member Author

Cleared the Scala 2.12 strict-warning failures in 137241bd8b5e3f8500480204f943b640ff1594c1: renamed the broadcast RDD constructor parameter to avoid shadowing RDD.name and made the join fixture numeric conversion explicit. The complete Spark 3.5 reactor now passes test-compile -Pstrict-warnings -DskipTests, including Spotless/Scalastyle. The documented DataFusion prepared-build API dependency remains the native CI blocker for this draft.

@sunchao
sunchao force-pushed the codex/oss-broadcast-build-reuse-6013 branch from 137241b to 2b94d42 Compare October 5, 2026 21:42
@sunchao

sunchao commented Oct 5, 2026

Copy link
Copy Markdown
Member Author

Rebased and fixed the native dependency blocker in 2b94d426d107187d4ade50c3fd88d557396b3bf9. The committed manifest now pins the existing public Arrow 59-compatible companion c9a142d834b6b47c2ac15524171bbf84f3022f73 across the DataFusion 55.1 crate family. That companion is only two commits ahead of official 55.1, with eight relevant files changed. No Arrow or other dependency versions change, and all DataFusion types come from one source. The temporary pin can be removed after an official compatible release is adopted.

The native build and all 12 focused native tests pass; all five Spark 3.5 broadcast/lifecycle runtime tests pass with the matching library. Strict full-reactor Scala compilation and formatting/style pass. Broader CI is pending, and draft status is preserved.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:joins Join operators and dynamic filter pushdown area:memory Memory pools, reservations, OOM handling area:scan Parquet scan / data reading enhancement New feature or request performance run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant