gortiz opened a new pull request, #19668:
URL: https://github.com/apache/pinot/pull/19668

   Follow-up to #19419, which added the streaming DISTINCT leaf operator.
   
   ## Problem
   
   `StreamingDistinctCombineOperator` bounds leaf-stage memory for DISTINCT by 
flushing the
   accumulated `DistinctTable` every `streamingDistinctFlushThreshold` values. 
The cost was the
   cross-segment early exit: `DistinctResultsBlockMerger#isQuerySatisfied` can 
only fire once the
   accumulated table reaches LIMIT, and flushing empties it long before that, 
so every segment got
   scanned.
   
   The cost is invisible until the feature is turned on: the affected 
population is exactly
   `SELECT DISTINCT ... LIMIT > streamingDistinctFlushThreshold` over a column 
whose cardinality
   reaches LIMIT — precisely the queries that used to stop early. They get a 
lower memory ceiling and a
   full scan, silently, with no error and no diagnostic. Anyone enabling
   `streamingDistinctFlushThreshold` broadly pays it on exactly those queries.
   
   ## Why a row counter does not work
   
   Flush windows overlap. A window emitting `{a,b,c}` followed by one emitting 
`{b,c,d}` is 6 rows but
   4 distinct values. A counter of emitted rows overshoots LIMIT and lets the 
leaf exit having emitted
   fewer than LIMIT distinct values — a silently truncated result. Restoring 
the exit needs cumulative
   *cardinality* across windows, not a row count.
   
   ## Approach
   
   `DistinctCardinalityTracker` is carried across flush windows as 
main-thread-only state, alongside
   the existing `_numDocsScanned`, and each flushed table is folded into it 
before the accumulator is
   dropped. Values are enumerated through a new 
`DistinctTable#forEachValueHash(LongConsumer)`, which
   walks each subtype's own primitive value set — no boxing, no `Object` 
dispatch, and adding a new
   `DistinctTable` subtype is a compile error rather than a silent 
mis-measurement.
   
   ## Guarantees
   
   Three behaviours, selected by the leaf-stage LIMIT against
   `streamingDistinctMaxTrackedCardinality` (default 16384) and by whether the 
estimated exit has been
   enabled:
   
   | Mode | Applies when | Exit condition | Guarantee |
   |---|---|---|---|
   | **Exact** | `LIMIT <= bound` | `\|hashes\| >= LIMIT` | **100%.** Equal 
values hash equally, so the tracked count can never exceed the number of 
distinct values emitted. Reaching LIMIT is a counted fact. |
   | **Full scan** | `LIMIT > bound` and `stdDev = 0` *(default)* | none — 
every segment is read | **100%**, by reading everything. This is the behaviour 
on master. |
   | **Estimated** | `LIMIT > bound` and `stdDev > 0` | 
`sketch.getLowerBound(stdDev) >= LIMIT` | **Bounded by the chosen `stdDev`** — 
a sketch lower bound is a confidence bound (DataSketches documents 
`getLowerBound` as the *"approximate lower error bound"*), so the exit can fire 
below LIMIT with the probability below. |
   
   Exact mode needs no opt-in. Above the bound the default is the full scan, so 
enabling the estimated
   exit is an explicit choice to trade a bounded probability of a short result 
for not scanning every
   segment.
   
   When it fires below LIMIT the query returns marginally fewer rows than it 
should, with no error
   raised and no partial-result flag. Measured per query over 20k trials per 
configuration, in the
   worst flush alignment (a boundary landing one value below LIMIT):
   
   | stdDev | P(short result) |
   |---|---|
   | 2 | ~2.2% |
   | 3 | ~0.12% |
   
   The probability does not grow with LIMIT, and raising the sketch's nominal 
entries does not reduce
   it — that only reduces overshoot (~1.02x LIMIT). Moving a query into exact 
mode, or leaving the
   estimated exit off, are the only levers that remove the risk.
   
   Method, so the numbers can be reproduced or challenged: per trial, feed 
distinct values into an
   `UpdatableThetaSketch` in `flushThreshold`-sized batches, evaluating 
`getLowerBound(stdDev) >= LIMIT`
   at each batch boundary as the operator does, and record whether the first 
crossing happened before
   the true count reached LIMIT. LIMIT is placed one value above a batch 
boundary, which is the worst
   alignment — when LIMIT sits far above the last boundary the crossing needs a 
many-sigma error and
   never occurs. The per-window evaluations share one cumulative sketch and are 
strongly correlated, so
   the single evaluation closest below LIMIT dominates rather than compounding 
