Skip to content

feat: carry completed join filters across Union inputs - #6437

Draft
sunchao wants to merge 7 commits into
apache:mainfrom
sunchao:codex/oss-union-runtime-filters
Draft

sunchao wants to merge 7 commits into
apache:mainfrom
sunchao:codex/oss-union-runtime-filters

Conversation

@sunchao

@sunchao sunchao commented Sep 29, 2026 •

Copy link
Copy Markdown
Member

Which issue does this PR close?

No linked issue. This is a stacked follow-up to #6037 and includes its rebased prerequisite commits. The final two commits contain the Union implementation and strict Scala warning fixes. It remains a draft while using the temporary compatible DataFusion companion described below.

Rationale for this change

A selective join can often tell a reader which rows will never match. That information currently stops at a Spark UNION ALL input, even when each branch uses a native Parquet reader. The readers can therefore load data that the join immediately rejects.

For example, consider finding transactions for a small set of selected accounts across current and archived data:

SELECT /*+ BROADCAST(a) */ t.account_id, t.amount
FROM (
  SELECT account_id, amount FROM current_transactions
  UNION ALL
  SELECT account_id, amount FROM archived_transactions
) t
JOIN selected_accounts a ON t.account_id = a.account_id;

After building selected_accounts, the join already knows its completed set of matching keys. Sending that filter to both transaction readers lets them skip row groups that cannot contain a match. The saving depends on the data and file layout; this PR makes no query-speedup claim.

The difficulty is the execution boundary. Spark owns the Union iterator, and its native branch plans are created separately from the native join above it. Simply attaching a filter to the join's native probe cannot reach those readers. Opening the branches eagerly also starts them before the build filter is ready.

What changes are included in this PR?

The native join now opens an eligible Union input lazily, after preparing its build. It passes the completed filter through that input to the separately created branch plans. A branch can attach the filter to its reader, forward it to another eligible lazy Union, or apply it after producing its native output when reader propagation is unsafe. A missing or incomplete filter lets the branch proceed immediately.

flowchart TD
    A[Selected accounts] --> B[Prepare the original join build]
    B --> C[Publish a completed filter and build lease]
    C --> D[Open Spark Union lazily]
    D --> E[Current-data native reader]
    D --> F[Archive native reader]
    E --> G[Original join verifies matches]
    F --> G
    B --> G
Loading

Spark keeps the original Union partitions, broadcast exchange and join. This preserves UNION ALL duplicates and avoids duplicating the join in every branch. The same prepared build supplies the filter and the final hash probe. A lease keeps the build storage alive and charged while any branch or enclosing plan can still reference it. Filter handles are scoped to one task attempt and explicitly authorized native plan roots, so unrelated plans cannot consume them.

The setting spark.comet.exec.join.dynamicFilter.union.enabled defaults to false and also requires spark.comet.exec.join.dynamicFilter.enabled=true. The first version supports broadcast inner joins with one direct signed integer key and no join residual, with fixed-width or plain UTF-8 build columns. Branches may rename or reorder columns or retain direct-column null checks. Computed projections, other residuals, limits and intermediate joins stop reader propagation. Existing per-file schema-conversion safeguards remain intact. For example, a branch containing a failing cast must still evaluate that cast before the transported filter can reject its output.

Lazy execution also changes resource ownership. A branch reaching EOF can still have buffers retained by its parent. Its Union owner now keeps the branch plan until the enclosing native plan releases those buffers, then closes child plans and filter leases before waiting for native memory to return. This ordering also covers early termination and branch failures.

The implementation uses the prepared-build support from #6037, but does not require enabling executor-wide broadcast reuse. This draft inherits the immutable Arrow 59-compatible prepared-build companion from that prerequisite, so ordinary native builds and CI can exercise the feature. The companion is two commits ahead of official DataFusion 55.1.0 and changes only eight prepared-join implementation/test and serialization-test files. All 32 DataFusion crates use that same source at version 55.1.0; Arrow and all other dependency versions remain unchanged. Remove this temporary pin when Comet adopts an official compatible release containing the API.

How are these changes tested?

Current head 234ebbbc0caa2820f291579650eed3f58dee1116 is rebased on Apache main 0ac4dadae70838d8675bfd846c2d1c7617c25c78, including the updated #6037 prerequisite.

  • All 15 selected native tests pass using the committed dependency manifest and lockfile: three Union filter-handoff/authorization tests and 12 broadcast/cache/filter tests.
  • All ten selected Spark 4.1 runtime tests pass with the native library built from this head. They cover both build-side/AQE configurations, early termination, cancellation, fallible branches, branch limits, sequential bucketed readers, nested Unions, intermediate joins, and lazy broadcast task-thread ownership.
  • In the two basic Union cases, the selected reader-byte metric decreases from 174890 to 3498 while preserving query results. This is focused fixture evidence, not a general performance benchmark.
  • Full root-reactor Spark 3.5 compilation with strict Scala warnings, the Spark 4.1 runtime reactor, Spotless, Scalastyle, Rust formatting, and whitespace checks pass.
  • Independent 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 the enhancement New feature or request label Sep 29, 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 29, 2026
@github-actions github-actions Bot added area:scan Parquet scan / data reading area:memory Memory pools, reservations, OOM handling area:joins Join operators and dynamic filter pushdown and removed run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue labels Sep 29, 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 Oct 4, 2026
@sunchao

sunchao commented Oct 4, 2026

Copy link
Copy Markdown
Member Author

Rebased in 25d3a5a896dd9072eba1fd14726e678e53472583. Rebased the Union change onto the rebased #6037 prerequisite. Kept main's TopK support, operators-crate imports and release of never-consumed Arrow streams while preserving Union-owned cleanup.

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 force-pushed the codex/oss-union-runtime-filters branch 2 times, most recently from 25d3a5a to 75dcf5e Compare October 4, 2026 18:00
@sunchao

sunchao commented Oct 4, 2026

Copy link
Copy Markdown
Member Author

Pushed 75dcf5e777148f23e1ec2f3b2ecedd9162b50b67 on top of the latest #6037 follow-up. This clears the Scala 2.12 strict warnings (RDD-name shadowing, discarded completion-listener return, and implicit integer widening) and the redundant test interpolation reported by Scalafix. The complete Spark 3.5 reactor passed test-compile -Pstrict-warnings -DskipTests, including Spotless/Scalastyle. The validated Scala files were verified byte-for-byte after refreshing the stack. The documented DataFusion prepared-build API dependency remains the native CI blocker for this draft.

@sunchao
sunchao force-pushed the codex/oss-union-runtime-filters branch from 75dcf5e to 234ebbb Compare October 5, 2026 21:49
@sunchao

sunchao commented Oct 5, 2026

Copy link
Copy Markdown
Member Author

Rebased on current main and the updated #6037 prerequisite in 234ebbbc0caa2820f291579650eed3f58dee1116. The native dependency blocker is fixed by the committed immutable Arrow 59-compatible companion pin c9a142d834b6b47c2ac15524171bbf84f3022f73; all DataFusion crates remain at 55.1.0 with a single source, and all other dependency versions are unchanged. Remove this temporary pin after adopting an official compatible release.

All 15 focused native tests and all ten Spark 4.1 Union/lifecycle runtime tests pass with the matching current-head library. The two basic Union fixtures preserve results while reader bytes decrease from 174890 to 3498. The full strict Spark 3.5 compilation and formatting/style checks also 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 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