Conversation
## 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.
40bdc29 to
7a930fa
Compare
|
Rebased in Native core/test compilation passes with the same validation-only public DataFusion companion ( 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. |
|
Cleared the Scala 2.12 strict-warning failures in |
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
eventsto the same broadcastuserstable. Today, the executor can decodeuserseight times and build eight copies of its hash table. Those tasks could instead probe one immutable build. Switching to DataFusion'sCollectLeftalone 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.maxMemorydefaults to1gper 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
137241bd8b5e3f8500480204f943b640ff1594c1is rebased on Apache mainac9ae94d057013095fdf1b70e2d4bc5b24a58d90.datafusion-physical-plansource from the same public companion commitc9a142d834b6b47c2ac15524171bbf84f3022f73, applied through a validation-only Cargo patch. The committed dependency manifest and lockfile remain unchanged.test-compile, including the Scala regression tests.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_buildandwith_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 -DskipTestspassed locally.