across windows, which is
   why the figure is flat in LIMIT.
   
   ## Configuration
   
   | Query option | Default | Meaning |
   |---|---|---|
   | `streamingDistinctMaxTrackedCardinality` | 16384 | Largest LIMIT tracked 
exactly (~256 KiB at the default), and the sketch's nominal entries above that. 
`0` disables the early exit entirely. Clamped internally at 2^20 (~16 MiB), so 
a client-set value cannot drive unbounded server heap; a LIMIT above the clamp 
falls through to the estimated regime, which is off by default and therefore 
yields no tracker at all. |
   | `streamingDistinctEstimatedExitStdDev` | 0 (off) | Std deviations for the 
sketch lower bound. Must be 1, 2 or 3 — `BinomialBoundsN` accepts nothing else, 
so anything larger is rejected where the option is parsed rather than surfacing 
as a mid-query server exception. `0` disables the estimated regime. |
   
   Matching cluster properties, both injected with `putIfAbsent` so a per-query 
`SET` always wins:
   
   | Cluster property | Default | Meaning |
   |---|---|---|
   | `pinot.broker.mse.streaming.distinct.max.tracked.cardinality` | unset | 
Cluster-level control for the regime that is on by default, so it can be 
switched off (`0`) without touching every client. |
   | `pinot.broker.mse.streaming.distinct.estimated.exit.std.dev` | unset (off) 
| Opts the cluster into the estimated regime. Validated at broker startup, so 
an out-of-range value fails there rather than failing every DISTINCT query at 
plan time. |
   
   ## Other gates
   
   The tracker is not built at all when it cannot help or would be wrong:
   
   - `LIMIT <= flushThreshold` — the accumulator reaches LIMIT inside one 
window, so
     `DistinctTable#isSatisfied()` already short-circuits.
   - `LIMIT == Integer.MAX_VALUE` — an MSE leaf with no LIMIT pushed down; 
nothing can reach it.
   - **ORDER BY** — an ordered distinct must never exit early: the top-LIMIT is 
unknown until every
     segment is seen, so stopping at LIMIT distinct values returns an arbitrary 
LIMIT rather than the
     ordered top-LIMIT. Both modes only establish that LIMIT distinct values 
were emitted, which is
     sufficient for unordered DISTINCT but not for ordered — so this is a hard 
gate, not a tuning knob.
   
   ## Semantics
   
   Unchanged from the non-streaming path: the exit is per-server, and a server 
that has emitted at
   least LIMIT distinct values guarantees the global result holds at least 
LIMIT, so applying LIMIT
   downstream yields a full result. That relies on the FINAL stage 
de-duplicating, which the hash
   exchange over the distinct columns provides.
   
   ## Testing
   
   - `DistinctCardinalityTrackerTest` — regime selection, the default-off gate, 
exactness at and below
     the bound, one-sidedness above it, per-type hashing pinned in **both** 
directions (200 distinct
     values must satisfy LIMIT 200 *and not* LIMIT 201 — this fails if any type 
hashes two equal values
     apart), null marker including multi-column nulls with null handling on, 
unhashable stored types,
     overlapping windows.
   - `StreamingDistinctCombineOperatorTest` — emission stops at LIMIT; scan 
work is actually avoided
     (`maxStreamingPendingBlocks=1`, 16 segments: 250 of 800 docs scanned, vs 
800 of 800 without the
     tracker); overlapping windows do not trigger an early exit; unbounded leaf 
LIMIT never exits.
   - `QueryOptionsUtilsTest` — both new getters, including the out-of-range 
rejection for the std dev
     (`4`, `10`, `MAX_VALUE`) and its exact message; without this the range 
check had no coverage at all.
   - `MultiStageBrokerRequestHandlerTest` — cluster property injection, 
per-query override, the `0` kill
     switch reaching the servers, and an out-of-range cluster config failing at 
broker startup.
   - `StreamingDistinctQueriesTest` (integration) — end to end through the 
broker, asserting on two signals per mode
     because neither suffices alone. `numDocsScanned < totalDocs` shows the 
exit *fired*, which the row count cannot
     see; `rows == LIMIT` shows it was not premature. Covers exact mode, 
estimated mode, the default above the bound
     (where the full scan is exact: without a tracker the only satisfaction 
check left reads the accumulated table,
     which flushing empties below LIMIT, so it can never fire — precisely the 
regression this repairs), and parity with
     the blocking path. The class Javadoc records the one thing the row count 
still cannot catch: the exit is
     per-server but the assertion is on the broker's union, so a leaf stopping 
marginally short is masked — a property
     of multi-server deployments, not of the test.
   


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