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]
