sunchao opened a new pull request, #6433: URL: https://github.com/apache/datafusion-comet/pull/6433
## Which issue does this PR close? No linked issue. ## Rationale for this change In a chain of joins, the final join may accept only a small part of the original input. Comet already filters rows immediately before that join, but by then earlier joins have performed lookups and produced intermediate rows that will be discarded. For example, suppose 1,000 fact rows join to a detail table with two matches per key, then to a selection table containing only key 42. Filtering only at the final join lets the detail join produce 2,000 rows. Applying the completed selection filter earlier lets it produce just the two rows for key 42. The final join still preserves every duplicate match. ## What changes are included in this PR? When `spark.comet.exec.join.dynamicFilter.enabled=true`, a completed join filter can now run on decoded batches before an eligible intermediate join's probe work: ```text Before: decoded rows → intermediate join → final join's filter → final join After: decoded rows → early filter → intermediate join → final join's filter → final join ``` The early filter follows direct probe-side columns through ordinary inner joins already eligible for runtime filtering (a single signed-integer key), column-only projections, and direct null checks. It preserves column positions when projections rename or reorder fields. Placement stops at computed expressions, other residual predicates, limits, unsupported joins, and execution boundaries. It does not introduce ancestor-filter propagation into Parquet readers; decoding and schema conversion still happen before this new consumer. Each intermediate join retains its own existing reader-filter safeguards. The final join remains responsible for matching rows. To limit redundant work on broad filters, the early consumer stops evaluating after two consecutive nonempty batches remove no rows. An inactive producer does not count toward that decision, and each execution starts fresh. If later batches become more selective, bypassing early evaluation loses an optimization but leaves results unchanged. Separate `dynamic_filter_early_*` metrics report evaluated, rejected, and bypassed rows and evaluation time on the join supplying the filter. Execution-local rewrites retain the metric owners of intermediate joins and projections, so their input/output counters continue to describe the work they perform. ## How are these changes tested? Passed 44 focused native runtime-filter tests in both debug and optimized builds, workspace Clippy across all targets with warnings denied, Rust formatting, and the root Maven formatting checks. All 15 Spark 4.1.3 join runtime-filter tests passed, including the new two-join regression with each intermediate build side and with runtime filtering off/on. The tested JNI library hash matches the library staged in the JVM test classpath. Native regressions cover reordered projections, duplicate and null keys, original metric ownership, repeat execution, cancellation, limits, computed-expression errors, intermediate build-side columns, join residuals, and adaptation with inactive and empty batches. A Spark regression compares two-join results and intermediate row counts with runtime filtering off/on. A follow-up commit will add a reproducible native component benchmark and matched public-base/head measurements, including an all-matching domain. This draft does not make a query-speed claim. -- 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]
