Repository navigation
Conversation
|
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. |
25d3a5a to
75dcf5e
Compare
|
Pushed |
75dcf5e to
234ebbb
Compare
|
Rebased on current main and the updated #6037 prerequisite in 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. |
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 ALLinput, 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:
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 --> GSpark keeps the original Union partitions, broadcast exchange and join. This preserves
UNION ALLduplicates 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.enableddefaults tofalseand also requiresspark.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
234ebbbc0caa2820f291579650eed3f58dee1116is rebased on Apache main0ac4dadae70838d8675bfd846c2d1c7617c25c78, including the updated #6037 prerequisite.Broader current-head CI is pending. No current-head benchmark was run. The PR remains a draft while using the temporary companion dependency.