Skip to content

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

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

letian-jiang wants to merge 16 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.

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).

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Agreed. I moved handler.getPlan(...) outside the cache fallback block, so cache lookup or binding failures can fall back before planning, while planning failures propagate without a cache-induced retry.

}
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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Thanks for identifying the impact beyond cached plans. I addressed this in the expression parser: the BIGINT cast emitted when serializing an ordinary integer literal is now restored as a LongExpression during deserialization. Explicit casts involving parameter slots remain intact.

@letian-jiang
letian-jiang requested a review from cgivre October 6, 2026 08:18
@letian-jiang

Copy link
Copy Markdown
Contributor Author

Hi @cgivre, thanks again for your review! I’ve pushed updates addressing your feedback on parameterized planning and BIGINT literal round trips. When you have a chance, could you take another look? Please let me know if there’s anything else you’d like me to address. Thanks!

@cgivre

cgivre commented Oct 8, 2026

Copy link
Copy Markdown
Contributor

Suggestion: make the plan cache size and expiration configurable

PlanCache hard-codes its limits (PlanCache.java:98-102): a 32 MB maximum weight and expireAfterWrite(10, MINUTES). Operators can't tune either one. A deployment with many distinct query templates may want a larger cache. One where tables change often, or where memory on the planning Drillbit is tight, may want a smaller cache or a shorter lifetime.

Could these be boot options in drill-module.conf, read in DrillbitContext and passed to the PlanCache constructor? The generated-code cache already does this with drill.exec.compile.cache_max_size. For example:

drill.exec.plan_cache: {
  # Upper bound on the total size of cached plan JSON, per Drillbit
  max_size_bytes: 33554432,
  # How long an entry may live after it is written; 0 disables time-based expiry
  expire_after_write: 10m,
  # Optional: evict entries that have not been used for this long; 0 disables
  expire_after_access: 0
}

A few related points:

  • The weigher uses value.json.length(), which counts characters, not bytes. If the setting is named in bytes, the weigher should measure bytes (or the setting should say characters).
  • expireAfterAccess would let rarely used templates age out sooner, while a hot query isn't replanned every 10 minutes just because it was written 10 minutes ago.
  • It would help to log the settings in use when the Drillbit starts, and to document them in the plan-cache design doc.
  • Startup settings are a good fit because the cache is created once per Drillbit, and Guava can't change these limits on a cache that's already built. planner.enable_plan_cache stays the runtime on/off switch.

@cgivre

cgivre commented Oct 8, 2026

Copy link
Copy Markdown
Contributor

Question: plans to support the other storage plugins?

Right now only HBase and Iceberg opt in (PLAN_CACHE_DESIGN.md:75). What are the plans for expanding this to the rest of the plugins? I'd like to see it cover as many as possible. The most heavily used ones (Parquet and the other dfs formats, and JDBC) are the ones that would benefit most, and they're currently excluded.

I understand why they can't simply be switched on. On a hit, the cache rebinds literals in Drill expressions. But most plugins turn those literals into native state during planning and keep only the result:

Plugin Value-dependent state built at planning time
Parquet Row-group / file pruning based on filter values
All dfs formats Directory (partition) pruning, and a fixed file list
JDBC Pushed-down SQL text with literals inlined
Mongo, Elasticsearch, Splunk, HTTP Native query / request built from the filter
Kafka Offset / timestamp ranges

Reusing those plans without rebuilding that state would silently return wrong results. So each plugin needs to keep the source expression in its scan JSON and redo pushdown/pruning on bind, as PLAN_CACHE_PLUGIN_GUIDE.md describes.

A possible path, roughly in order of effort:

  1. dfs formats without filter pushdown (CSV/TSV, JSON, Avro, etc.): filtering stays in Drill's Filter operator, so the plugin only needs to re-list files and redo directory pruning on bind, plus a table version that catches schema changes.
  2. Parquet: redo row-group/partition pruning for the new values, reusing the existing Parquet metadata cache to keep it cheap.
  3. JDBC and other pushdown plugins: a shared "rebuild pushdown on bind" hook in the plugin API, so each plugin doesn't reinvent it.

It would be worth measuring planning time against bind-time rebuild cost for dfs/Parquet, since re-listing and re-pruning eat into what a hit saves.

Given how large this PR already is, I'm fine with plugin expansion coming in follow-ups. But I'd like to understand the roadmap, and to make sure the plugin API in this PR makes opting in as cheap as possible. Could you add the plan (or JIRAs for each plugin) to the design doc?

@cgivre

cgivre commented Oct 8, 2026

Copy link
Copy Markdown
Contributor

Review findings

Line numbers refer to the PR head. The first three are correctness issues I'd like to see fixed before merge.

Correctness

  1. No fallback when planning the parameterized SQL fails (DrillSqlWorker.java:354). On a miss, the query is planned once from candidate.sql. The try/catch covers lookup and preparation but not handler.getPlan(planningSql). Any query that only plans with real literals then fails, but only when the cache is enabled. Examples: JOIN ... ON TRUE (becomes ON ?), PERCENTILE_CONT(0.5) WITHIN GROUP, and functions needing a literal operand that aren't on the hard-coded list. Please re-plan from the original literal SqlNode on failure. The comment about validation mutating the tree suggests keeping an unmodified copy (or re-parsing) for the fallback.
  2. HAVING/QUALIFY are parameterized while GROUP BY and the select list are not (PlanCacheParameterizer.java:130). GROUP BY x || 'a' HAVING x || 'a' = 'ba' becomes HAVING x || ? = ?, which no longer matches the grouped expression. Calcite then rejects it with "Expression 'x' is not being grouped", and because of Mavenized query-parse subproject #1 the query fails.
  3. Global change to BIGINT literal serialization (ExpressionStringBuilder.java:191). visitLongConstant now always emits cast(N as BIGINT), even with caching off. That changes EXPLAIN output, scan digests and filter strings (e.g. AbstractGroupScanWithMetadata.getFilterString) for everyone, and relies on a special case in ExprParser.createCast to undo it. Could this be limited to bound parameters?

Performance / design

  1. Context mismatch evicts instead of coexisting (DrillSqlWorker.java:334). The options/config fingerprint is checked inside the entry rather than included in the key. Two sessions running the same template with different options (e.g. planner.slice_target) invalidate each other's entry on every run, so the hit rate is 0 and each query also pays to republish. Including the fingerprints in the key fixes this.
  2. Remote metadata I/O on every lookup, including hits (HBaseStoragePlugin.java:77, IcebergFormatPlugin.java:124). HBase opens an Admin and makes two master RPCs. Iceberg runs HadoopTables.load, and then loads the table again in IcebergGroupScan during bind. A hit can cost more than the planning it skips. Benchmarks of hit latency against planning time would help.
  3. Entries are published without checking they can be rebound (PlanCache.java:187). put() only validates that the JSON reads back. An entry whose bind always fails is published, fails on every hit, is invalidated, gets replanned and is republished, forever. A trial bind before publishing would catch this.
  4. The writer queue holds live plan graphs (PlanCache.java:162). Up to 64 full PhysicalPlan object graphs (region maps, Iceberg task lists, plugin references) stay alive until serialized. They're also serialized after execution, which may have changed them. Serializing to a string right after planning and queuing only the string avoids both.
  5. Hard-coded allowlist of literal-only function operands (PlanCacheParameterizer.java:142). Any function not on the list (format strings, scale arguments, new UDFs) gets a slot and fails planning because of Mavenized query-parse subproject #1. Could this come from operator operand metadata instead?
  6. On a miss, the executed plan loses literal-based optimization (DrillSqlWorker.java:345). WHERE 1 = 1 and flag = TRUE become ? = ? / flag = ?, so constant reduction and selectivity estimates are lost for that run. This happens even if the plan is never published.

Observability and operations

  1. No logging of cache activity. Hits, misses, evictions and context-mismatch invalidations aren't logged; only some failures are logged at debug level. PlanCache.getHitCount() exists but nothing outside tests reads it. Please add debug logging for hits, misses and invalidation reasons, and expose hit/miss/eviction counts through Drill metrics or a sys table.
  2. No way to flush the cache. Entries are only removed on expiry or when found stale on lookup. Turning off planner.enable_plan_cache doesn't clear them, and there's no admin command or REST endpoint. (Configurable size and expiration are covered in my earlier comment.)

Docs

  1. PLAN_CACHE_DESIGN.md and PLAN_CACHE_PLUGIN_GUIDE.md are at the repo root. Developer docs belong under docs/dev/ and should be linked from docs/dev/DevDocs.md.

@letian-jiang

Copy link
Copy Markdown
Contributor Author

@cgivre Thank you for the detailed and thoughtful review! I’ve learned a lot from your comments.

