Advanced Spark Ch.2: From DataFrame Call to Physical Plan
Outline
- 0:13 The structural deception
- 0:54 Four representations along the way
- 1:19 Why four representations at all
- 2:05 Laziness preserves the freedom to rewrite
- 2:40 The unresolved logical plan
- 3:13 The analyzed logical plan
- 3:51 The optimized logical plan
- 4:33 Isn't it scary that Spark rewrites my query?
- 5:25 Physical plan: where logical becomes distributed
- 6:09 The physical plan is the bill, not the menu
- 6:46 Why not just always broadcast?
- 7:43 Exchanges and shuffle boundaries
- 8:22 Physical plan maps to the DAG of stages
- 8:59 Why is there an AQE then?
- 9:51 Read Spark like a collaborator
Transcript
0:00 Welcome to Learning Podcasts. Advanced Spark, chapter 2: From DataFrame Call to Physical Plan. Today we trace what happens between the code you typed and the job Spark actually ran on the cluster. Every Spark engineer has had this moment. You write a chain of methods. A filter on top of a join, on top of a scan. You hit run, and the job behaves nothing like the code looked. The filter ends up sitting at the data source, not at the top of the tree. Only four columns ever come off disk, even though you wrote `SELECT star`.
0:33 A subquery you wrote turned into a join. An order you specified got reshuffled. You wrote one program, and Spark ran a different one. That gap is where almost every surprising behavior in Spark lives. Most engineers never look inside it. Today we are opening it up. From the moment you call a DataFrame method to the moment Spark has something runnable, your job passes through four internal representations. Unresolved logical plan. Analyzed logical plan. Optimized logical plan. Physical plan. Each 1 is a tree.
1:06 Each 1 describes the same answer. And each 1 is doing different work. So the engine is rewriting your job four times before a single executor sees it. But this is the obvious question. Couldn't Spark just execute the code I wrote? Simpler engines do exactly that. Type a command, run a command. They are easier to reason about, and they are slower at scale, because they have no chance to rewrite anything. Think about a house being built. The architect's sketch, the engineering blueprint, the permitted plan, and the construction schedule are all the same house.
1:40 But each artifact answers a different question. The sketch says what the owner wants. The blueprint resolves it against load-bearing reality. The permitted plan optimizes for code and cost. The construction schedule tells the crew what to do on Monday. If you skip straight from sketch to construction, you build something correct that costs 10 times more than it had to. Same building, four artifacts, four different jobs. This is also why DataFrames are lazy. Every method you chain just builds the tree.
2:09 Spark does not execute anything until you ask for results. Think about a grocery shopper writing the full list before they leave the house. With the full list in hand, they can reorder the route, drop duplicates, and skip aisles entirely. If they shopped the moment each item crossed their mind, they would walk the same store 6 times. Spark collects the full plan 1st so the optimizer has something worth optimizing. Laziness is not "deferring work for later." Laziness is "preserving the freedom to rewrite."
2:37 The full list before the trip. Start at the top of the pipeline. The 1st tree Spark builds is the unresolved logical plan. It captures the shape of what you wrote. A `Filter` node, a `Join` node, a `Project` node. But it does not yet know whether the column names you used are real, whether the table references exist, or whether the types make sense. The references are placeholders. It is a wish list, not a program. The plan knows you want to filter on `region`, but it does not yet know which `region` column, in which table, with which type.
3:09 So at this stage, your code is just intent. The next step is analysis. Spark hands the unresolved tree to the analyzer, which walks the tree against the catalog. Every column reference gets bound to a real schema slot. Types propagate up through expressions. Function calls resolve to actual function definitions. Aliases get expanded. If a column does not exist, this is where the `AnalysisException` fires, before any executor has been contacted. The output is the analyzed logical plan. It is the same tree shape, but every node is grounded in real schema information.
3:45 The plan still describes what you asked for, but Spark is now confident the request is coherent against your tables. Now the analyzed plan goes to the optimizer. This box is called Catalyst. Catalyst takes the analyzed tree and applies a sequence of rewrite rules until the plan stops changing. Filters slide closer to the data source. Projections collapse. Subqueries become joins. The result is the optimized logical plan. Same answer, smaller shape, less work. Catalyst itself is the whole subject of chapter 3, so we will not unpack the rules here.
4:20 The structural point for this chapter is that the plan that comes out of Catalyst is still logical. It describes what to compute, not how to distribute it. So three trees in, and the cluster has not been told anything yet. Pause on that for a 2nd. How am I supposed to trust a plan I did not write? If Catalyst is rewriting my query into something else, how do I know it is still the right query? That is the right question, and it has a clean answer. The contract Catalyst signs with you is equivalence.
4:53 Same input rows, same output rows, every time. Always. The deal is performance. The rewriter is allowed to change the shape of the plan, but never the answer. The safety valve is that `explain` shows you exactly what Catalyst did. You can compare the analyzed plan against the optimized plan and see every rewrite. If you ever doubt a rewrite, the diff is right there. So yes, the plan can be unrecognizable next to the SQL you typed. That is the price of competitive performance, and Spark hands you the receipt.
5:27 Now the optimized logical plan goes to physical planning. This is where Spark turns the logical description into something that resembles a real distributed program. Physical planning answers questions the logical plan never asked. Which join algorithm should Spark use for this particular join? Should it broadcast the small side or sort-merge both sides? Where does a shuffle boundary need to appear? Which operators can be fused into a single stage? The output is the physical plan, a tree of concrete operators with names like file scan, broadcast hash join, sort merge join, hash aggregate, broadcast exchange, shuffle exchange.
6:06 Every one of those nodes maps to something the engine knows how to execute. Here is the cleanest way to think about the difference. The optimized logical plan is the menu. It says "join these two tables, filter the result, group by user." The physical plan is the kitchen ticket the line cook reads. It says "broadcast the small side. Scan the parquet file with the filter pushed into the reader. Run a hash join. No shuffle. Aggregate locally." Same dish on the plate at the end, but the menu and the ticket answer different questions.
6:39 The menu says what the diner wants. The ticket says what the kitchen will actually do. The physical plan is the 1st place where you can see what the cluster is about to spend. If broadcast hash join is fast and shuffle is slow, why does Spark even pick? Why not just always broadcast the small side? Because broadcast has a hard memory constraint. To broadcast a table, Spark has to send a full copy of that table to every executor and hold it in each 1's memory. If the small side is genuinely small, that is a great trade.
7:10 No shuffle, no network reorganization, just a local hash lookup. But the threshold matters. So broadcast is only safe up to a point. Above roughly 10 megabytes by default, every executor has to hold a copy of the broadcast side in memory while doing its own work. Push past what the executor can spare, and the job goes red with out-of-memory errors. So the planner has to choose. Below the threshold, broadcast every time. Above it, fall back to shuffle and accept the cost. The configuration knob is `autoBroadcastJoinThreshold`, and tuning it is one of the most common ways to change a Spark job's runtime by 10 times.
7:45 Look at any non-trivial physical plan and you will see nodes labeled `Exchange`. An exchange is Spark's name for a shuffle boundary, the point where data has to be repartitioned across the cluster. Exchanges show up as shuffle exchange exec for hash repartitioning and broadcast exchange exec for broadcasts. Every exchange in your plan means data is about to move across the network. Reading a physical plan well starts with counting the exchanges and noticing where they sit. Two exchanges in the wrong place can dominate the runtime of a job whose logical plan looked fine.
8:22 So the cost of the job is mostly the cost of the exchanges. The physical plan also has a direct mapping to the DAG of stages you see in the Spark UI. The plan is the compiler's view. The DAG is the scheduler's view. When Spark hits an exchange in the physical plan, it draws a stage boundary. Everything before the exchange becomes one stage. Everything after becomes a dependent stage that cannot start until the shuffle data from the prior stage materializes. So if your plan has two exchanges, your job has roughly three stages.
8:54 That mapping is mechanical, which means you can read the physical plan and predict the stage shape before the job ever runs. That is debugging without an execution. If the physical plan is already this thoughtful, why does Spark also have Adaptive Query Execution? What is left to adapt? Reality. Physical planning is best-effort given the statistics Spark had at compile time. The planner thought the small side was 8 megabytes, but after a filter ran, it turned out to be 8 hundred kilobytes. The planner thought a shuffle would produce two hundred partitions; the real data was small enough to coalesce into 50.
9:31 AQE replans on real numbers between stages. Joins can flip from sort-merge to broadcast mid-job. Skewed partitions can split. The physical plan is the plan with the best information available before any data has moved. AQE is the plan with the best information available after the previous stage finished. Both layers exist because both layers do real work, and chapter 9 goes inside AQE. When you can read the physical plan fluently, Spark stops being a black box and starts being a collaborator with opinions.
10:04 The plan is the engine telling you exactly what it is about to do. Next chapter we open the optimizer itself. How Catalyst rewrites queries. Thanks for listening to Learning Podcasts.