apache / apache/druid

Large numeric IN predicates cause slow SQL planning on Druid 27

Open
#20,326 5 comments 0 reactions 0 assignees View on GitHub
Area - SQL Performance
Dominant language
Java
Stars
14.1k
Forks
3.8k
Avg merge
2d 58m
Merged PRs (30d)
233

Description

### Affected version

Druid 27 with Calcite 1.21.0.

### Problem

A numeric `IN` predicate containing 11,481 literals causes a severe SQL planning slowdown:

```sql
maxLongUniform IN (11481 literal values)
```

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 SQL `IN` into one equality per literal joined by `OR`:

```text
SQL IN
-> 11,481 equality OR terms
-> RexSimplify.simplifyOrTerms
-> native IN filter optimization
```

Druid 27's legacy null-replacement mode exposes numeric columns as non-nullable to Calcite. This makes `simplifyOrTerms` accumulate prior equality terms as predicates while simplifying later terms, which scales poorly for thousands of values. Druid eventually combines the filters into a native `InDimFilter`, 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](https://issues.apache.org/jira/browse/CALCITE-3178).

### Version comparison

| Druid version | Calcite | Planning time |
| --- | --- | ---: |
| Druid 27 | 1.21.0 | approximately 7,160 ms |
| Druid 32.0.0 | 1.37.0 | 25.322 +/- 5.677 ms |
| Current master | 1.42.0 | 27.814 +/- 6.485 ms |

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 with `inSubQueryThreshold = 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

- On Druid 27, set `inSubQueryThreshold` below the literal count so Calcite uses an inline `VALUES` relation instead of generating the large OR expression; validate the resulting execution plan.
- Upgrade to Druid 32 or later for the long-term resolution of this ordinary-datasource case.
- Retain a large-`IN` planning benchmark as a regression guard.

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.

Contributor guide

Open the contributing guide

Research direction

Start by reproducing the planning-only JMH scenario described in the issue with 11,481 numeric IN literals and compare the Druid 27 and later Calcite behavior. Review the suggested inSubQueryThreshold workaround and the large-IN regression benchmark; done means validating a practical mitigation or documenting a regression guard without reintroducing the Druid 32 behavior.

Written by the indexing model from the issue text.

Assessment

Tech stack
java, sql
Domain
databases
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Active
Clarity
Needs clarification
Newbie friendliness
35/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.