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]