asolimando commented on PR #21122:
URL: https://github.com/apache/datafusion/pull/21122#issuecomment-5955682180

   @2010YOUY01 resuming this now that #21815, #23051 and #23651 have landed. 
You asked for concrete examples before the design, so here they are first: 
seven queries over one dataset, each with the estimate it needs, the operator 
that consumes it, and the design point it tests. The decisions they lead to 
follow.
   
   ### Examples
   
   Dataset: `sales`, 10,000 rows, values uniform within their range, `c` 
strongly correlated with `a`.
   
   | column | type | min | max | distinct | nulls |
   |---|---|---|---|---|---|
   | `a` | INT | 0 | 99 | 100 | 1,000 |
   | `b` | INT | 200 | 300 | 10 | 0 |
   | `c` | INT | 0 | 99 | 100 | 0 |
   | `ts` | TIMESTAMP | 2025-01-01 | 2025-12-31 | ~10,000 | 0 |
   | `s` | VARCHAR | | | 5,000 | 0 |
   
   **1. `SELECT a + b AS k, count(*) FROM sales GROUP BY k`.** The range of `a 
+ b` is [200, 399], so there are at most 200 distinct values plus one NULL 
group: at most 201 groups, where main estimates 10,000. Point: one propagated 
unit. The distinct count follows from the range, without a distribution 
assumption.
   
   **2. `... GROUP BY a > 50`.** At most 3 groups (true, false, NULL); main 
estimates 10,000. Point: a predicate is also a value, so there is no `Predicate 
| Value` split.
   
   **3. `sales JOIN dims ON sales.b = dims.b WHERE NOT (sales.a > 50)`, `dims` 
with 5,000 rows.** `NOT p` keeps `1 - 0.441 - 0.1` of the rows, about 4,590, so 
`sales` is the build side. Ignoring the null share gives about 5,590 and picks 
`dims`. Point: three-valued logic needs the null share, which comes from the 
predicate's null count.
   
   **4. `... GROUP BY date_trunc('month', ts)`.** 12 groups, where main 
estimates 10,000, through a provider that knows `date_trunc`. In a two-phase 
plan the partial stage publishes the distinct count on its output column and 
the final stage reads it back. #25611 reports the same shape on TPC-H SF1: 
6,001,215 groups estimated for 7 actual.
   
   **5. `... GROUP BY upper(s)`.** A consumer that must not underestimate (for 
example future hash table presizing) wants the bound, 5,000. A consumer that 
compares alternatives (join ordering) wants a best guess, such as 3,000. Point: 
the consumer's need belongs to the request.
   
   **6. `sales JOIN dims ON sales.b = dims.b WHERE sales.a > 50 AND sales.c > 
50`, `dims` with 3,000 rows.** Independence gives about 2,160 rows, so `sales` 
builds. With `c` equal to `a` the answer is up to 4,410, so `dims` builds. 
Point: no built-in correlation, but a provider holding correlation metadata can 
answer the `AND`.
   
   **7. `orders JOIN customers ON customer_id = id JOIN regions ON 
customers.region = regions.id`** (keys 0..99 over 1,000 rows, keys 50..149 over 
100 rows, `regions` with 700 rows). The distinct-count formula estimates the 
first join at 1,000 rows, so `regions` builds the second join. A sketch 
intersection estimates 500, which is the truth, so the first join's result 
builds. Point: custom statistics travel per column. A built-in rule drops one 
it does not understand: with `customer_id + 1 = id` the distinct count carries 
through, but a hash sketch does not.
   
   Two more edge cases: `a - a` and `a = a` (a rule sees its children, so it 
can recognise the same column on both sides), and `random() < 0.1 AND random() 
< 0.1` (0.01, because the two calls are independent; a rule can check 
volatility).
   
   ### Decisions
   
   **You were right about the propagated unit.** What propagates through the 
expression tree should be one synopsis per expression, not N independent 
answers. Your duplication argument holds: with separate methods every 
expression implements N rules, with one synopsis it implements one. The 
examples add a second reason. In a two-phase aggregate (case 4), the group 
expression exists only in the partial stage, and the final stage groups on a 
bare column of the partial output. The partial stage can pass the expression's 
distinct count upward only as the `distinct_count` of that output column, so 
the synopsis has to carry `ColumnStatistics`. With four independent getters 
there is nothing to write into that column.
   
   **My laziness concern from May does not need separate methods.** Built-in 
rules are local and cheap: each combines the synopses of its children and never 
walks the tree itself, and the walk caches one synopsis per structurally equal 
sub-expression. The expensive cases I worried about (sampling, external 
lookups) belong to providers, which decide for themselves what to compute. Case 
5 shows when a request would matter (a bound for one consumer, a best guess for 
another). I would not add it yet: the per-call inputs go in an arguments struct 
with private fields, the same pattern as `StatisticsArgs`, so the request can 
become a field later without changing the `PhysicalExpr` method.
   
   **A flat struct instead of `Predicate | Value`.** I agree `a + 1` must not 
report a selectivity, but the partition is not total. A predicate is also a 
value: `GROUP BY a > 50` (case 2) needs the distinct count of a predicate. And 
a selectivity alone is not enough for three-valued logic (case 3): `NOT p` is 
`1 - selectivity(p) - null_count(p) / rows`. The type-safety part can be 
enforced in one place instead: the synopsis carries its data type, set from the 
schema, and a synopsis that sets a selectivity for a non-Boolean expression is 
rejected.
   
   **Estimation only, next to the existing mechanisms.** The synopsis never 
drives a correctness decision such as pruning, because an estimate can be 
wrong. It does not compute ranges: its min/max come from `evaluate_bounds`, so 
the careful floating-point rounding stays in one place. The dependency runs one 
way: the synopsis consumes `evaluate_bounds`, and the correctness users of 
`evaluate_bounds`, such as the symmetric hash join, never consume the synopsis. 
For a predicate that interval analysis supports, `analyze()` stays the 
estimate, and the synopsis rules combine the rest, so `p_size = 15 AND p_type 
LIKE '%BRASS%'` (#25612) keeps the interval estimate of its first conjunct. It 
does not overlap with your #19609, which computes sound bounds per container 
for pruning. Given that work, does this split between interval analysis and the 
synopsis match how you see the two fitting together?
   
   **No distributions.** I would keep the unification and leave out parametric 
distributions (the Statistics V2 shape). Skew, correlation or histograms can 
travel as extensions, supplied by a provider that understands them: case 6 for 
correlation, case 7 for a set sketch that gets a join right where the 
distinct-count formula is off by a factor of two.
   
   The shape these decisions give, as a sketch:
   
   ```rust
   pub struct ExprSynopsis {
       pub column: ColumnStatistics,   // distinct count, min/max, null count
       pub selectivity: Option<f64>,   // Boolean expressions only
       pub extensions: Extensions,     // custom statistics, from providers
       pub data_type: DataType,
   }
   
   // On PhysicalExpr, with a default that returns None:
   fn synopsis_from_inputs(&self, args: &SynopsisArgs, child_synopses: 
&[ExprSynopsis])
       -> Option<ExprSynopsis>;
   ```
   
   If the decisions work for you, I will open the framework PR (the synopsis, 
the walk and the provider chain, with the aggregate as the first consumer), as 
mentioned on #25610.


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