@letian-jiang

Copy link
Copy Markdown
Contributor Author

Suggestion: make the plan cache size and expiration configurable

PlanCache hard-codes its limits (PlanCache.java:98-102): a 32 MB maximum weight and expireAfterWrite(10, MINUTES). Operators can't tune either one. A deployment with many distinct query templates may want a larger cache. One where tables change often, or where memory on the planning Drillbit is tight, may want a smaller cache or a shorter lifetime.

Could these be boot options in drill-module.conf, read in DrillbitContext and passed to the PlanCache constructor? The generated-code cache already does this with drill.exec.compile.cache_max_size. For example:

drill.exec.plan_cache: {
  # Upper bound on the total size of cached plan JSON, per Drillbit
  max_size_bytes: 33554432,
  # How long an entry may live after it is written; 0 disables time-based expiry
  expire_after_write: 10m,
  # Optional: evict entries that have not been used for this long; 0 disables
  expire_after_access: 0
}

A few related points:

  • The weigher uses value.json.length(), which counts characters, not bytes. If the setting is named in bytes, the weigher should measure bytes (or the setting should say characters).
  • expireAfterAccess would let rarely used templates age out sooner, while a hot query isn't replanned every 10 minutes just because it was written 10 minutes ago.
  • It would help to log the settings in use when the Drillbit starts, and to document them in the plan-cache design doc.
  • Startup settings are a good fit because the cache is created once per Drillbit, and Guava can't change these limits on a cache that's already built. planner.enable_plan_cache stays the runtime on/off switch.

I’ve made the cache size and both expiration policies configurable.

  1. The defaults are 32 MiB, expire_after_access = 10m, and expire_after_write = 0 (disabled). This keeps actively used plans available while allowing idle entries to expire. Since either expiration condition can expire an entry, retaining a 10-minute write expiry would still force hot plans to expire regardless of access. Operators who want periodic re-optimization can enable write-based expiration.
  2. The size limit now counts the combined UTF-8 byte lengths of the cache key, physical-plan JSON, and optional text plan, rather than Java character counts. This bounds cached text content; it does not include context metadata, Java object overhead, or Guava’s internal structures, so it is not an exact heap-memory limit.

@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.

Line-by-line follow-up to my earlier review comment, with suggested changes where there's a concrete fix.

I applied all the suggestions together on top of 53306c1, and exec/java-exec (with its dependencies, including logical) builds cleanly, checkstyle included. Several depend on each other, so it's easiest to add them to a batch and commit them together:

  • the PlanCache constructor, its three imports, and new PlanCache(config) in DrillbitContext
  • writeAfterSuccess / put in PlanCache and the serialization change in DrillSqlWorker
  • keyFingerprint() in PlanCache and the key in DrillSqlWorker

The drill-module.conf defaults are in the comment on the PlanCache fields, because the right spot in that file is outside the diff.

Tests: I couldn't find any plan-cache tests in this PR. PLAN_CACHE_PLUGIN_GUIDE.md describes the checks a plugin should pass (results with caching off vs. a hit with different literals, schema changes, plugin config changes, joins with unsupported plugins, etc.), but none are included for HBase, Iceberg or the engine. Could you add them? They should cover at least:

  • the fallback path
  • HAVING with GROUP BY
  • two sessions with different options sharing a template
  • a stale table version
  • the new boot options

Comment on lines +352 to +358
// Plan exactly once, outside cache fallback. Validation can mutate the SQL tree,
// so retrying a failed planning attempt could reuse partially rewritten nodes.
PhysicalPlan planned = handler.getPlan(planningSql);
if (prepareCacheInsert != null) {
prepareCacheInsert.accept(planned);
}
return planned;

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.

[Correctness] Fall back to ordinary planning when the parameterized SQL fails to plan.

On a miss, the query is planned once from candidate.sql, and the try/catch above only covers lookup and preparation. Any query that only plans with real literals fails, but only when the cache is enabled. Examples: JOIN ... ON TRUE (becomes ON ?), PERCENTILE_CONT(0.5) WITHIN GROUP, and functions that need a literal operand but aren't in configurationOperand.

The comment above is right that the mutated tree can't be retried. This suggestion avoids that by re-entering getQueryPlan from the SQL string with planner.enable_plan_cache set to false at query level, the same pattern the Metastore retry in getPlan uses. The fresh parse, converter and handler start clean. A query that fails for a real reason (bad column, etc.) fails again on the second attempt and reports the error from the original SQL.

