sunchao opened a new pull request, #6437:
URL: https://github.com/apache/datafusion-comet/pull/6437

   ## Which issue does this PR close?
   
   No linked issue. This is a stacked follow-up to #6037 and remains a draft 
until the prepared-build API is available in a compatible DataFusion release. 
It includes that prerequisite's commits plus their merge with current main; the 
Union-specific change is isolated in [the final 
commit](https://github.com/sunchao/arrow-datafusion-comet/commit/8a6fc78ea6f550b7c9449e65beff884be15718fd).
   
   ## 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:
   
   ```sql
   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.
   
   ```mermaid
   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
   ```
   
   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. The released dependency pins 
remain unchanged, so a build with those unmodified dependencies does not yet 
compile this draft. The native results below use the same [public DataFusion 55 
companion 
port](https://github.com/sunchao/arrow-datafusion/commit/c9a142d834b6b47c2ac15524171bbf84f3022f73)
 as that prerequisite, applied through validation-only Cargo patches.
   
   ## How are these changes tested?
   
   Validated head `8a6fc78ea6f550b7c9449e65beff884be15718fd` with the public 
DataFusion companion identified above. Native execution tests passed **306 
tests**, with four existing HDFS-dependent tests ignored. Workspace Clippy 
across all targets with warnings denied, the default-feature native build, Rust 
formatting, and the root Maven Spotless/Scalastyle checks passed.
   
   All **nine focused Spark tests** passed on Spark 4.1.3 / Scala 2.13.17 / JDK 
21. The JNI library built from this source matched the SHA-256 of the library 
staged in the JVM test resources. The tests cover both join build sides, AQE on 
and off, renamed/reordered columns, duplicate/null semantics, branch-specific 
reader pruning, fallible projections, limits, early termination, cancellation 
followed by a fresh query, sequential bucketed readers, and nested Union 
ownership and forwarding.
   
   In the two-reader fixture, enabling transport reduced measured scan bytes 
from 174,890 to 3,498 while preserving the same six result rows, including 
duplicate matches. Both build-side cases produced those results. This is a 
feature-toggle regression check on one JNI build, not an elapsed-time benchmark 
or a whole-query speedup claim.
   
   Native regressions separately cover branch/root and task-attempt 
authorization, live filter remapping and completion, incompatible key metadata, 
chunking of large build-copy handoffs, and reader preparation with renamed 
columns and stable Spark metrics.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to