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]

Reply via email to