yashmayya opened a new pull request, #19674:
URL: https://github.com/apache/pinot/pull/19674
## Summary
Multi-stage (MSE) planning time for queries with large IN lists grows much
faster than the list size. This has caused production planning timeouts for
years (#13617 and the follow-ups below). The cause keeps coming back with each
Calcite upgrade and each new planner rule.
This PR hides large IN lists from Calcite while the query is optimized
("sealing"), and restores them after the last rule phase. Calcite sees one
small, deterministic function call instead of the value list. For all tested
query shapes, EXPLAIN and the plans that servers get are the same as without
sealing. The one known difference is in "What changes and what does not".
Measured on the same build with sealing off (`0`, same as master) and on
(default `20`):
| Query | Values | Sealing off | Sealing on |
|---|---|---|---|
| Fact table with 2 LEFT JOINs, 3 IN filters, GROUP BY 16 keys, FILTER
aggregates, ORDER BY LIMIT | 100,000 | 3,363 ms | 795 ms |
| Same query, V2 physical optimizer | 100,000 | 2,896 ms | 852 ms |
| `JOIN ... ON a.k = b.k AND b.x IN (...)` | 400 | 2,505 ms | 23 ms |
| `JOIN ... ON a.k = b.k AND b.x IN (...)` | 800 | 24.3 s | 43 ms |
| `SUM(CASE WHEN x IN (...) THEN 1 ELSE 0 END)` | 800 | 21.4 s | 35 ms |
| Single table, `WHERE x IN (...)` | 100,000 | 650 ms | 542 ms |
The optimize phase of the first query takes 12–16 ms at every list size with
sealing, against 2,514 ms at 100k values without it. Full tables and the method
are in the Benchmarks section below.
## Root cause
- Calcite has two regimes for IN lists. A list with fewer than
`inSubQueryThreshold` values (default 20) becomes `SEARCH(col, Sarg[...])`,
which Calcite optimizes in depth. A longer list becomes a join against an
inline `VALUES` table, so the values are data that the simplifier never reads.
- Pinot sets `withInSubQueryThreshold(Integer.MAX_VALUE)` (#10168), so every
list stays one `SEARCH` call through every rule and metadata call.
- `RexSimplify` rebuilds the whole range set of the Sarg on every AND or OR
that holds it, even when nothing merges. At 37k values this costs about 29 ms
and 26 MB per call.
- The number of such calls grows with each release. CALCITE-5036 (Calcite
1.31) and CALCITE-7302 (1.42) added calls in predicate inference. #18554 and
#18720 added rules that call `getPulledUpPredicates`. On 1.42, one
`getPulledUpPredicates(Join)` call costs 122 ms at 37k values, compared with 2
ms on 1.40.
- HepPlanner clears the metadata cache of a node and all its ancestors after
each rewrite, so the same inference runs many times per query.
- Outside WHERE and HAVING (JOIN ON, CASE, `FILTER (WHERE ...)`, the SELECT
list, GROUP BY and window expressions), SqlToRel builds an OR of N equalities,
and `RelBuilder` simplifies each term against the negation of all earlier
terms. This is cubic: on master, 800 values take 20+ s to plan. It is not a
1.42 regression, and no rule config avoids it.
## How it works
1. **Mark** (`SearchSealer#markInLists`). After validation, each `IN` / `NOT
IN` with at least T values gets a marker operator in place
(`PinotInListOperator`). The marker keeps the name, type derivation and unparse
of the original, and it is restored right after SqlToRel. Its kind is not `IN`,
so SqlToRel does not expand the list into an OR.
2. **Convert** (`PinotConvertletTable`). Each value becomes `operand =
value` through the standard convertlets, exactly like `convertInToOr`, so the
types and casts are the same as today. The points of one operand are folded
once into one `SEARCH` with `RexBuilder#makeIn`. `NOT IN` uses the negated
Sarg. NULL and non-literal values stay as exact OR / AND terms. There is no OR
of N terms, so the cubic path is gone.
3. **Seal** (`SearchSealer#seal`). After SqlToRel and field trimming, each
filter and join condition that holds a large Sarg is simplified once, like
`RelBuilder#filter` does. So the other predicates on the same operand fold into
the Sarg exactly as today: `x IS NOT NULL`, ranges, a second list, and
contradictions such as `x IN (...) AND x = 5`. Then every `SEARCH` with at
least T ranges is sealed, including one that Calcite built from a user-written
`x = 1 OR x = 2 ...` chain.
4. **Unseal** (`SearchSealer#unseal`). After the trait program, every sealed
call becomes `SEARCH(operand, sarg)` again, and `NOT(sealed)` becomes a
`SEARCH` with the negated Sarg, as `RexSimplify` builds today. Copied nodes
keep their traits: `LogicalJoin#copy` drops the trait set it is given, so joins
are created again with their distribution. `RexExpressionUtils` also accepts a
sealed call, in case a conversion runs before the end.
`PinotSealedSearchOperator` is the sealed form. Each distinct Sarg of a
query gets its own instance, which holds the Sarg literal:
| Property | Value | Why |
|---|---|---|
| Digest, `equals`, `hashCode` | `$SEARCH#<id>(operand)`, identity | O(1).
Different sets never merge, and equal sets share one instance. |
| Kind | `OTHER_FUNCTION` | Calcite reads operand 1 of every `SEARCH` as a
Sarg literal. |
| Deterministic, safe | yes | Rules can push it down, pull it up and copy it
across equi-joins. |
| `Strong` policy | `ANY` for `NULL AS UNKNOWN`, `NOT_NULL` otherwise | The
same answers as `Strong` gives for `SEARCH`, so null-rejecting filters still
turn outer joins into inner joins. |
| Return type | `BOOLEAN`, nullable like `SqlSearchOperator` | Row types
stay stable across rewrites. |
| Folding | not folded during optimization | Pinot's HepPlanner has no
executor. Unseal and `RexExpressionUtils` fold a literal operand, as for
`SEARCH`. |
The state is in the operator, not in a `RexCall` subclass, because
`RexCopier` and `RexBuilder#makeCall` rebuild calls from `call.getOperator()`.
This PR also changes `RexExpressionUtils#handleSearch` for a Sarg that mixes
points and ranges, for example `x IN (<list>) OR x > 10` or `x NOT IN (<list>)
AND x > 0`. Today each point becomes `x >= v AND x <= v`, so the servers scan
the segment once per value. With a 10k-value list this shape costs 2.7–8.8 s
per 250k-row segment on master, against 3 ms for the same query in the
single-stage engine, and the stage is 8x larger. Now, when the ranges (or their
complement) hold at least 20 single points, the points become one `IN` (or
`NOT_IN`) plus the other ranges. Smaller Sargs keep today's shape.
`BIG_DECIMAL` keeps today's shape too, because an intermediate stage matches
`IN` values with `equals`, which depends on the scale.
## What changes and what does not
- No change for the tested shapes: the leaf filter, broker segment and
partition pruning, EXPLAIN output, serialized stage plans, exchange
distributions, and the V2 physical optimizer input.
- No wire change: the sealed operator never leaves the broker, so mixed
broker and server versions are safe.
- Kept: filter push-down, transitive predicates across equi-joins,
outer-join simplification, and hints.
- Lost for sealed lists (T or more values): reasoning across plan nodes
during optimization. When a rule moves a predicate next to a sealed list on the
same column (for example a `WHERE` filter pushed into a join input that has a
list from the `ON` clause), Calcite does not merge it into the Sarg.
Contradictions across clauses are not folded during planning, and the query
returns 0 rows at runtime.
- One visible case of this: an explicit `x IS NOT NULL` in another clause
than a sealed list on `x` is dropped as redundant when the two meet. This is
correct in SQL, and plans with null handling on return the same rows. With
per-column nullable schemas and null handling off, the leaf no longer filters
the rows where `x` is null. Master already does this today for `x < 5 AND x IS
NOT NULL`, even in one clause. The same case in one clause keeps its null check
(tested).
## Config
- Broker config `pinot.broker.multistage.sealed.in.list.threshold` and query
option `sealedInListThreshold`. Default: 20, the same size at which Calcite
itself stops treating an IN list as a scalar predicate.
- `0` or a negative value turns sealing off (master behavior). This is the
kill switch, per broker or per query.
- Planning outside the broker's multi-stage request handler (for example the
controller `/sql` endpoint) does not read the broker config and uses the
default. The query option works everywhere.
## Why not the alternatives
- **Calcite's default (IN list → join against `VALUES`).** Measured with
`inSubQueryThreshold=20` on a 2M-row table: `NOT IN`, `IN ... OR ...` and CASE
shapes shuffle 2.0M–8.1M rows where 15 rows move today. V2 always shuffles the
fact rows, and broker pruning is lost. Since Calcite 1.38 (CALCITE-6599),
`RelMdPredicates` turns the VALUES table back into one big Sarg, so planning is
1.3–4x slower. With `expand=true` it also returns NULL rows for `NOT IN` on
nullable columns, and it fails for IN lists in `JOIN ... ON`. Earlier attempts
(#13605, #15027) hit the same problems.
- **Fix Calcite.** A 54-line `RexSimplify` patch removes the 1.42 cost (122
ms → 2 ms per join call). It does not remove the cubic SqlToRel path, and it
does not protect against the next caller. It is worth doing upstream as well.
- **Disable rules.** This removes one caller at a time and loses the
optimizations of those rules for all queries.
## Benchmarks
- Method: planning only (parse to dispatchable plan), on the same build with
`SET sealedInListThreshold=0` (off, same as master) and the default (on).
Synthetic tables: a 22-column fact table on 4 servers with 16 segments, and 2
dimension tables. Warm medians on an Apple Silicon laptop with JDK 25. Compare
ratios, not absolute times.
- IN plus a range on the same column (`x NOT IN (<10k values>) AND x > 0`):
master ships 795,076 bytes in the leaf stage, one range per value. This PR
ships 98,047 bytes (one `NOT_IN` plus the range), as SSE does.
<details>
<summary>Full tables</summary>
#### Linear shapes (V1 planner, default rules)
| Query | N | Off | On | Speed-up | Optimize off → on |
|---|---|---|---|---|---|
| Fact + 2 LEFT JOINs, 3 IN filters, GROUP BY 16 keys, FILTER aggs, ORDER BY
LIMIT | 1,000 | 66 ms | 43 ms | 1.5x | 30 → 12 ms |
| Fact + 2 LEFT JOINs, 3 IN filters, GROUP BY 16 keys, FILTER aggs, ORDER BY
LIMIT | 10,000 | 322 ms | 135 ms | 2.4x | 219 → 12 ms |
| Fact + 2 LEFT JOINs, 3 IN filters, GROUP BY 16 keys, FILTER aggs, ORDER BY
LIMIT | 40,000 | 1,326 ms | 426 ms | 3.1x | 922 → 15 ms |
| Fact + 2 LEFT JOINs, 3 IN filters, GROUP BY 16 keys, FILTER aggs, ORDER BY
LIMIT | 100,000 | 3,363 ms | 795 ms | 4.2x | 2,514 → 15 ms |
| Same, wrapped in an outer COUNT / SUM | 1,000 | 62 ms | 46 ms | 1.3x | 31
→ 12 ms |
| Same, wrapped in an outer COUNT / SUM | 10,000 | 309 ms | 120 ms | 2.6x |
215 → 12 ms |
| Same, wrapped in an outer COUNT / SUM | 40,000 | 1,303 ms | 437 ms | 3.0x
| 895 → 15 ms |
| Same, wrapped in an outer COUNT / SUM | 100,000 | 3,298 ms | 793 ms | 4.2x
| 2,512 → 15 ms |
| Same with NOT IN | 1,000 | 70 ms | 47 ms | 1.5x | 37 → 12 ms |
| Same with NOT IN | 10,000 | 330 ms | 128 ms | 2.6x | 227 → 11 ms |
| Same with NOT IN | 40,000 | 1,320 ms | 450 ms | 2.9x | 891 → 15 ms |
| Same with NOT IN | 100,000 | 3,594 ms | 923 ms | 3.9x | 2,636 → 16 ms |
| Single table, WHERE x IN (...) | 1,000 | 20 ms | 21 ms | 0.9x | 2 → 2 ms |
| Single table, WHERE x IN (...) | 10,000 | 94 ms | 78 ms | 1.2x | 2 → 2 ms |
| Single table, WHERE x IN (...) | 40,000 | 334 ms | 285 ms | 1.2x | 4 → 2
ms |
| Single table, WHERE x IN (...) | 100,000 | 650 ms | 542 ms | 1.2x | 10 → 3
ms |
| Single table, two lists on two columns | 1,000 | 39 ms | 37 ms | 1.0x | 7
→ 2 ms |
| Single table, two lists on two columns | 10,000 | 145 ms | 123 ms | 1.2x |
31 → 3 ms |
| Single table, two lists on two columns | 40,000 | 674 ms | 526 ms | 1.3x |
145 → 3 ms |
| Single table, two lists on two columns | 100,000 | 1,670 ms | 1,341 ms |
1.2x | 400 → 3 ms |
| CTE with the list, referenced twice | 1,000 | 30 ms | 34 ms | 0.9x | 4 → 4
ms |
| CTE with the list, referenced twice | 10,000 | 113 ms | 109 ms | 1.0x | 6
→ 3 ms |
| CTE with the list, referenced twice | 40,000 | 448 ms | 356 ms | 1.3x | 7
→ 4 ms |
| CTE with the list, referenced twice | 100,000 | 966 ms | 815 ms | 1.2x | 9
→ 4 ms |
#### Prod shape with the V2 physical optimizer
| Query | N | Off | On | Speed-up | Optimize off → on |
|---|---|---|---|---|---|
| Fact + 2 LEFT JOINs, 3 IN filters, GROUP BY 16 keys, FILTER aggs, ORDER BY
LIMIT | 40,000 | 1,206 ms | 427 ms | 2.8x | 811 → 10 ms |
| Fact + 2 LEFT JOINs, 3 IN filters, GROUP BY 16 keys, FILTER aggs, ORDER BY
LIMIT | 100,000 | 2,896 ms | 852 ms | 3.4x | 2,065 → 10 ms |
#### Positions that are cubic today
| Query | N | Off | On | Speed-up |
|---|---|---|---|---|
| JOIN ... ON a.k = b.k AND b.x IN (...) | 100 | 50 ms | 18 ms | 2.8x |
| JOIN ... ON a.k = b.k AND b.x IN (...) | 200 | 294 ms | 21 ms | 14.2x |
| JOIN ... ON a.k = b.k AND b.x IN (...) | 400 | 2,505 ms | 23 ms | 108.9x |
| JOIN ... ON a.k = b.k AND b.x IN (...) | 800 | 24.3 s | 43 ms | 570.8x |
| LEFT JOIN ... ON ... AND a.x IN (...) | 100 | 48 ms | 16 ms | 3.0x |
| LEFT JOIN ... ON ... AND a.x IN (...) | 200 | 296 ms | 21 ms | 14.4x |
| LEFT JOIN ... ON ... AND a.x IN (...) | 400 | 2,413 ms | 23 ms | 104.5x |
| SUM(CASE WHEN x IN (...) ...) | 100 | 46 ms | 14 ms | 3.2x |
| SUM(CASE WHEN x IN (...) ...) | 200 | 277 ms | 18 ms | 15.1x |
| SUM(CASE WHEN x IN (...) ...) | 400 | 2,434 ms | 19 ms | 127.4x |
| SUM(CASE WHEN x IN (...) ...) | 800 | 21.4 s | 35 ms | 607.1x |
| COUNT(*) FILTER (WHERE x IN (...)) | 100 | 45 ms | 13 ms | 3.4x |
| COUNT(*) FILTER (WHERE x IN (...)) | 200 | 279 ms | 17 ms | 16.8x |
| COUNT(*) FILTER (WHERE x IN (...)) | 400 | 2,437 ms | 18 ms | 131.7x |
| SELECT x IN (...) | 100 | 42 ms | 11 ms | 3.7x |
| SELECT x IN (...) | 200 | 277 ms | 15 ms | 18.1x |
| SELECT x IN (...) | 400 | 2,443 ms | 18 ms | 135.7x |
| GROUP BY CASE WHEN x IN (...) ... | 100 | 48 ms | 16 ms | 3.1x |
| GROUP BY CASE WHEN x IN (...) ... | 200 | 274 ms | 18 ms | 14.9x |
| GROUP BY CASE WHEN x IN (...) ... | 400 | 2,782 ms | 20 ms | 140.5x |
| OVER (PARTITION BY CASE WHEN x IN (...) ...) | 100 | 46 ms | 14 ms | 3.2x |
| OVER (PARTITION BY CASE WHEN x IN (...) ...) | 200 | 282 ms | 18 ms |
15.8x |
| OVER (PARTITION BY CASE WHEN x IN (...) ...) | 400 | 2,361 ms | 21 ms |
111.4x |
| HAVING SUM(CASE WHEN x IN (...) ...) | 100 | 46 ms | 15 ms | 3.2x |
| HAVING SUM(CASE WHEN x IN (...) ...) | 200 | 274 ms | 19 ms | 14.3x |
| HAVING SUM(CASE WHEN x IN (...) ...) | 400 | 2,455 ms | 22 ms | 109.1x |
#### Cubic positions with sealing, large N
| Query | 10,000 | 100,000 |
|---|---|---|
| JOIN ... ON a.k = b.k AND b.x IN (...) | 106 ms | 852 ms |
| LEFT JOIN ... ON ... AND a.x IN (...) | 114 ms | 901 ms |
| SUM(CASE WHEN x IN (...) ...) | 73 ms | 566 ms |
| COUNT(*) FILTER (WHERE x IN (...)) | 82 ms | 595 ms |
| SELECT x IN (...) | 82 ms | 591 ms |
| GROUP BY CASE WHEN x IN (...) ... | 87 ms | 723 ms |
| OVER (PARTITION BY CASE WHEN x IN (...) ...) | 79 ms | 598 ms |
| HAVING SUM(CASE WHEN x IN (...) ...) | 86 ms | 616 ms |
#### Stage size for IN plus a range on the same column (10k values)
| Query | Off: plan / stage bytes | On: plan / stage bytes |
|---|---|---|
| NOT IN (...) AND x > 0 | 103 ms / 98,047 | 106 ms / 98,047 |
| IN (...) OR x > c | 85 ms / 98,046 | 93 ms / 98,046 |
#### Phase split, prod shape at 100k values (ms)
| Sealing | parse | validate | sql2rel | optimize | plan | serialize | total
|
|---|---|---|---|---|---|---|---|
| Off | 147 | 66 | 606 | 2,514 | 11 | 19 | 3,363 |
| On | 139 | 63 | 552 | 15 | 9 | 18 | 795 |
</details>
## Tests
- `SearchSealerTest` (363 cases): seal and unseal round trips for 8 types ×
4 Sarg shapes × 3 null modes. Type and `Strong` answers equal to `SEARCH`.
Threshold, operator identity and literal operands.
- `SealedInListPlanningTest` (274 cases): 60 query shapes × V1 and V2 give
the same EXPLAIN, implementation plan, serialized stages and worker segments
with and without sealing, and EXPLAIN never shows a sealed call. 11 more shapes
on nullable columns (with and without null handling, V1 and V2) cover null
checks, `NOT IN` with NULL and outer joins. A guard rule in every planner phase
fails if any rule sees a `SEARCH` or a comparison list with 20 or more values,
and a negative test shows that the guard fires with sealing off. Cubic
positions plan 5,000 values within a time limit.
- `LargeInLists.json` (43 H2-checked queries): WHERE, NOT IN, IN plus a
range, CASE, FILTER, SELECT list, GROUP BY, HAVING, JOIN ON, LEFT JOIN, CTE,
UNION ALL, window, NULL in the list, nullable columns, and mixed Sargs in
intermediate stages.
- `RexExpressionUtilsTest`: mixed-Sarg conversion (IN and NOT_IN forms,
bounded complements, both null modes, small lists, no points, BIG_DECIMAL).
`QueryOptionsUtilsTest`: option parsing.
- Robustness run with every IN list sealed
(`DEFAULT_SEALED_IN_LIST_THRESHOLD = 1`): 2,593 planner tests and 3,922
H2-checked runtime queries pass, with 0 failures. The same suites pass at the
default threshold.
- An adversarial corpus of 458 H2-checked queries (not committed) gives
identical results at T = 0, 1 and 20 for both planners, with and without null
handling.
## Follow-ups (not in this PR)
- Many small IN lists, or a long `x = 1 OR x = 2 ...` chain, OR-ed inside
CASE or JOIN ON still take the cubic SqlToRel path, because no single list
reaches T. The same mechanism can cover this by marking the whole OR tree.
- Runtime: `FilterOperand.In` scans the list for every row, and it compares
BIG_DECIMAL values with `equals`, which gives wrong results in intermediate
stages.
- Planning is not cancellable: Pinot supplies no Calcite `CancelFlag`, so a
timed-out plan keeps its planner thread.
- Upstream Calcite: the `RexSimplify` patch for CALCITE-7302, and bug
reports for the `expand=true` issues above.
## References
- Pinot: #10168, #13605, #13614, #13617, #14615, #15027, #18554, #18720,
#17481
- Calcite: CALCITE-3178, CALCITE-5036, CALCITE-6467, CALCITE-6599,
CALCITE-7226, CALCITE-7232, CALCITE-7302
- Prior art: Druid rewrites IN lists of more than 100 literals before
SqlToRel. Hive uses an opaque `HiveIn` above 16 values. Spark uses `InSet`
above 10 values.
--
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]