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]

Reply via email to