Abhisheklearn12 opened a new issue, #25976: URL: https://github.com/apache/datafusion/issues/25976
### Is your feature request related to a problem or challenge? #25771 evaluates `AND` chains as one n-ary conjunction. Before each conjunct it filters the batch down to the rows that are not yet `false`, but only when at most `PRE_SELECTION_THRESHOLD` (20%) of the batch's rows are left and the conjunct is not cheap and infallible. Cheap conjuncts, like `x < 70`, always run on every remaining row, because filtering the batch and scattering the result back usually costs more than a cheap comparison. That rule loses when very few rows survive and they sit in a few long runs, since filtering is cheap then. In review, @jayzhan211 reproduced the `selectivity_q21` loss from `predicate_eval`: on a 3-column `Int32` batch of 8192 rows, a clustered `a < t` prefix with ~1% survivors followed by cheap `x < 70` conjuncts is 1.14x slower with one cheap conjunct (2.03 → 2.32 µs) and 1.24x slower with two (2.81 → 3.48 µs). At 5 to 20% survivors, or with randomly placed survivors, #25771 is 0.26x to 0.99x, so simply lowering the threshold would hurt the random case. ### Describe the solution you'd like Before a cheap conjunct, also filter when the undecided rows form a few long runs. The slice count from `SlicesIterator` is a natural signal: `pre_selection_scatter` walks the mask slice by slice, so its cost depends on the number of runs, not just the survivor count. Open questions: - How to count runs cheaply, without slowing down the common case where nothing gets filtered. - What cutoff to use, chosen from measurements across survivor fraction, run count, batch width and the number of cheap conjuncts left. Filtering before an infallible conjunct can only shrink the rows it sees, so this changes performance only, not results or errors. ### Describe alternatives you've considered - A lower survivor threshold for cheap conjuncts: hurts randomly placed survivors, as above. - A faster `pre_selection_scatter`, which currently appends values one at a time. That would make filtering cheaper everywhere and may shift this trade-off, so it's worth doing separately. - Keeping the current rule: the loss is small in absolute terms (under 1 µs per 8192-row batch in the case above). ### Additional context - Discussion: #25771 (review by @jayzhan211) and #25035. - The `conjunction` benchmarks in `binary_op` (#25870) all use randomly placed survivors, including the selective-first case, so a clustered-survivor case needs to be added to measure this. I'm happy to work on this once #25771 merges. -- 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]
