Repository navigation
Predicates outside an OR do not restrict its branches, and adding an index there can cause a large regression #19339
Description
Activity
I've opened PR #19350 to fix this issue. The fix adds a DistributeConjunctsIntoOrFilterOptimizer that implements the planner-based approach (option 2 from the issue): selective conjuncts (EQUALS/IN on single columns) from an enclosing AND are distributed into each OR branch, so that P AND (A OR B) becomes P AND ((P AND A) OR (P AND B)). This ensures the OR subtree's scan-based predicates are only evaluated against documents matching the selective predicate P.
We hit the same cliff, from a different direction. Instead of a subquery inside the OR, the OR is on two time ranges (day-over-day comparison) and the scan appears only on segments that lack the configured range index.
Query shape (v1 engine, realtime table, ~134K segments):
SELECT <10-min bucket>, FUNNELCOUNT(...)
FROM events
WHERE (
(name = 'A' AND flow_id = 'F' AND action_id = 'X')
OR (name = 'B' AND flow_id = 'F')
)
AND device_os IN ('android', 'ios')
AND (
(ts >= :a0 AND ts < :a1)
OR (ts >= :b0 AND ts < :b1) -- b = a - 24h
)
GROUP BY 1name, flow_id, action_id, device_os have inverted indexes. ts is the time column and has rangeIndexColumns configured, but ~3% of segments are missing the index on disk (separate issue, being fixed by reload).
Numbers
| window A only | A OR B | |
|---|---|---|
numEntriesScannedInFilter |
101,187 | 441,757,659 |
numDocsScanned |
66,800 | 154,975 |
numSegmentsProcessed |
1,384 | 4,541 |
realtimeThreadMemAllocatedBytes |
684 MB | 1.83 GB |
~4,400× more filter work for ~2.3× more matched rows. Window B alone behaves like window A alone.
Verbose EXPLAIN on the OR query (all cohorts, 4,588 segments after pruning):
| plan on the time predicate | segments |
|---|---|
single FILTER_RANGE_INDEX under FILTER_AND |
3,449 |
FILTER_OR → 2× FILTER_RANGE_INDEX |
1,004 |
single FILTER_FULL_SCAN under FILTER_AND |
44 |
FILTER_OR → 2× FILTER_FULL_SCAN |
83 |
FILTER_EMPTY |
8 |
The 83 segments account for essentially all of the 441M: ~2.7M docs × 2 predicates each, i.e. a full pass over the column per range predicate. The 44 unindexed segments whose scan sits directly under the AND cost almost nothing, because AndDocIdSet hands them the intersected inverted-index bitmap via ScanBasedDocIdIterator.applyAnd. The 83 get no candidate set because OrDocIdSet has no equivalent, and AndDocIdSet#iterator() calls iterator() on the OR child before it has computed anything.
Two details that made this particularly sharp for us:
- Per segment, a range predicate that can't match the segment's min/max becomes
EmptyFilterOperatorand is dropped from the OR, so the OR only survives on segments that overlap both windows. - Each window on its own never exposes the missing index (~100K entries), so this was invisible until both windows were combined in one query.
On the two proposed fixes
+1 to option 1 (push applyAnd down into composite doc-id sets). For our shape, #19350 as currently scoped would help only marginally: the only conjunct it can distribute is device__os IN ('android','ios'), which matches most rows. The selective part of the outer AND is the OR(AND(name, flow_id, action_id), AND(name, flow_id)) block, which isn't a single-column EQ/IN and so wouldn't be pushed into the time-range branches.
@gortiz if you push the intersection down to the range index you can make use of RangeBitmap’s context parameter, which allows the range index evaluation to skip over unnecessary regions of the range index itself. This is typically much faster; RangeBitmap is being used for observability at Yandex and making use of the context parameter to push intersections down is critical to performance there.
I've opened #19408 implementing option 1, the execution-time push-down.
applyAnd is generalized from ScanBasedDocIdIterator up to BlockDocIdSet, so an AND hands the document ids it has matched to its composite children: OrDocIdSet unions the restricted branches, AndDocIdSet intersects starting from the candidate set, NotDocIdSet subtracts. No heuristic, no duplicated index lookups, any nesting depth, and the eager/lazy cliff goes away rather than being routed around.
It ships disabled. The push-down materializes the filter result instead of streaming it, so a query that stops at its LIMIT ends up doing the full filter work — measured on the existing suite at 21–41% more entries scanned for a selection with LIMIT, against 7–19% fewer for aggregations. Hence andRestrictionPushdownMode: NEVER (default), AUTO (only queries that read every matching document, i.e. aggregation and group-by), ALWAYS.
@alexch2000 — I think this covers your shape, and I'd appreciate a check against your data. Your query is an aggregation with a GROUP BY, so AUTO enables it. The 83 segments where the plan is FILTER_OR → 2× FILTER_FULL_SCAN should now receive the candidate set from OR(cohorts) ∩ device_os instead of evaluating each range predicate over the whole column. Worth spelling out why it is so bad today: an OR left lazy is driven through advance(target), and SVScanDocIdIterator.advance scans forward from the target until it finds a match — so on a segment where the predicate matches nothing it walks to the end. The push-down replaces that with a batched applyAnd bounded by the candidate set, which is where your ~441M goes.
@richardstartin — I like the RangeBitmap context suggestion a lot, and I think this PR is a good short-term fix that yours can be built on top of rather than an alternative to it. Two reasons:
- Something has to deliver a bitmap to the leaf, and today nothing carries one across an OR boundary —
AndDocIdSetonly pushes into its own scan children, andStarTreeFilterOperatorandFilterPlanNode#wirePreFilterForVectorOperatorsare both AND-only. This PR is what makes a range index nested under an OR reachable at all. - The wiring already exists for another index type.
wirePreFilterForVectorOperatorsdoes exactly the shape you're describing — partition the AND's children, checkFilterAwareVectorIndexReader#supportsPreFilter(), check the producers can cheaply produce bitmaps, apply a selectivity cost model, then push the bitmap in. AFilterAwareRangeIndexReaderforwarding torangeBitmap.between(min, max, context)looks like a fairly mechanical generalization of it.
The real blocker for the context parameter is one level below this PR: AndFilterOperator#getTrues() evaluates every child eagerly, so BitSlicedRangeIndexReader has already produced its full bitmap by the time any BlockDocIdSet exists. Consuming a context needs the restriction to arrive at the operator level, which means ordered rather than eager evaluation of the AND's children. Happy to open that as a separate issue if you think that framing is right — it is also the one that would speed up a plain AND(a = 1, ts BETWEEN …) with no OR anywhere, which this PR does nothing for.
For completeness on the other option in the description: I don't think the planner rewrite in #19350 is the right primary fix. It only distributes single-column EQ/IN, so it cannot reach @alexch2000's case; it costs an extra index read per branch on the common all-indexed OR, where there is nothing to gain; and it leaves the cliff in place — adding an index to a column inside an OR branch still flips that branch to the eager path, just with P inside it now.
@gortiz Some historical context that seems relevant to #19408:
- Performance degradation case in evaluation of nested AND-OR operators. #10396 describes essentially the same missing candidate-bitmap pushdown through an OR: selective indexed predicates combined with
OR(isSubnetOf(...), ...), whereadvance()can scan almost the entire segment. It proposes wiring the outer bitmap into the nested scans. I previously linked it to Improve AND filter operator #9839, which tracks restricted scans, switching from indexes to scans when few candidates remain, and lazy evaluation / limit pushdown. - [Bug] Full scan is happening on all docs instead of just those returned by the inverted index #9402 is an earlier production report of
P AND ((A AND B) OR C), with roughly 42.8M versus 554 filter entries scanned after changing a nested predicate. The discussion also points to Provide bitmap from previous stages of queries toRangeIndexReader#7597 for passing previously matched documents into the range index. - [WIP] Enhance applyAnd() for ScanBasedDocIdIterator #5833, addressing Use efficient search for the intersection of ImmutableRoaringBitmap and ScanBasedDocIdIterator #5596, experimented with changing how multiple scan predicates are evaluated against bitmap candidates. It was closed without merging after benchmarks showed no benefit with enough appropriate indexes and mixed gains/losses with fewer indexes (author's explanation). This was a different scan strategy, so it is a caution about workload dependence rather than evidence against OR candidate pushdown itself.
- Fix the ConcurrentModificationException for And/Or DocIdSet #12611 accidentally stopped registering bitmap iterators in
OrDocIdSet; Fix Bitmap Filter Execution Order in OrDocIdSet #15756 fixed the resulting loss of bitmap-based scan restriction. One workload saw almost 100x higher latency in Pinot 1.2 (report), and some combinations of sorted and inverted indexes also returned incorrect results. This was an implementation regression, but it illustrates how sensitive the surrounding iterator contracts are. - In Evaluate cheap filters before expensive filters #8453, we already discussed that applying a bitmap eagerly can be appropriate for aggregations while selections should evaluate lazily and stop at their limit (discussion). This is directly relevant to the tradeoff measured in Push AND restrictions into composite filter children #19408.
My reading is that the history supports fixing the execution-layer gap. It does not establish that an earlier implementation of this same OR pushdown was tried and failed. Restricting which documents a filter examines and deciding how eagerly to examine them are separate decisions; improving the former can still regress early-terminating queries through the latter.
The disabled default and conservative AUTO mode in #19408 acknowledge that tradeoff. Could you evaluate the history above against the current proposal and see whether any adjustments are needed? In particular:
- Can candidate propagation preserve lazy/batched evaluation and early termination, or should that remain an explicit follow-up with the current mode gating?
- Should deferral be limited to subtrees that contain scan/expression predicates, given the acknowledged extra bitmap work for entirely indexed subtrees?
- Can we validate with latency and allocation benchmarks, in addition to entries scanned, across aggregations and limited selections, sparse and dense candidate sets, and mixed versus entirely indexed branches?
Please summarize which historical cases are already covered and whether they suggest changes to the implementation, mode eligibility, or regression tests.
Thanks @Jackie-Jiang — that history is exactly the right frame, and I agree with your reading. I have worked through all three questions; #19408 has been updated.
2. Should deferral be limited to subtrees containing scan/expression predicates?
Yes. Done.
Deferral is now driven by BlockDocIdSet#isScanBased(), which the composites compute over their children; isApplyAndDeferrable() follows it. An index-only subtree returns the same bitmap either way, so handing it a candidate set only added an intersection per branch.
Reading this off the DocIdSet rather than the operator tree turned out to be both simpler — no new constructor parameters anywhere — and more accurate. An operator that scans internally and hands back a bitmap, such as a non-exact range index or an H3 index, has already paid for that scan by the time its DocIdSet exists, so restricting it later would gain nothing; the DocIdSet-level view correctly reports false, where an operator-level flag would have said "scan-based". It also takes the last piece of policy out of isApplyAndDeferrable(), so the mode is enforced only by the parent AndDocIdSet.
Measured neutral: INDEXED_OR below is 0.99–1.02x time and 1.00x allocation.
3. Can we validate with latency and allocation benchmarks?
Done — BenchmarkAndRestrictionPushdown in pinot-perf. 2M documents, a real SVScanDocIdIterator over a fixed-bit forward index, -prof gc for allocation. Ratios of push-down on to push-down off, lower is better:
| shape | consume | sel | time | alloc |
|---|---|---|---|---|
| SCAN_IN_OR | DRAIN | 0.01 | 0.03x | 0.10x |
| SCAN_IN_OR | DRAIN | 0.5 | 0.67x | 1.14x |
| SCAN_IN_OR | LIMIT | 0.01 | 0.02x | 0.10x |
| SCAN_IN_OR | LIMIT | 0.5 | 0.64x | 1.14x |
| SCAN_ONLY_OR | DRAIN | 0.01 | 0.28x | 21.95x |
| SCAN_ONLY_OR | DRAIN | 0.5 | 0.34x | 231.14x |
| SCAN_ONLY_OR | LIMIT | 0.01 | 144.77x | 25.04x |
| SCAN_ONLY_OR | LIMIT | 0.5 | 9431.50x | 383.65x |
| INDEXED_OR | all four | 0.99–1.02x | 1.00x |
SCAN_IN_OR is a scan next to an index-based predicate inside an OR branch — the shape reported here. SCAN_ONLY_OR is scans directly under the OR with no index-based sibling. INDEXED_OR has no scan anywhere.
You were right to push on this: the entries-scanned proxy understated the risk by three orders of magnitude. What the PR previously described as "21–41% more entries scanned" for a limited selection is, in the worst shape, 9431x slower and 384x more allocation. That shape's OR is already lazy on master, so materializing it replaces a scan of a few dozen documents with a scan of the whole segment.
Two other things fell out that I did not expect:
SCAN_IN_ORwith a LIMIT is faster with the push-down, not slower. That OR is not lazy on master either — it is driven byadvance(), andSVScanDocIdIterator.advancescans forward from the target until it finds a match. So the loss is not "selections" in general; it is specifically ORs that are genuinely lazy today.- The gain on the reported shape is largest exactly where the enclosing AND is selective, which is the production case in this issue.
1. Can candidate propagation preserve lazy evaluation and early termination?
Yes, and there is now a spike with numbers rather than a design argument: #19589 (draft, stacked on #19408).
It adds a re-entrant per-chunk kernel, ScanBasedDocIdIterator#matchDocIds(int[], int), and a RestrictedScanDocIdIterator that drives a scan from an upstream candidate stream a chunk at a time. The scan still never sees a document outside the candidate set, but nothing is materialized. The only thing that made applyAnd non-re-entrant was a trailing close() in SVScanDocIdIterator releasing one ForwardIndexReaderContext; MVScanDocIdIterator#applyAnd was already re-entrant.
Same benchmark, three strategies, with a setup assertion that all three return identical documents:
| shape | consume | sel | eager time | streaming time | eager alloc | streaming alloc |
|---|---|---|---|---|---|---|
| SCAN_IN_OR | DRAIN | 0.01 | 0.03x | 0.03x | 0.10x | 0.06x |
| SCAN_IN_OR | DRAIN | 0.5 | 0.66x | 1.21x | 1.14x | 0.27x |
| SCAN_IN_OR | LIMIT | 0.01 | 0.02x | 0.01x | 0.10x | 0.06x |
| SCAN_IN_OR | LIMIT | 0.5 | 0.58x | 0.04x | 1.14x | 0.26x |
| SCAN_ONLY_OR | DRAIN | 0.01 | 0.32x | 0.40x | 21.97x | 1.41x |
| SCAN_ONLY_OR | DRAIN | 0.5 | 0.34x | 0.53x | 232.27x | 1.17x |
| SCAN_ONLY_OR | LIMIT | 0.01 | 138.85x | 2.23x | 25.04x | 1.33x |
| SCAN_ONLY_OR | LIMIT | 0.5 | 7925.64x | 3.27x | 383.31x | 1.33x |
Streaming removes the cliff — worst case 3.27x instead of 7926x, allocation 1.33x instead of 383x — and beats eager wherever the consumer stops early, because it keeps the early termination eager destroys.
But eager is still the better strategy for a full drain: 0.66x against 1.21x on SCAN_IN_OR at high candidate density, where streaming is actually slower than master. Chunking costs more than it saves when every matching document will be read anyway.
So my answer to your first question is: streaming does not let us drop the mode, it changes what the mode chooses — eager when the query will drain the filter, streaming when it may stop early. AUTO stops being a safety gate and becomes a strategy selector, and its excluded branch goes from merely safe to fast.
I would still keep it as an explicit follow-up rather than folding it into #19408, for two reasons: the spike is not wired into AndDocIdSet, and the scan in the benchmark is an SVScanDocIdIterator. An ExpressionScanDocIdIterator — the IN_SUBQUERY case that opened this issue — builds a whole ProjectionOperator per call, so per chunk it will want a much larger chunk, and a larger chunk erodes the early-termination granularity that is the entire point. Measuring that trade-off is the entry criterion for doing it properly, and it is still open.
The historical cases
- Performance degradation case in evaluation of nested AND-OR operators. #10396 and [Bug] Full scan is happening on all docs instead of just those returned by the inverted index #9402 are both fixed by Push AND restrictions into composite filter children #19408 and both still open. [Bug] Full scan is happening on all docs instead of just those returned by the inverted index #9402's
P AND ((A AND B) OR C)isSCAN_IN_ORabove; Performance degradation case in evaluation of nested AND-OR operators. #10396'sadvance()scanning nearly the whole segment is the mechanism I describe under question 3. With this issue that is three independent reports of the same gap. - Provide bitmap from previous stages of queries to
RangeIndexReader#7597 is not covered. Pushing the restriction into the range index itself viaRangeBitmap's context parameter is a layer below this, and needs the restriction to arrive at the operator level —AndFilterOperator#getTrues()evaluates every child eagerly, soBitSlicedRangeIndexReaderhas already produced its bitmap by the time anyBlockDocIdSetexists.FilterPlanNode#wirePreFilterForVectorOperatorsalready does exactly that shape for vector indexes, so aFilterAwareRangeIndexReaderlooks like a fairly mechanical generalization. Happy to open that separately. - [WIP] Enhance applyAnd() for ScanBasedDocIdIterator #5833 I read as you do: a different scan strategy, closed for lack of time rather than on the evidence. The caution about workload dependence is fair, and the matrix above is my attempt to answer it directly — which is how the
SCAN_ONLY_ORcliff surfaced at all. - Fix the ConcurrentModificationException for And/Or DocIdSet #12611 / Fix Bitmap Filter Execution Order in OrDocIdSet #15756 is the one I took most seriously, since it is the same iterator contracts. Push AND restrictions into composite filter children #19408 adds a
release()for children abandoned by a short-circuit (those iterators holdForwardIndexReaderContexts and are otherwise closed only by reaching EOF), rejects a second consumption of a composite DocIdSet, and preserves the[0, numDocs)bound thatNotDocIdIteratorapplies — the candidate set can carry ids pastnumDocsbecauseBitmapDocIdIterator#getDocIds()hands out the raw bitmap while itsnext()clamps. An intersection cannot introduce such an id, but a complement can. The 200-tree randomized differential test is there to catch exactly the class of regression Fix Bitmap Filter Execution Order in OrDocIdSet #15756 had to fix. - Evaluate cheap filters before expensive filters #8453 is the one that most directly supports the design. The eager-for-aggregations, lazy-for-selections split you discussed there is what
AUTOencodes, and the numbers above put a size on it.
Adjustments made
The default stays NEVER. AUTO is unchanged. What changed as a result of your comment: deferral is now limited to scan-bearing subtrees, the benchmark exists, and the PR description now states the SCAN_ONLY_OR risk in latency and allocation rather than in entries scanned — ALWAYS is genuinely dangerous on a selection workload, and the PR previously undersold that.
Description
When a filter has the shape
P AND (A OR B), Pinot evaluates the(A OR B)subtree independently ofP. IfPis highly selective, and a branch of the OR contains a predicate with no index — for exampleIN_ID_SET(...), produced byIN_SUBQUERY, which is evaluated per document throughExpressionScanDocIdIterator— that predicate is evaluated for every document matching the branch, not only for the documents that satisfyP.Duplicating
Pinside the OR branch by hand is a sound rewrite (P ∧ (A∨B) ≡ P ∧ ((P∧A) ∨ (P∧B)), and it holds under three-valued logic because the filter only passes TRUE). Doing so makes such queries several times faster, which shows the restriction is simply not being applied.The cost also depends on whether the branch columns have indexes
AndDocIdSet#iterator()chooses between two strategies atAndDocIdSet.java:121:OrDocIdSetkeeps it lazy, the outerAndDocIdSetplaces it inremainingDocIdIterators, and the outer merged bitmap — which includesP— drives it. The expensive predicate is only evaluated at documents that already matchP.P, before the outer AND intersects.So adding an index to a column that appears inside an OR branch can make a query orders of magnitude slower. In a production deployment, the same query returning the same two rows went from ~500 to ~17,000,000
numEntriesScannedInFilter, and allocated 4.6 GB, after the only change was two columns in the branch gaining a range index and an inverted index respectively.Reproduction sketch
Table
events(tenant_id INT, ts LONG, kind STRING, id LONG, value DOUBLE), one selective tenant among many.Run it twice: once with no index on
tsorkind, once with a range index ontsand an inverted index onkind. The second configuration is far slower and scans far more entries in the filter.Suggested fixes
applyAndfromScanBasedDocIdIteratortoBlockDocIdSet, soAndDocIdSetcan pass its merged index bitmap into composite children:OrDocIdSet.applyAnd(b)= union ofchild.applyAnd(b),AndDocIdSet.applyAnd(b)= intersect starting fromb. This needs no heuristic, adds no duplicated index lookups, works at any nesting depth, and removes the eager/lazy cliff. Needs care withnumEntriesScannedInFilteraccounting,NotDocIdSet, and null handling.FilterOptimizerthat distributes selective conjuncts (EQ/IN on dictionary-encoded columns) from an enclosing AND into each OR branch. Cheaper to implement and gateable behind a query option, but it is a heuristic: duplicating a predicate on a raw column doubles a scan.