Ch.3: Spark Catalyst: Why Your Query Gets Rewritten

Outline

Transcript

0:00 Welcome to Learning Podcasts. Advanced Spark: How Catalyst Rewrites Queries. Today we crack open the rewriter that takes the query you typed and quietly turns it into the query Spark actually ran. Sometimes those two queries do not look like the same query at all. Every Spark engineer has had this moment. You write a clean filter on top of a join. You scan the plan in the UI, and the filter is gone from the top. It is sitting all the way down at the source. Or you write a subquery, and Spark has quietly turned it into a join.

0:34 Yeah. I have stared at one of those plans for, 20 Minutes wondering where my filter went. Right. Or you write `SELECT star`, and only four columns ever come off disk. You wrote one query. Spark ran a different one. That is not a betrayal. That is the entire product. The thing doing the rewriting is called Catalyst, and that is what we are unpacking today. Quick orientation before we go inside. Sure. Spark turns your code into four trees: parsed, analyzed, optimized, and physical. Right. We walked through that pipeline in the previous chapter.

1:10 Catalyst owns the move from analyzed to optimized. So logical in, logical out. Same answer, same schema, cheaper shape. The cluster has not been told anything yet. Catalyst is a tree rewriter. The plan is a tree of operators. A rule is a function that pattern-matches a sub-tree and returns a new sub-tree. But isn't that what every database optimizer has done since the seventies? What makes Catalyst different? Honestly, the model itself is not new. What is different is that it is exposed as plain Scala objects on a tree you can read in the source.

1:45 There is no SQL parser magic at this layer. No machine-learned optimizer. It is a tree, and there are rules, and you can inspect every one of them. So it is not one big optimization pass. It is a loop of small, mechanical rewrites. Yeah. Stacked on top of each other. A rule has two pieces. A pattern that says, "if you see a node like this." And a rewrite that says, "replace it with a node like that." Think of it as search-and-replace in your editor, except the search is structural and the replacement, Knows the rules of SQL.

2:18 If a `Filter` sits above a `Project`, and the predicate only references columns the project preserves, swap them. The filter slides under the project. Same answer, less data flowing through. And the rule itself does not know what plan it is in? Right. The rewrite is local. The rule does not know what plan it is in, and it does not need to. It just finds its pattern wherever it appears and produces a new sub-tree. Local match, local rewrite. Multiply that across roughly two hundred built-in rules and you have the bulk of Catalyst.

2:48 Wait. If rules can re-enable each other, can they fight each other too? Could one rule undo what another just did? They could in theory. In practice the batches are designed not to. Each batch repeats until no rule in the batch changes the plan anymore. Spark calls that condition a fixpoint. Right. So the batch is the unit of guaranteed termination. Exactly. Then the next batch runs in order. The reason batches loop is that one rule often unlocks another. Resolving an attribute exposes a constant. Folding the constant simplifies a predicate.

3:22 The simplified predicate is now eligible to be pushed down. None of those rules know about each other. The fixpoint loop is what stitches them together into something that looks intentional. Some rules are mechanical and run on every plan. Attribute resolution binds every column reference to a real schema field with a real type. That is what catches the typo in your column name and throws an `AnalysisException` instead of a runtime error. Yeah. Constant folding evaluates expressions that do not depend on a row.

3:50 `WHERE 1 equals 1 AND region equals 'EU'` collapses to `WHERE region equals 'EU'`. And boolean simplification removes the obvious redundancy. `NOT NOT x` becomes `x`. A filter that ANDs a predicate with `true` drops the `true`. Null propagation, expression normalization, dead-branch elimination on `CASE WHEN` are all in this family. Cheap, deterministic, easy to miss. They are the reason the more interesting rules have something clean to work on. Okay, so those clean up the plan. What about the rules that actually move work?

4:24 The next family moves work toward the data. Predicate pushdown: a filter slides as close to the data source as possible, sometimes all the way down into the Parquet reader. We saw this rule do the heavy lifting in the previous chapter's example. And column pruning. If your query only reads four columns out of two hundred, the scan only reads four. The other one hundred 96 never come off disk. Two hundred down to four, before any data leaves the disk. Projection collapse: two `Project` nodes back-to-back become one combined projection, so the engine evaluates each expression once.

