Skip to content

DRILL-8555: Physical plan cache for parameterized SQL queries - #3086

Open
letian-jiang wants to merge 4 commits into
apache:masterfrom
letian-jiang:physical-plan-cache
Open

letian-jiang wants to merge 4 commits into
apache:masterfrom
letian-jiang:physical-plan-cache

Conversation

@letian-jiang

Copy link
Copy Markdown
Contributor

DRILL-8555: Physical plan cache for parameterized SQL queries

Description

This PR adds a Drillbit-scoped physical plan cache with HBase and Iceberg support, reusing plans across connections to reduce repeated validation and optimization. It is disabled by default; enable it with ALTER SESSION SET planner.enable_plan_cache = true.

Method

flowchart LR
    SQL[SQL] --> Template[SQL template]
    Template --> Cache{Plan cache}
    Cache -->|Hit| Bind[Bind literals]
    Cache -->|Miss| Planner[Plan query]
    Bind --> Plan[Physical plan]
    Planner --> Plan
Loading

Eligible literals become parameter slots in the SQL template. A hit binds current values to a fresh copy of the cached plan; a miss plans the query normally and populates the cache after successful execution.

Safety guarantees

  • Only supported read queries are cached. Volatile or query-context functions and unsupported scans bypass caching. Structural literals and function configuration arguments stay in the cache key.
  • Reuse requires matching effective options, plugin configurations and table compatibility versions. Binding checks parameter types and numeric ranges; compatibility or reconstruction failures fall back to normal planning.
  • Cached plans are immutable, and each execution gets a fresh operator graph. Plugins opt in explicitly and must rebuild scan state from current parameters and metadata while preserving residual filters.

Benchmark

Measured on one local Drillbit with a Ryzen 7 9700X, 30 GiB RAM and OpenJDK 21. HBase used a 1,000-row mini-cluster for point reads, 50-row range scans and column filters. Iceberg ran all 22 TPC-H queries over eight SF0.01 tables (Q15 used a derived table; Q19 exposed the common equijoin).

Cache hits

Workload Planning off → hit (reduction) End-to-end off → hit (reduction)
HBase point read 50 → 14 ms (72.0%) 66 → 29 ms (56.1%)
HBase range scan 40 → 12 ms (70.0%) 55 → 26 ms (52.7%)
HBase column filter 32 → 10 ms (68.8%) 44 → 22 ms (50.0%)
Iceberg TPC-H SF0.01 2,161 → 355 ms (83.6%) 18,831 → 16,906 ms (10.2%)

The benefit is largest when planning dominates latency: it accounts for roughly 73–76% of the reported HBase baseline latency, and hits reduce end-to-end latency by 50–56%. Analytical queries also benefit: Iceberg planning drops 83.6%, reducing aggregate end-to-end latency by 10.2%.

HBase values are medians of three run medians (15 pairs per workload per run). Iceberg values are sums of per-query medians (three pairs per query), not suite wall time. Percentages are latency reductions relative to cache-off execution.

Cache misses

A separate warmed comparison cleared the plan cache before each cache-on query and drained background writes before both modes. Misses added a median 3–5 ms of paired end-to-end latency for HBase (45 pairs per workload). For Iceberg, the sum of query end-to-end medians changed from 20,624 to 21,806 ms (+5.7%). Miss overhead was modest in these local measurements, while hits provided the largest benefit for short queries.

Documentation

  • PLAN_CACHE_DESIGN.md: basic principles and supported scope.
  • PLAN_CACHE_PLUGIN_GUIDE.md: plugin APIs and scan reconstruction requirements.

@cgivre

cgivre commented Oct 2, 2026

Copy link
Copy Markdown
Contributor

@shfshihuafeng Thank you for this PR. Did you take a look at #3023? I realize they are somewhat different, but it seems the overall goal is similar in trying to cache plans. Are there components of that PR that you could incorporate in yours?

@letian-jiang

Copy link
Copy Markdown
Contributor Author

@cgivre Thanks for pointing me to #3023. I have reviewed it, and I agree that both PRs share the same underlying goal: reusing planning work to reduce the overhead of repeated queries.

Before considering #3023 production-ready, I think we would need to establish the conditions for safe plan reuse, even for identical SQL. For example:

  1. Time-dependent or non-deterministic expressions, such as date/time and random functions, need special handling when expressions are evaluated or folded during planning.
  2. Changes to options that affect planning may invalidate a cached plan.
  3. Physical plans contain mutable state, such as node assignments, which should be isolated between executions.
  4. Changes to the underlying data files may require rebuilding scan state.
  5. Changes to the underlying table schema may invalidate the plan.

There is also a practical consideration: production workloads often repeat the same query shape with different literals. Restricting reuse to identical SQL text would substantially limit the benefit for those workloads.

My implementation focuses on these correctness requirements through eligibility checks, context and table compatibility checks, and reconstruction of a fresh plan and scan state for each execution. It also preserves typed parameter slots so queries with different literals can reuse compatible cached plans.

@cgivre cgivre left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

A few issues with parameterized planning and the BIGINT literal change; details inline.

}
snapshot = snapshot.withOptionsFingerprint(optionsFingerprint);
}
PhysicalPlan planned = handler.getPlan(candidate.sql);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This plans from the parameterized SQL even when snapshot == null, i.e. for queries that can never be cached. Two consequences once planner.enable_plan_cache=true:

  1. JDBC: JdbcExpressionCheck.visitDynamicParam returns true, so WHERE id = ? is pushed down and JdbcBatchReader executes it with no bound value. SELECT * FROM mysql.db.t WHERE id = 5 fails at execution time, past the fallback.
  2. Plugins that require literals lose pushdown: FindPartitionConditions returns NO_PUSH (no dir0 pruning), Mongo, etc. These queries pay the cost with no cache benefit.

Fix: only plan from candidate.sql when snapshot != null and all tables are cacheable; otherwise use handler.getPlan(sqlNode).

templateTextPlan, cacheReader, cacheContext));
}
return planned;
} catch (Exception e) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This also swallows genuine planning failures from handler.getPlan(candidate.sql) (validation errors, planner timeouts), so failing queries are planned twice. The retry reuses sqlNode subtrees that the first validation already rewrote in place.

Fix: catch only cache-specific failures here (or rethrow ValidationException/UserException/planner timeouts).

}
if (lExpr.getDynamicParamIndex() < 0) {
// Small BIGINT values otherwise parse back as INT during a plan round trip.
sb.append("cast(").append(lExpr.getLong()).append(" as BIGINT)");

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This applies to all plans, not just cached ones: after a JSON round trip, BIGINT literals come back as CastExpression. DrillExprToPaimonTranslator has no visitCastExpression, so in multi-fragment plans the Paimon predicate is dropped after Drill has already removed the Filter, and bigint_col = 5 returns unfiltered rows. Only the Iceberg translator was updated.

Fix: either unwrap cast-of-literal in the pushdown translators (Paimon and any others), or have the parser read the round-tripped value back as a BIGINT literal instead of emitting a cast.

This branch has not been deployed

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

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants