Skip to content

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

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

sunchao wants to merge 4 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 remains a draft until DataFusion's prepared-build API is available in a compatible release. Dependency pins remain unchanged in this PR.

How are these changes tested?

Current head 137241bd8b5e3f8500480204f943b640ff1594c1 is rebased on Apache main ac9ae94d057013095fdf1b70e2d4bc5b24a58d90.

  • Native core and test compilation passed with the datafusion-physical-plan source from the same public companion commit c9a142d834b6b47c2ac15524171bbf84f3022f73, applied through a validation-only Cargo patch. The committed dependency manifest and lockfile remain unchanged.
  • The full Spark 4.1 Maven reactor passed test-compile, including the Scala regression tests.
  • Rust formatting, Spotless, the workflow suite-registration check, and base-relative whitespace checks passed. The broadcast lifecycle suite is registered in both Linux and macOS CI matrices.

Runtime tests have not been rerun locally on this rebased head, and no new benchmark was run. The unmodified released DataFusion 55.1.0 dependency still lacks PreparedHashJoinBuild, prepare_build and with_prepared_build, so native CI remains blocked on the existing dependency prerequisite. Upstream DataFusion #25491 has merged, but those APIs are not in Comet's pinned release. This PR remains a draft.

The CI follow-up clears Spark 3.5/Scala 2.12 strict warnings in the broadcast RDD and duplicate/null join fixture. The full Spark 3.5 reactor test-compile -Pstrict-warnings -DskipTests passed locally.

@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.

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