When is a custom Catalyst rule the right fix instead of rewriting the Spark query?
answer
- last resort, not first tool
- a property of the platform, not the query
- every session pays for it
- a broken rule returns answers, not errors
- internal APIs move between minors
basics
~20 sOnly when the behaviour must apply to every query in the session and cannot live in user code — governance predicates, lineage capture, a new SQL grammar. For a single slow query, fix statistics, the predicate or the layout instead.
solid answer
~50 sSpark exposes `SparkSessionExtensions`, registered through the `spark.sql.extensions` config, with injection points for resolution rules, optimizer rules, planner strategies, parser extensions and functions. That machinery is how Delta Lake and Apache Iceberg add their own SQL grammar and resolution behaviour, so it is a legitimate tool — but the bar is a **platform-wide invariant**, not a slow query. Good reasons: injecting a mandatory security or partition predicate on every scan of governed tables, capturing lineage from the analyzed plan, rejecting cartesian products before they run. The costs are real: the rule runs on every query in the session, a semantics-breaking bug produces silently wrong results rather than an error, and it binds you to Catalyst's internal plan classes, which change across Spark minors. Exhaust the cheaper levers first — `ANALYZE TABLE`, a hint, a rewritten predicate, a view, or `spark.sql.optimizer.excludedRules` to disable one built-in rule.
code
scala · 6 linesclass GuardrailExtensions extends (SparkSessionExtensions => Unit) {
override def apply(ext: SparkSessionExtensions): Unit = {
ext.injectOptimizerRule { _ => RequirePartitionFilter }
ext.injectResolutionRule { spark => CaptureLineage(spark) }
}
}go deeper
You will not be asked to write one. Know only that Catalyst is rule-based and that connectors like Delta Lake and Iceberg plug into it, so the optimizer is extensible rather than fixed.
Be able to say that SparkSessionExtensions can inject analyzer, optimizer and planner rules, and that ordinary performance problems are solved with statistics, hints and data layout instead.
Judge the alternatives out loud and reject the rule for a single query. If you have shipped one, describe how it was tested against real plans and how it was disabled in a hurry.
Own it as a governance decision: what invariant justifies engine-level enforcement, who owns the rule across Spark upgrades, how differential testing protects correctness, and what the kill switch is.
## The mechanism Catalyst is deliberately extensible. `SparkSessionExtensions` offers injection points including: - `injectResolutionRule` and `injectPostHocResolutionRule` — analyzer rules, run while or after names are resolved. - `injectOptimizerRule` — a logical-plan rewrite rule added to the optimizer's batches. - `injectPlannerStrategy` — a strategy that can produce physical operators. - `injectParser` — extend or wrap the SQL parser, which is how new statements are added. - `injectFunction` — register a native Catalyst function rather than a UDF. - `injectColumnar` — supply columnar physical rules. An extension is a class implementing a function from `SparkSessionExtensions` to `Unit`, registered by setting `spark.sql.extensions` to its fully qualified name — usually in the cluster's default configuration so it applies to every session on the platform. Delta Lake and Apache Iceberg both ship such extensions to add their own SQL statements and resolution rules; this is not an exotic corner of the API. A rule itself is a function over plan trees: match a pattern, return a transformed tree. It runs inside the optimizer's fixed-point loop, so it must be *stable* (repeated application converges) and, above all, **semantics-preserving**. ## When it is the right call The test is whether the behaviour is a property of the **platform** rather than of a query. Cases that pass: - **Policy injection.** Every scan of a governed table must carry a tenant or region predicate, or a row-level filter derived from the caller's identity. Asking a hundred analysts to remember it is not a control; a resolution rule that adds it to the plan is. - **Lineage and governance capture.** Walking the analyzed logical plan gives exact table-and-column-level lineage for every query, with no parsing of SQL text and no cooperation from the author. - **Guardrails.** A planner-side check that fails a query containing `CartesianProduct`, or that refuses a scan of a huge table with no partition filter, turns an expensive incident into an immediate error message. - **A new SQL surface.** A parser extension adding statements a connector needs — the reason the lakehouse formats use this API. - **A rewrite that must survive everyone's queries.** Redirecting references from a retired table to its replacement during a migration, rather than chasing every job. ## When it is the wrong call For a single slow query, essentially always. The cheaper levers, roughly in order: 1. **Fix the estimate.** `ANALYZE TABLE t COMPUTE STATISTICS FOR COLUMNS ...` so the planner works from real numbers rather than compressed file sizes. 2. **Fix the predicate.** A cast or function wrapping the column is the usual reason a filter is not in `PushedFilters`; unwrap it and the built-in rules do the work. 3. **Use a hint.** `BROADCAST`, `MERGE`, `SHUFFLE_HASH` express your knowledge for that query only, and are visible to the next reader. 4. **Change the data.** Partitioning, file sizes and clustering fix more query plans than any rule will. 5. **`spark.sql.optimizer.excludedRules`.** If a built-in rule is genuinely misbehaving on your workload, excluding it by name is far cheaper than adding one, and it is reversible in a config change. 6. **A view or a shared function.** If the goal is consistency of expression rather than of execution, this is the boring, correct answer. ## What it costs you **Blast radius.** The rule runs on every query in every session that loads the extension. A pattern that matches more broadly than intended affects work you have never seen. **Silent wrongness.** An optimizer rule that violates semantics does not throw; it returns a different, plausible-looking result. This is the failure mode that makes custom rules a governance decision rather than an engineering preference. Mitigations are plan-level unit tests over the trees you expect to match, differential testing that runs a corpus of queries with and without the rule and compares results, and a config flag that disables the rule instantly in production. **API drift.** `LogicalPlan`, `Expression`, `SparkStrategy` and the pattern-matching helpers are developer APIs. They shift between Spark minor versions, and your rule is a hard dependency on internals during every upgrade — including the pattern-matching and tree-traversal changes that landed as Catalyst evolved through Spark 3.x into 4.x. Budget the upgrade tax as ongoing cost, not one-off. **Ownership.** A rule is a piece of the query engine that your team now owns. It needs a named owner, tests that run against each Spark upgrade, and a documented kill switch. Without those, it becomes the thing nobody dares touch and everybody blames. ## How to answer the question A strong answer starts by declining: name what you would try first and why those fix most problems. Then give the narrow class of cases where a rule is the only mechanism that scales — a cross-cutting invariant that must hold for every query, enforced by the engine rather than by convention — and finish with the operational conditions you would attach to shipping one: a flag to disable it, differential tests, an owner, and an upgrade plan.
- What is the most dangerous failure mode of a custom optimizer rule?Returning wrong results quietly. A rule that does not preserve semantics still produces a valid plan and a plausible answer, so nothing fails and nobody is alerted. That is why such a rule needs differential testing — run a query corpus with and without it and compare outputs — plus plan-level unit tests and a configuration flag that disables it without a redeploy.
- When would you use spark.sql.optimizer.excludedRules instead of writing a rule?When a built-in rule is the problem rather than a missing one. Excluding it by name is a config change, reversible in seconds, with no code to own and no dependency on internal APIs. It is blunt — the rule is off for every query in the session — so it suits a targeted workaround while a Spark bug is fixed upstream, not a permanent design.
- How does a custom rule affect Spark upgrades?It pins you to Catalyst internals. Plan and expression classes, tree traversal helpers and strategy interfaces are developer APIs that change across minor versions, so every upgrade needs the rule recompiled, retested against real plans, and sometimes rewritten. Treat that as a recurring cost with an owner rather than a one-time build.
Amending building code is the right response to a hazard in every building in the city, and the wrong response to a door that sticks in one office.
saying these in an interview costs you the question
- Reaches for a custom rule to fix one slow query
- Assumes a buggy optimizer rule will fail loudly
- Ignores that the rule applies to every query in the session
- Treats Catalyst plan classes as a stable public API
- Cannot name a cheaper lever such as statistics, hints or excluded rules