Lab 02: Spark Internals & Catalyst Optimizer
Before we write code, let's understand our Goal: You need to understand exactly what happens under the hood when you press "Run" on a Spark query. If you don't know how Spark translates your code into physical work, you cannot fix it when it breaks.
The Tool: The Catalyst Optimizer is the "brain" of Apache Spark. It takes your Python code and rewrites it to be as mathematically efficient as possible before executing it.
1. Tracing the Execution Lifecycle
Follow the Data: Let's track a single query from the moment you press "Run" to the moment the data is saved.
The Code:
df = spark.read.parquet("s3://data/").filter("age > 18").groupBy("city").count()
df.write.format("delta").save("s3://output/")
Phase 1: Logical Plan (The Brainstorming Phase)
- Unresolved Logical Plan: Spark parses your code into a tree. It doesn't even check if the columns
ageorcityactually exist yet. - Resolved Logical Plan (Analyzer): Spark checks the internal Catalog to verify that
ageandcityactually exist in the Parquet files. - Optimized Logical Plan (Catalyst Optimizer): The brain steps in and applies rules.
- Predicate Pushdown: It moves the
filter("age > 18")as close to the storage as possible so it doesn't load invalid rows into memory. - Column Pruning: It ignores all columns except
ageandcityto save memory.
- Predicate Pushdown: It moves the
Phase 2: Physical Plan (The Battle Plan)
Catalyst generates multiple Physical Plans. For example: "Should I sort the data first, or should I hash it?" It assigns a cost to each plan and selects the cheapest, most efficient one.
Phase 3: Execution (Anatomy Breakdown)
This is where the actual compute clusters do the work.
Real-World Analogy Mapping: Imagine building a house.
- Job (The Entire Project): Triggered by an Action (like
write,show,count). Building the house is the Job. - Stages (The Phases): A Job is divided into Stages. You must pour the foundation before framing the walls. In Spark, a Stage boundary happens every time a Shuffle occurs (like
groupBy("city")). Records for the same city must be physically moved across the network to the same computer. - Tasks (The Workers): Each Stage is divided into Tasks. A Task is the smallest unit of work (e.g., one worker hammering one nail). In Spark, one Task equals one Core on a Worker Node reading one block of data.
← Previous: Lab 01: Advanced PySpark & Spark SQL | Next: Lab 03: Structured Streaming & Auto Loader →**