4:58 Limit pushdown moves a `LIMIT 100` down past a sort or a scan when it can. And every one of these preserves the answer. Then there are the rewrites that change the shape of the query, not just its order. A correlated `EXISTS` subquery becomes a `LEFT SEMI JOIN`. Wait, but couldn't you just write the join yourself if you knew that was faster? You could. Most engineers do not, because the subquery reads more naturally. So Catalyst meets you where you wrote the readable version and gives you the performant one. A long `IN` list becomes a broadcast hash join against a tiny in-memory relation.

5:36 A `DISTINCT` on top of a `UNION ALL` becomes a single aggregation. So you did not write a join. Catalyst wrote one. The plan you read in the UI looks unrecognizable next to the SQL you typed. The query is unrecognizable. The answer is the same. That is the system working as intended. Here is the wall a pure rule-based optimizer hits. Two plans can be provably equivalent. Same rows, same columns, same answer, every time. Yeah. But one of them might run in 30 seconds and the other in 30 minutes, depending on table sizes, key distributions, and which side of a join is smaller.

6:13 Same answer, very different bill. So rules can prove equivalence. Rules cannot answer which equivalent plan is actually cheaper on this cluster, with this data, right now. That is a different question, and it needs different inputs. So why not just write more rules until the rule-based optimizer always picks the right plan? Because rules cannot read the table. They can only read the plan. Right. So Catalyst has a 2nd layer riding on top: cost-based optimization. CBO takes a set of equivalent candidate plans and uses statistics about your tables to pick one.

6:49 Row counts, column cardinalities, value distributions, average row size, null fractions. I mean, I always thought of the cost-based optimizer as, This giant black box. But that list is not crazy long. Well, it is not. The two big wins are join reordering, where CBO picks the order that builds the smallest intermediate result across a chain of joins, and broadcast threshold decisions, where Spark broadcasts the small side instead of shuffling both. So rules know what is correct. Costs know what is cheap. Catalyst needs both layers.

7:25 So how good is that cost model in practice? As good as your statistics. Same query, two clusters. On the 1st cluster, the stats say users is 5 gigabytes, so the optimizer picks a sort-merge join. Two big shuffles, slow. Mm-hmm. On the 2nd cluster, you ran `ANALYZE TABLE users` last night, the stats say users is 5 megabytes, and Spark broadcasts. No shuffle, fast. Same query, same engine, 10 times the runtime. All from one stale row count. Yeah. Documented. In a runbook. Which nobody reads. So CBO is not a button you turn on.

8:02 It is a contract you have to maintain. There is a runtime safety net for bad stats coming up later in this series, called AQE. So back to where we started. You wrote one query. Catalyst delivered a cheaper, equivalent version. The contract is equivalence. The deal is performance. The job of a rule-based plus cost-based optimizer is to honor the 1st while improving the 2nd, every time, without you asking and without you having to inspect the plan. When the plan looks nothing like your SQL, that is the system doing its job.

8:36 The fact that you can read the plan and see what it did is the safety valve. One more thing worth knowing. Catalyst is not a closed box. It exposes extension points. You can add custom logical rules through the optimizer. Connectors register pushdown rules so their data sources can absorb filters, projections, aggregations, even joins. So most of the connector ecosystem is, you know, just rules underneath? Pretty much. Honestly, the 1st time I tried to read the source for an Iceberg pushdown rule, I had to re-read it three times before it clicked.

9:12 But every Spark connector you use, Iceberg, Delta, Parquet, JDBC, plugs into Catalyst through the same API a custom rule would use. The reason connector quality varies so much across data sources is mostly about how seriously the connector takes its Catalyst integration. Catalyst is conservative for a reason. It will never apply a rule that could change the answer, even if the rule would help in 99 percent of cases. Right. So when a rewrite you expected does not happen, the cause is almost always one of four things. Stale stats: CBO is making a join decision based on a row count from last quarter, and it is wrong by orders of magnitude.

9:52 A non-deterministic predicate, like a filter that calls `rand`, cannot be pushed down because the source might evaluate it differently than the engine, breaking equivalence. Yeah. Randomness breaks the contract. A Python UDF in the middle of your filter: the UDF is opaque to Catalyst, so the filter cannot move past it, and column pruning often cannot reach the columns the UDF reads inside its body. And a wide subquery that returns more columns than the outer query consumes: pruning gets stuck at the subquery boundary.

10:21 So opaque code blocks rewrites. Stale numbers mislead them. So Catalyst is two optimizers riding on one tree. Rules for what is correct, costs for what is cheap. Coming next, we drop a level: how a single executor runs that physical plan as fast as the hardware allows. Tungsten and whole-stage codegen. Thanks for listening to Learning Podcasts.