neilconway opened a new pull request, #25761:
URL: https://github.com/apache/datafusion/pull/25761

   ## Which issue does this PR close?
   
   - Related to #15177
   
   ## Rationale for this change
   
   If a highly selective `FilterExec` feeds a Top-K or ungrouped min/max node, 
there is an unfortunate interaction between batching and dynamic filtering: the 
selective filter means that it may take many input rows before a batch of 
output from the `FilterExec` is ready to be consumed by the TopK/aggregate 
node. That delay means that dynamic filtering may be much less effective.
   
   This pattern can be observed in ClickBench Q23:
   
   ```
   SELECT *
   FROM hits
   WHERE "URL" LIKE '%google%'
   ORDER BY "EventTime"
   LIMIT 10;
   ```
   
   But it should be reasonably common in practice.
   
   This PR aims to address this problem by arranging for the `FilterExec` to do 
a one-time early flush of "startup rows", for plan shapes like those described 
above. This can allow dynamic filtering to be applied sooner, potentially 
improving query performance. The extra small output batch and earlier 
dynamic-filter processing can add overhead when the initial bound eliminates 
little work.
   
   On an M4 Max machine with 16 partitions, this optimization reduced 
ClickBench Q23 by 33% (2,351 ms -> 1,575 ms).
   
   ## What changes are included in this PR?
   
   * During physical plan optimization, annotate eligible filters with 
`startup_rows`. Specifically, we look for a path from Top-K or ungrouped 
min/max, to a `FilterExec`, and finally to a scan with a matching dynamic 
filter. A round-robin partition is allowed below the filter (to match a common 
plan shape), but other plan nodes (e.g., hash partitioning) disable the 
optimization.
   * The number of `startup_rows` is determined by the consumer: ungrouped 
`min`/`max` ask for 1 row, Top-K asks for k rows (larger values might be worth 
considering but I think these values are a reasonable starting point).
   * Teach filters to request one early flush after accepting `startup_rows` 
qualifying rows across input batches. On reaching the threshold, flush any 
buffered output, then continue with normal batching.
   * Update protobuf and JSON serialization for `startup_rows`
   * Add tests
   
   ## What is the testing strategy for this PR?
   
   Existing tests pass; new tests added. Added tests cover one-time flushing, 
preserved values and fetch limits, null predicates, thresholds spanning 
multiple batches, and preservation of natural batch boundaries. Optimizer tests 
cover eligible and excluded plan shapes, round-robin placement, unsupported 
sources, disabled dynamic filtering, and MIN/MAX eligibility.
   
   ## Are there any user-facing changes?
   
   `EXPLAIN VERBOSE` displays `startup_rows` for annotated filters; default 
EXPLAIN output is unchanged.
   


-- 
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