sunchao opened a new pull request, #6431: URL: https://github.com/apache/datafusion-comet/pull/6431
## Which issue does this PR close? No linked issue. ## Rationale for this change A join can often rule out most probe-side rows once its build side is complete. Comet already uses that information to prune eligible Parquet readers, but an intervening predicate such as `quantity > 0` stops propagation. The reader may therefore scan many rows that cannot join, even though the predicate is deterministic. Crossing that predicate is an execution-policy choice. Consider a residual `10 / key > 0`: a row with `key = 0` can raise a division-by-zero error even when the build contains only key 3. If reader pruning removes that row first, its expression is never evaluated. Determinism alone does not make that change invisible, so broader propagation should require an explicit opt-in. ## What changes are included in this PR? This PR adds `spark.comet.exec.join.dynamicFilter.allowDeterministicFilterPushdown`, defaulting to **false**. With both this setting and `spark.comet.exec.join.dynamicFilter.enabled` enabled, a join runtime filter may cross a deterministic probe-side residual and reach the existing Parquet reader path. For example, a build containing item IDs 17 and 42 can let the reader skip a row group containing only item 99 before evaluating its quantity predicate. The residual stays in the plan and still checks every retained row. A row that passes the runtime filter can still fail the residual, and errors on retained rows still propagate. The new setting explicitly permits skipping expression errors on rows eliminated earlier. With the setting disabled, propagation remains limited to direct-column `IS NOT NULL` checks combined with `AND`. Spark certifies that the complete residual is deterministic and sends that permission with the filter. Native planning uses the metadata from the actual probe side, including build-right joins, and preserves it when rebuilding the filter for an execution. Nondeterministic predicates, limits, computed projections, and execution boundaries remain barriers. Existing per-file Parquet schema-conversion safeguards are preserved, including when the new setting is enabled. ## How are these changes tested? - All 15 `CometJoinSuite` tests selected by `join dynamic filter` pass on Spark 4.1.3. The new regression covers both build sides, compares the option off/on, checks residual rejection of a key present in the build, and verifies reader attachment, row-group pruning, scan bytes, and metric ownership. The suite also checks seeded-random evaluation and Parquet schema errors with the option enabled. - All 40 native tests selected by `filter` pass. A real Parquet regression verifies that opt-in pruning skips division by zero only for an eliminated row; the default path and a retained zero-key row still fail. Additional tests cover metadata selection and preservation during plan rebuilding. - `cargo clippy --locked --offline --all-targets --workspace -- -D warnings`, the native build, Maven Spotless, Markdown formatting, and `git diff --check` pass. The Spark run used the newly built JNI library. The scan-byte assertion is regression coverage for the fixture, not an end-to-end performance claim. Broader Spark SQL coverage is requested through the `run-spark-4.1-tests` CI label. -- 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]
