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]