Suggested change
// Plan exactly once, outside cache fallback. Validation can mutate the SQL tree,
// so retrying a failed planning attempt could reuse partially rewritten nodes.
PhysicalPlan planned = handler.getPlan(planningSql);
if (prepareCacheInsert != null) {
prepareCacheInsert.accept(planned);
}
return planned;
if (prepareCacheInsert == null) {
return handler.getPlan(planningSql);
}
PhysicalPlan planned;
try {
planned = handler.getPlan(planningSql);
} catch (Exception e) {
// Some queries only plan with real literal values (e.g. ON TRUE, or functions that
// need a literal operand). Validation mutates the SQL tree, so re-parse and plan the
// original SQL with the cache off for this query, like the Metastore retry above.
logger.debug("Parameterized planning failed; replanning without the plan cache", e);
context.getOptions().setLocalOption(PlannerSettings.ENABLE_PLAN_CACHE_OPTION, false);
return getQueryPlan(context, sql, textPlan);
}
prepareCacheInsert.accept(planned);
return planned;

Comment on lines +130 to +131
select.getGroup(), visitNullable(select.getHaving()),
select.getWindowList(), visitNullable(select.getQualify()),

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.

[Correctness] Keep HAVING and QUALIFY verbatim when GROUP BY / the select list are.

When structuralProjection is true, GROUP BY and the select list keep their literals, but HAVING and QUALIFY still get slots. So they no longer match the grouped or windowed expressions:

SELECT x || 'a', COUNT(*) FROM hbase.t GROUP BY x || 'a' HAVING x || 'a' = 'ba'

HAVING becomes x || ? = ?, and Calcite raises "Expression 'x' is not being grouped". Without a fallback, the query fails only when the cache is on. Using the same condition for HAVING and QUALIFY keeps them consistent with what they have to match.

Suggested change
select.getGroup(), visitNullable(select.getHaving()),
select.getWindowList(), visitNullable(select.getQualify()),
select.getGroup(),
structuralProjection ? select.getHaving() : visitNullable(select.getHaving()),
select.getWindowList(),
structuralProjection ? select.getQualify() : visitNullable(select.getQualify()),

Comment on lines +189 to +194
if (!lExpr.isDynamicParam()) {
// Preserve integer width in JSON; the parser restores this as a BIGINT literal.
sb.append("cast(").append(lExpr.getLong()).append(" as BIGINT)");
} else {
sb.append(lExpr.getLong());
}

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.

[Correctness / compatibility] Narrow the BIGINT cast to values that need it.

This now writes cast(N as BIGINT) for every non-parameter BIGINT literal in every plan, even with the cache off. That changes EXPLAIN output, scan digests and filter strings (e.g. AbstractGroupScanWithMetadata.getFilterString) for all users.

Only values in int range are ambiguous on read-back, because ValueExpressions.getNumericExpression tries Integer.parseInt first and then Long.parseLong. Values outside int range already round-trip as LongExpression, so they don't need the cast. This suggestion limits the cast to the ambiguous case. The createCast special case in ExprParser.g4 still handles it.

Is the cast needed at all outside the plan cache? Physical plans were already serialized to JSON for remote fragments before this PR. If only cached plans need it, could it be applied only when serializing for the cache?

Suggested change
if (!lExpr.isDynamicParam()) {
// Preserve integer width in JSON; the parser restores this as a BIGINT literal.
sb.append("cast(").append(lExpr.getLong()).append(" as BIGINT)");
} else {
sb.append(lExpr.getLong());
}
long value = lExpr.getLong();
if (!lExpr.isDynamicParam() && value >= Integer.MIN_VALUE && value <= Integer.MAX_VALUE) {
// Only int-range values are ambiguous: the parser reads them back as INT.
sb.append("cast(").append(value).append(" as BIGINT)");
} else {
sb.append(value);
}

this.optionsFingerprint = Objects.requireNonNull(optionsFingerprint, "optionsFingerprint");
this.tableVersions = Collections.unmodifiableMap(tableVersions);
this.pluginConfigs = Collections.unmodifiableMap(pluginConfigs);
}

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.

[Performance] Put the options and plugin-config fingerprints in the key (1 of 2; the other suggestion is in DrillSqlWorker).

