Visitar URL original
Predicates outside an OR do not restrict its branches, and adding an index there can cause a large regression · Issue #19339 · apache/pinot · GitHub
Skip to content

Predicates outside an OR do not restrict its branches, and adding an index there can cause a large regression #19339

Description

@gortiz

Description

When a filter has the shape P AND (A OR B), Pinot evaluates the (A OR B) subtree independently of P. If P is highly selective, and a branch of the OR contains a predicate with no index — for example IN_ID_SET(...), produced by IN_SUBQUERY, which is evaluated per document through ExpressionScanDocIdIterator — that predicate is evaluated for every document matching the branch, not only for the documents that satisfy P.

Duplicating P inside 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 at AndDocIdSet.java:121:

if ((numIndexBasedDocIdIterators > 0 && numScanBasedDocIdIterators > 0) || numIndexBasedDocIdIterators > 1) {
  // eager: merge the index bitmaps, then scanIterator.applyAnd(docIds)
} else {
  return new AndDocIdIterator(allDocIdIterators);  // lazy
}
  • With no index-based child in the AND under the OR, the subtree stays lazy. OrDocIdSet keeps it lazy, the outer AndDocIdSet places it in remainingDocIdIterators, and the outer merged bitmap — which includes P — drives it. The expensive predicate is only evaluated at documents that already match P.
  • With one or more index-based children, the eager path is taken. The branch is fully materialized as a bitmap over the whole segment, ignoring 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.

SELECT sum(value) FROM events
WHERE tenant_id = 42
  AND ( ( ts >= :t0 AND ts <= :t1 AND kind = 'a'
          AND IN_SUBQUERY(id, 'SELECT ID_SET(id) FROM events WHERE tenant_id = 42 AND ts >= :t2') = 0 )
        OR ( ts >= :t2 AND kind = 'b' ) )

Run it twice: once with no index on ts or kind, once with a range index on ts and an inverted index on kind. The second configuration is far slower and scans far more entries in the filter.

Suggested fixes

  1. Push the restriction down at execution time (preferred). Generalize applyAnd from ScanBasedDocIdIterator to BlockDocIdSet, so AndDocIdSet can pass its merged index bitmap into composite children: OrDocIdSet.applyAnd(b) = union of child.applyAnd(b), AndDocIdSet.applyAnd(b) = intersect starting from b. This needs no heuristic, adds no duplicated index lookups, works at any nesting depth, and removes the eager/lazy cliff. Needs care with numEntriesScannedInFilter accounting, NotDocIdSet, and null handling.
  2. Rewrite in the planner. Add a FilterOptimizer that 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.

Activity

added
bugSomething is not working as expected
queryRelated to query processing
performanceRelated to performance optimization
on Aug 24, 2026

waterWang commented on Aug 24, 2026

@waterWang

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.

alexch2000 commented on Aug 25, 2026

@alexch2000
Contributor

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 1

name, 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 EmptyFilterOperator and 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.

richardstartin commented on Aug 28, 2026

@richardstartin
Member

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

gortiz commented on Aug 31, 2026

@gortiz
ContributorAuthor

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:

  1. Something has to deliver a bitmap to the leaf, and today nothing carries one across an OR boundary — AndDocIdSet only pushes into its own scan children, and StarTreeFilterOperator and FilterPlanNode#wirePreFilterForVectorOperators are both AND-only. This PR is what makes a range index nested under an OR reachable at all.
  2. The wiring already exists for another index type. wirePreFilterForVectorOperators does exactly the shape you're describing — partition the AND's children, check FilterAwareVectorIndexReader#supportsPreFilter(), check the producers can cheaply produce bitmaps, apply a selectivity cost model, then push the bitmap in. A FilterAwareRangeIndexReader forwarding to rangeBitmap.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.

Jackie-Jiang commented on Sep 8, 2026

@Jackie-Jiang
Contributor

@gortiz Some historical context that seems relevant to #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:

  1. Can candidate propagation preserve lazy/batched evaluation and early termination, or should that remain an explicit follow-up with the current mode gating?
  2. Should deferral be limited to subtrees that contain scan/expression predicates, given the acknowledged extra bitmap work for entirely indexed subtrees?
  3. 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.

gortiz commented on Sep 17, 2026

@gortiz
ContributorAuthor

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_OR with a LIMIT is faster with the push-down, not slower. That OR is not lazy on master either — it is driven by advance(), and SVScanDocIdIterator.advance scans 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

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.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    bugSomething is not working as expectedperformanceRelated to performance optimizationqueryRelated to query processing

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions