Large numeric IN predicates cause slow SQL planning #20326
Description
Activity
- changed the title
[-]Add a configurable timeout for SQL query planning on the Broker[/-][+]Large numeric IN predicates cause slow SQL planning on Druid 27[/+]on Sep 11, 2026 #16388 (Druid 31) was also relevant to speeding things up. It updated our planning to divert most INs and long OR chains into
SCALAR_IN_ARRAY, a function that has an optimized implementation.@gianm Thanks for the info.
@aho135 Based on above investigation, from Druid 32 where incompatible null-handling is removed, I think such problem should not exist. We suffered from this problem because most of our clusters are still running on Druid 27.
I'm wondering in which case do you still suffer from this problem that drives you to add the timeout setting to the planner ?Thanks for the additional references @gianm @FrankChen021 This part of the code is pretty foreign to me, still taking some time to get a better grasp on the issue
We are encountering this problem in Druid 36, so I'm wondering if the previous optimizations don't handle all query shapes. The only difference I'm seeing with the query we're running is that the large IN clause is wrapped in a CTE. Let me try running some benchmarks to see if this could effect the planning time
@aho135 Can you share the SQL snippet so we can also investigate it
@FrankChen021 The query itself has a simple pattern like:
WITH datapoints AS ( SELECT * FROM <datasource> WHERE dim1 IN (v1, v2, ... vn) )I ran the jmh benchmarks with some small query modifications to include the CTE but it did not have any material impact on planning time.
I did some more analysis around the degradation window and actually found that even some simple queries had extremely long planning times. Turns out the long planning times for those were actually just a side effect of stop the world GC pauses caused by a flurry of queries with large IN clauses. The planning timeout should help alleviate some of the GC pressure in this scenario, but curious if you have any other thoughts on this scenario
Investigation update: regression isolated to the Calcite 1.35 → 1.37 upgrade
I reproduced the exact
InPlanningBenchmark.queryStringFunctionInSqlworkload introduced by #16388:EXPLAIN PLAN FOR SELECT COUNT(*) FROM foo WHERE long1 = 8 OR LOWER(string1) IN ('1', '2', ..., '1000000')
Parameters:
inClauseLiteralsCount=1000000inSubQueryThreshold=2147483647rowsPerSegment=500000
Results across Druid’s Calcite upgrades:
Druid commit Calcite version Result #16388 merge commit 72432c2e781.35.0 10.39 s/op, 14.94 GB allocated/op ac19b148c2from #165041.37.0 No completed operation after 20 minutes 43ab654f851.41.0 No completed operation after 20 minutes 35e9e436a71.42.0 No completed operation after 20 minutes “No completed operation” means the benchmark was stopped after 20 minutes; it did not produce a JMH score.
Regression boundary
#16388 changed large
INpredicates to useSCALAR_IN_ARRAY. At its merge commit, using Calcite 1.35, the one-million-string workload completed in approximately 10 seconds.#16504 upgraded Calcite from 1.35 to 1.37. This brought in CALCITE-5948 via PR apache/calcite#3395, a Calcite 1.36 correctness fix that explicitly casts array elements when their types differ from the derived component type.
The implementation loops over the array operands and calls:
call.setOperand(i, castTo(operands.get(i), elementType));
SqlBasicCall.setOperandcreates a new copy of the complete immutable operand list for every replacement. With a large string array, many literals require adjustment to the common string type. This results in approximately quadratic copying.A JFR profile of the 100,000-string case on the later Calcite path attributed most allocation pressure to these copies:
ImmutableList.copyOf: 50.29%Platform.copy: 41.28%
The dominant stack was:
SqlArrayValueConstructor.inferReturnType -> SqlValidatorUtil.adjustTypeForArrayConstructor -> SqlValidatorUtil.adjustTypeForMultisetConstructor -> SqlBasicCall.setOperand -> ImmutableNullableList.copyOfRelease impact
Both #16388 and #16504 are included in Druid 31.0.0. Therefore, the fast combination of #16388 with Calcite 1.35 existed temporarily on master but was not included in a Druid release.
This is distinct from the original
simplifyOrTermsproblem addressed by #16388. It is a later Calcite array-validation regression, primarily affecting large stringINlists. Numeric lists generally do not require the same per-element string type adjustment.Suggested fix direction
The underlying fix should be made in Calcite, or carried as a narrow backport in Druid: compute all required casts first and replace or reconstruct the operand list once, instead of copying the full list for every array element.
Thanks for the analysis and putting up the Calcite patch @FrankChen021
Above fix in calcite side has been merged. Open a new PR ( #20362 ) to reduce allocation via 20362 which helps ease the GC although the improvement is not very big.
Larger memory allocation saving will be on Calcite side, not on Druid side.
#20367 further greatly reduces the memory allocation during planning which also reduces the planning time
- changed the title
[-]Large numeric IN predicates cause slow SQL planning on Druid 27[/-][+]Large numeric IN predicates cause slow SQL planning[/+]on Sep 27, 2026
Affected version
Druid 27 with Calcite 1.21.0.
Problem
A numeric
INpredicate containing 11,481 literals causes a severe SQL planning slowdown:A deterministic, planning-only JMH benchmark measured approximately 7,160 ms per plan on Druid 27. The delay occurs during SQL planning, before native query execution.
Root cause
For a literal list below
inSubQueryThreshold, Calcite converts SQLINinto one equality per literal joined byOR:Druid 27's legacy null-replacement mode exposes numeric columns as non-nullable to Calcite. This makes
simplifyOrTermsaccumulate prior equality terms as predicates while simplifying later terms, which scales poorly for thousands of values. Druid eventually combines the filters into a nativeInDimFilter, but only after Calcite has paid the simplification cost.The SQL-to-OR rewrite is intentional. The defect is the scaling of cross-term predicate simplification for a very large equality disjunction. This is related to #7904 and CALCITE-3178.
Version comparison
The exact Druid 32.0.0 result was verified using the official Calcite 1.37.0 artifact and the same
maxLongUniform IN (11,481 literals)query shape withinSubQueryThreshold = Integer.MAX_VALUE.Why Druid 32 and later are unaffected
Druid 32 removed legacy SQL-incompatible null handling in #17609. Ordinary numeric datasource columns are now exposed to Calcite as nullable. Calcite therefore skips the expensive predicate accumulation and can later combine the equality terms into
SEARCH/Sarg.The underlying Calcite weakness can still affect genuinely non-nullable schemas or plans where non-nullability is already established, but the original query against an ordinary Druid datasource does not reproduce the slowdown from Druid 32 onward.
Suggested actions
inSubQueryThresholdbelow the literal count so Calcite uses an inlineVALUESrelation instead of generating the large OR expression; validate the resulting execution plan.INplanning benchmark as a regression guard.Related Calcite allocation reductions
This is separate from the original Druid 27 numeric planning regression: these upstream changes reduce temporary allocation for large string
INpredicates on Calcite 1.42.0.The isolated reductions overlap and should not be added. With all five changes applied, allocation fell by 70.8% at 100,000 strings and 70.5% at 1,000,000 strings.
PR #20314 adds a separate planning-timeout guardrail that limits the Broker impact of pathological planning cases; it does not change the root cause described here.