Today, a mismatch in matchesContext invalidates the shared entry. Two sessions running the same template with different options (e.g. ALTER SESSION SET planner.slice_target = ...) keep evicting each other's entry, so the hit rate is 0 and every query also pays to republish. With the fingerprints in the key, each option set gets its own entry. Only a table-version change, which should invalidate, falls through to invalidate.

Suggested change
}
}
/**
* Session- and plugin-dependent part of the context. It belongs in the cache key, so
* sessions with different options keep separate entries instead of evicting each other.
* Table versions stay in {@link #matches}, where a change should invalidate the entry.
*/
String keyFingerprint() {
return optionsFingerprint + '\n' + pluginConfigs;
}

Comment on lines +317 to +319
String key = context.getQueryUserName() + '\n'
+ context.getSession().getDefaultSchemaPath() + '\n'
+ candidate.template;

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.

[Performance] Include the session/plugin fingerprint in the key (2 of 2; see ContextSnapshot.keyFingerprint() in PlanCache). This keeps sessions with different option values from evicting each other's entries.

Suggested change
String key = context.getQueryUserName() + '\n'
+ context.getSession().getDefaultSchemaPath() + '\n'
+ candidate.template;
String key = context.getQueryUserName() + '\n'
+ context.getSession().getDefaultSchemaPath() + '\n'
+ snapshot.keyFingerprint() + '\n'
+ candidate.template;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import com.fasterxml.jackson.databind.JsonNode;

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.

Import for the plan-cache metrics.

Suggested change
import com.fasterxml.jackson.databind.JsonNode;
import com.codahale.metrics.Gauge;
import com.fasterxml.jackson.databind.JsonNode;

ExecConstants.STORAGE_PLUGIN_REGISTRY_IMPL, StoragePluginRegistry.class, this);

reader = new PhysicalPlanReader(config, classpathScan, lpPersistence, endpoint, storagePlugins);
planCache = new PlanCache();

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.

