adriangb opened a new issue, #25778: URL: https://github.com/apache/datafusion/issues/25778
Part of https://github.com/apache/datafusion/issues/22883. Follow-up to the optional filter gate in https://github.com/apache/datafusion/pull/25674. ### Problem The gate pauses an optional filter when `eval_ns > removed_rows * saving_ns_per_row`. Today, `saving_ns_per_row` is the work that the producer of the filter (for example a hash join) measures for each probe row without a match. If the producer has no measurement, the gate uses the tunable estimate `optional_filter_min_saving_ns_per_row` (default 20 ns). A removed row can save much more than the producer work, or much less. The value depends on how far the row travels downstream: | Plan | What a removed row saves | Real value (TPC-DS SF1) | |---|---|---| | `scan(store_sales) -> HashJoin(date_dim)` (Q10, Q31) | One hash and one lookup | 1.3 ns | | `scan(lineitem) -> ... -> join -> aggregate` (TPC-H Q17) | Work of all operators up to the join | more than 12 ns | | `scan(catalog_sales) -> join(inventory)` (Q72) | Output rows of a join with many matches | the query goes from 14.9 s to 0.28 s | The producer cannot measure this: 1. **Starvation.** A selective filter removes the rows before the producer sees them. In Q10 the date join saw 839 probe rows in total, fewer than `MIN_OBSERVED_ROWS`, so the gate used the estimate. 2. **Scope.** The producer measures only its own work. It does not see the operators between the scan and the producer. ### Evidence Total CPU (EXPLAIN ANALYZE, ms, minimum of 3 runs). The columns show the leaf with the estimate set to 20, 5 and 0 ns: | Query | pruning_only | est 20 | est 5 | est 0 | |---|---|---|---|---| | DS Q72 | 26655 | 279 | 396 | **14903** | | DS Q54 | 60 | 21 | 58 | 61 | | H Q17 | 477 | 295 | 471 | 464 | | H Q20 | 159 | 110 | 160 | 151 | | CB Q23 | 4891 | 1791 | 1641 | 2573 | | DS Q10 | 68 | 75 | 80 | 75 | The wins depend on the estimate. The estimate is also wrong where the regressions are: | Filter | Saving that the gate used | Real work per removed row | Eval | |---|---|---|---| | Q10 `ss_sold_date_sk` | 20 (estimate; producer starved) | 1.31 | 2.5 | | Q10 `cd_demo_sk` | 20 (estimate) | 4.55 | 7.7 | | Q69 `ss_sold_date_sk` | 20 (estimate) | 1.26 | 2.8 | ### Prototype Branch: https://github.com/pydantic/datafusion/tree/optional-filter-downstream-saving-prototype | Part | Design | |---|---| | Measurement | For each input batch, the scan records the filter state (on or off), the input rows, the rows that the filter removed, and the time from the return of the batch to the next poll. | | Estimate | `saving = (median(off) - median(on)) / fraction removed while on`, in ns for each input row | | Probe | A site without an "off" sample pauses the filter once (`initial_pause_batches`), shared by all partitions and files | | Floor | `max(downstream, measured producer work)`, because the time stops at an exchange | | Constants | None new | Results (total CPU, ms, minimum of 3): | Query | leaf est 20 | est 20 | est 5 | est 0 | |---|---|---|---|---| | DS Q72 | 279 | 617 | 540 | 718 | | H Q17 | 295 | 347 | 304 | 346 | | H Q20 | 110 | 124 | 119 | 121 | | CB Q23 | 1791 | 1621 | 1611 | 1881 | | DS Q54 | 21 | 57 | 57 | 57 | | DS Q79 | 117 | 110 | 148 | 147 | | DS Q10 | 75 | 76 | 75 | 78 | For H Q17, H Q20, CB Q23 and DS Q72, the decisions no longer depend on the estimate. The prototype is not ready because of the failure modes below. ### Failure modes | Mode | Example | Numbers | |---|---|---| | The probe is not free | DS Q72: the probe lets `catalog_sales` rows into the join with `inventory` | 279 → 617 ms | | The probe is not free | DS Q54: a probe on a filter that removes all rows of `customer_address` puts 32,770 rows into a build side, so `store_sales` is scanned and not skipped | 21 → 57 ms | | The sample depends on the other filters of the scan | DS Q10: `ss_customer_sk` was sampled while the date filter was on (6.5 ns). The real value after the date filter paused is about 1.5 ns, against an eval of 3.1 ns | filter kept by mistake | | The sample depends on the other filters of the scan | DS Q79: the `ss_sold_date_sk` sample is 2.3 ns with est 20 and 34.8 ns with est 0 | scan 93 → 133 ms | | The sample depends on the other filters of the scan | DS Q31: date filters are paused correctly, but other filters change their state | 169 → 178 ms | | Work not visible | The work after an exchange (`RepartitionExec`) or into a hash join build side is not in the time | DS Q54 build side | ### Open design questions 1. Continuous sampling at a low rate, instead of one probe for each site. What rate, from what measurement? 2. Take a new sample when another filter of the same scan changes its verdict. 3. A bound on the cost of a probe, for example from the removed fraction and the producer work. It must come from a measurement, not from a new constant. 4. How to measure the work after an exchange or into a build side (for example, the producer could report the rows that it sees for each input row of the scan). 🤖 Generated with [Claude Code](https://claude.com/claude-code) -- 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]
