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]

Reply via email to