Pass the boot config so the cache size and expiry can be configured (see the suggestion on PlanCache's fields).

Suggested change
planCache = new PlanCache();
planCache = new PlanCache(config);

return null;
}
String identifier = ((HBaseScanSpec) selection).getTableName();
try (Admin admin = getConnection().getAdmin()) {

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.

[Performance] Remote metadata I/O on every lookup, including hits.

ContextSnapshot.resolve calls this before every lookup. Each call opens an Admin and makes two master RPCs (tableExists + getTableDescriptor). Iceberg is similar: IcebergFormatPlugin.planCacheTable runs HadoopTables.load, and IcebergGroupScan loads the table again during bind. A hit can end up costing more than the planning it skips.

Possible fixes:

  • Drop tableExists and call getDescriptor directly (it throws TableNotFoundException, which can map to null). That's one RPC instead of two.
  • Use Connection.getTable(name).getDescriptor() rather than opening an Admin per query.
  • Cache the descriptor hash briefly (a few seconds) per Drillbit, accepting a short window where a schema change is caught at bind time instead.

Could you include hit latency against full planning time for an HBase point lookup in the PR benchmarks?

// Keep this list aligned with the value-dependent rewrites in DrillOptiq
// and PreProcessLogicalRel; ordinary data operands still get slots.
switch (call.getOperator().getName().toUpperCase(Locale.ROOT)) {
case "DATE_PART":

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.

[Robustness] Hard-coded list of functions whose operands must stay literal.

Any function not on this list that needs a literal operand during validation, return-type inference or DrillOptiq conversion silently gets a slot and fails planning. That includes format strings, scale arguments, percentile fractions and new UDFs. The fallback suggested in DrillSqlWorker stops these from failing queries. It would still be better not to rely on the list: could literal-only operands come from the operator's operand type checker (e.g. SqlOperandTypeChecker / OperandTypes.LITERAL), or from an annotation on Drill UDFs, rather than a name switch?

@@ -0,0 +1,103 @@
<!--

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.

[Docs] Please move PLAN_CACHE_DESIGN.md and PLAN_CACHE_PLUGIN_GUIDE.md from the repo root to docs/dev/, and link them from docs/dev/DevDocs.md. It would also help to document the new boot options and metrics, and the plan for supporting other plugins (see my earlier comment).

@letian-jiang

Copy link
Copy Markdown
Contributor Author

Question: plans to support the other storage plugins?

Right now only HBase and Iceberg opt in (PLAN_CACHE_DESIGN.md:75). What are the plans for expanding this to the rest of the plugins? I'd like to see it cover as many as possible. The most heavily used ones (Parquet and the other dfs formats, and JDBC) are the ones that would benefit most, and they're currently excluded.

I understand why they can't simply be switched on. On a hit, the cache rebinds literals in Drill expressions. But most plugins turn those literals into native state during planning and keep only the result:

Plugin Value-dependent state built at planning time
Parquet Row-group / file pruning based on filter values
All dfs formats Directory (partition) pruning, and a fixed file list
JDBC Pushed-down SQL text with literals inlined
Mongo, Elasticsearch, Splunk, HTTP Native query / request built from the filter
Kafka Offset / timestamp ranges
Reusing those plans without rebuilding that state would silently return wrong results. So each plugin needs to keep the source expression in its scan JSON and redo pushdown/pruning on bind, as PLAN_CACHE_PLUGIN_GUIDE.md describes.

A possible path, roughly in order of effort:

  1. dfs formats without filter pushdown (CSV/TSV, JSON, Avro, etc.): filtering stays in Drill's Filter operator, so the plugin only needs to re-list files and redo directory pruning on bind, plus a table version that catches schema changes.
  2. Parquet: redo row-group/partition pruning for the new values, reusing the existing Parquet metadata cache to keep it cheap.
  3. JDBC and other pushdown plugins: a shared "rebuild pushdown on bind" hook in the plugin API, so each plugin doesn't reinvent it.

It would be worth measuring planning time against bind-time rebuild cost for dfs/Parquet, since re-listing and re-pruning eat into what a hit saves.

Given how large this PR already is, I'm fine with plugin expansion coming in follow-ups. But I'd like to understand the roadmap, and to make sure the plugin API in this PR makes opting in as cheap as possible. Could you add the plan (or JIRAs for each plugin) to the design doc?

  1. I chose HBase because I expect high-concurrency, short-query workloads to benefit most from plan caching: planning can account for a substantial share of their latency. Since Drill is also an analytical engine, I wanted to test a data lake format such as Iceberg to evaluate the benefit for analytical workloads as well.

  2. I agree that the progression from DFS formats without filter pushdown, to Parquet, and then JDBC is reasonable. I plan to add support for these commonly used plugins in follow-up PRs and document the roadmap in the design doc.

@letian-jiang

letian-jiang commented Oct 8, 2026 •

Copy link
Copy Markdown
Contributor Author
  1. No fallback when planning the parameterized SQL fails (DrillSqlWorker.java:354). On a miss, the query is planned once from candidate.sql. The try/catch covers lookup and preparation but not handler.getPlan(planningSql). Any query that only plans with real literals then fails, but only when the cache is enabled. Examples: JOIN ... ON TRUE (becomes ON ?), PERCENTILE_CONT(0.5) WITHIN GROUP, and functions needing a literal operand that aren't on the hard-coded list. Please re-plan from the original literal SqlNode on failure. The comment about validation mutating the tree suggests keeping an unmodified copy (or re-parsing) for the fallback.
  2. HAVING/QUALIFY are parameterized while GROUP BY and the select list are not (PlanCacheParameterizer.java:130). GROUP BY x || 'a' HAVING x || 'a' = 'ba' becomes HAVING x || ? = ?, which no longer matches the grouped expression. Calcite then rejects it with "Expression 'x' is not being grouped", and because of Mavenized query-parse subproject #1 the query fails.

I've updated the parameterizer to preserve HAVING and QUALIFY expressions whenever the SELECT expressions are preserved for grouping, ordering, or window definitions. This keeps matching expressions consistent across clauses. Eligible WHERE literals can still be rebound.
I don't think replanning the original SQL after a planning failure is a reliable general fallback. Exception types alone cannot distinguish failures introduced by dynamic parameters from errors in the original query. Also, parameterization can produce incorrect results without throwing an exception.
For example, before this fix:

SELECT x + 1 AS k,
       ROW_NUMBER() OVER (ORDER BY x + 1) AS rn
FROM (VALUES (1, 1), (1, 2), (2, 3)) AS t(x, v)
WHERE v > 0
QUALIFY x + 1 = 2
    AND ROW_NUMBER() OVER (ORDER BY x + 1) = 1
ORDER BY x + 1;

With plan caching disabled, this returned one row: (k = 2, rn = 1). With caching enabled, it returned zero rows without any exception, both on initial planning and on a confirmed cache hit. An exception-triggered fallback would not catch this regression.

@letian-jiang

Copy link
Copy Markdown
Contributor Author
  1. Hard-coded allowlist of literal-only function operands (PlanCacheParameterizer.java:142). Any function not on the list (format strings, scale arguments, new UDFs) gets a slot and fails planning because of Mavenized query-parse subproject #1. Could this come from operator operand metadata instead?

This is a difficult problem, and I agree that a hard-coded list carries a maintenance risk. Given Drill’s current architecture, though, I don’t see a clearly better practical approach: DrillOptiq itself uses function-name switches for value-dependent rewrites, so the parameterizer needs explicit knowledge of which operands must remain literals. Moving this information into function metadata could improve organization, but would still require manually declaring and maintaining those rules. For now, I think keeping this list aligned with the existing rewrites and adding targeted regression tests is a reasonable approach, while acknowledging that it does not establish safety for every SQL pattern.

@letian-jiang

Copy link
Copy Markdown
Contributor Author

Following up on the remaining points:

  • 3. BIGINT serialization: Ordinary BIGINT literals now retain their existing serialization format. Bound dynamic parameters preserve their type through the parameter wrapper. I also removed the parser special case that converted BIGINT casts back into literals, so explicit casts retain their normal behavior.

  • 4. Context fingerprints: Effective option and storage-plugin configuration fingerprints are now included in the cache key. Different configurations can coexist instead of repeatedly invalidating each other. Table versions remain compatibility checks so incompatible table changes invalidate the cached entry.

  • 5. Metadata I/O: The [benchmark section](DRILL-8555: Physical plan cache for parameterized SQL queries #3086 (comment)) already compares ordinary planning with confirmed cache hits, including planning time and end-to-end latency. For example, HBase point-read planning drops from 50 to 14 ms, and end-to-end latency from 66 to 29 ms. I’ve also applied the two suggested HBase improvements: removed the separate tableExists check and switched to Connection.getTable(name).getDescriptor(), avoiding an Admin handle per lookup. TableNotFoundException returns no cache metadata, while other I/O failures retain the existing fallback behavior.

  • 6. Trial binding: My intended contract is that parameter binding should be reliable for every eligible plan. Deterministic binding failures indicate a gap in the binder or eligibility checks that should be fixed there. A trial bind with the first execution’s values can catch some defects, but does not establish correctness for subsequent values. I have retained JSON read-back validation without adding a trial bind. Scan reconstruction can still fail because of metadata or I/O errors; the existing invalidation and replanning path handles those failures.

  • 7. Writer queue: The plan is now serialized immediately after planning, before parallelization or execution can mutate its operators. Pending publication and the writer queue retain the immutable JSON string instead of the live physical-plan graph. Publication still happens only after successful execution.

  • 9. Plan quality: I agree that this is an expected trade-off. The cached plan must remain general enough for subsequent parameter values, so constant folding based on those values cannot specialize it to one execution. Value-specific selectivity, join ordering and distribution choices may also be less effective, including on the initial miss. This limitation is now documented.

  • 10. Observability: Added DEBUG logging for hits, misses and invalidation reasons. The hits, misses, invalidations, evictions and entries gauges are registered in DrillMetrics under drill.plan_cache., using the existing JMX and log reporters.

  • 11. Cache clearing: Added ALTER SYSTEM CLEAR PLAN CACHE. It clears the receiving Drillbit’s cache and returns its address. With authentication enabled, it follows the existing SYSTEM administrator policy. Clearing also prevents plans prepared before the command, including queued or pending publications, from restoring old entries.

  • 12. Documentation: Both documents have been moved to docs/dev/ and linked from DevDocs.md. They now cover configuration, metrics, the clearing command, plan-quality trade-offs and the plugin expansion roadmap.

@letian-jiang

Copy link
Copy Markdown
Contributor Author

I’ve rerun both benchmarks on the latest commit, [6ee1c592b](letian-jiang@6ee1c59), using the same machine, datasets, queries and sampling method as the original measurements.

Workload Planning off → hit Reduction End-to-end off → hit Reduction
HBase point read 48 → 14 ms 70.8% 63 → 29 ms 54.0%
HBase range scan 40 → 13 ms 67.5% 54 → 26 ms 51.9%
HBase column filter 34 → 10 ms 70.6% 47 → 22 ms 53.2%
Iceberg TPC-H SF0.01 1,950 → 616 ms 68.4% 18,691 → 17,101 ms 8.5%

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