Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing

PySpark for Distributed Processing

Progress0/21
x

PySpark Architecture

  • PySpark architecture: the one explanation35m
  • When do you need Spark?12m
  • What is a cluster?12m
  • Reading and writing data12m
  • SQL inside Spark12m

Mental model & DataFrames

  • PySpark architecture & mental model10m
  • DataFrame basics10m
  • Column operations & built-in functions10m

Aggregations, windows & nested data

  • Aggregations & groupings10m
  • PySpark window functions12m
  • Joins & optimization strategies10m
  • Handling complex & nested data12m
  • Partitioning, repartition & coalesce10m
  • Medallion pipeline project14m

Performance at Scale

  • Caching & persistence12m
  • Data skew detection and salting14m
  • Reading a Catalyst physical plan12m

Production Spark

  • UDFs in depth, and why to avoid them12m
  • Structured Streaming & watermarks14m
  • Delta Lake: MERGE, time travel, ACID12m
  • Capstone part 2: incremental MERGE16m
Back to track
  1. Learn
  2. PySpark for Distributed Processing
  3. PySpark Architecture
  4. PySpark architecture: the one explanation

Lesson 1 of 21 · Theory first, then run it

PySpark architecture: the one explanation

pysparkbeginner35 min

Overview

Who executes your PySpark code, where it runs, how data moves, and why Spark exists. The full picture from Python to partitions, tasks, and executors.

On this page17 sections›
  1. 1Why Spark exists
  2. 2Spark and PySpark
  3. 3Driver, workers, and executors
  4. 4Partitions and tasks
  5. 5Lazy evaluation
  6. 6DAG, shuffle, and stages
  7. 7Tracing a query
  8. 8Catalyst and Tungsten
  9. 9How tasks run
  10. 10Python and the JVM
  11. 11Memory, cache, and lineage
  12. 12Joins, skew, and Spark UI
  13. 13Run modes
  14. 14Where data lives
  15. 15The picture to memorize
  16. 16What comes next
  17. 17Practice

Why Spark exists

PySpark Architecture: the one explanation I would give a beginner

If you understand this properly, PySpark will stop feeling like magic. The verbs look familiar. The system underneath them is not a faster pandas. It is a distributed engine that plans work, splits data, and runs that work on many machines.

The biggest mistake beginners make is starting with this chain and never asking who actually executes it, where it runs, how data moves, and why Spark exists at all:

PythonThe chain beginners type first
df.filter(...).groupBy(...).agg(...)

That line is a description of a computation, not a loop that walks rows on your laptop. Until you can picture the driver, the executors, the partitions, and the moment an action turns a plan into a job, the API will keep feeling like a black box. So let's build the entire picture from zero, in the same order you will meet it on a real cluster.

This tab is not a cluster

This lesson teaches the architecture. The editor still simulates Spark on a small orders table. There is no live YARN or Kubernetes cluster here. Practice the names and the flow. Judge real scale on a real job.

Before you start

This track assumes you completed Core Python. The pandas track is strongly recommended so later translations land.

1. First: Why did Spark even come into existence?

Imagine you have a Python program that lives entirely on one computer. It loads a list, loops through it, and builds a new list in memory. Nothing is distributed. One process does the whole job.

PythonA single-machine Python loop
data = [1, 2, 3, 4, 5]

result = []
for x in data:
    if x > 2:
        result.append(x * 10)

This works perfectly. Your computer processes the data. The list is tiny, RAM is plenty, and if the process crashes you just run it again. That is how almost everyone starts, and it is the right tool until the data outgrows the machine.

But what happens when your data becomes 10 GB, then 100 GB, then 1 TB, then 1 PB? A single machine starts having problems. Not one dramatic problem. Three separate ones that show up together: memory, speed, and failure.

Problem 1: Memory

Your laptop might have 16 GB of RAM. Your data is 500 GB. You cannot comfortably load everything into one machine. Even if the file is compressed on disk, parsing it into Python objects can be several times larger than the file itself. The process starts swapping, the kernel freezes, or you get an out-of-memory error in the middle of the loop.

One machine vs the dataset
16 GB RAM

But your data is:

500 GB

You could try to stream the file line by line. That helps for simple filters. It falls apart the moment you need a group-by, a join, or a sort, because those operations need related rows to meet. Spark exists in part because those operations on large data need many machines' memory, not a cleverer loop on one laptop.

Problem 2: Processing speed

Suppose processing 1 TB takes one machine 10 hours. What if you had 100 machines? Ideally you divide the work, each machine processes a portion, and you finish much faster. Wall-clock time is what a nightly pipeline cares about. A job that finishes at 6am is useful. A job that is still running at noon is not.

Parallel speedup, in the ideal case
100 machines
→ divide the work
→ each processes a portion
→ finish much faster

You do not get a perfect 100x speedup. There is overhead to split the work, move data, and combine answers. The idea is still the reason Spark is worth running: many CPUs on many machines, each chewing a slice of the same job.

Problem 3: Machine failures

When you are processing huge data across many machines, something will fail. Disks die. Nodes get preempted. A JVM hits a bad record and exits. Your system needs to recover automatically, or every overnight job becomes a manual incident.

Partial failure in a cluster
Machine 1 OK
Machine 2 OK
Machine 3 ❌ crashed
Machine 4 OK
Machine 5 OK

A single-machine script has a simple failure story: restart it. A 100-machine job cannot afford to throw away 99 finished slices because one worker died. Spark keeps enough history of how each slice was produced that it can recompute the lost piece instead of starting from zero. That is not a nice extra. It is why people trust Spark with production lakes.

2. The Core Idea Behind Distributed Computing

Instead of sending the whole dataset into one computer's CPU and memory, we distribute the work. The data is split. Each machine sees a portion. The answers come back together. That shift, from one box to many, is the whole subject of this lesson.

Single-machine processing
DataCPU + MemoryResult

One computer holds the data, the CPU, and the memory. Everything waits on that one box.

We do this instead:

Distributed computing
Distribute workMachine 1Machine 2Machine 3

Split the work, run it on many machines, combine the result.

This is called distributed computing. The idea is simple, and you should be able to say it in one breath:

The fundamental idea

Take a huge problem, divide it into smaller problems, send them to multiple machines, process in parallel, combine the results. That is the fundamental reason Spark exists.

Everything else in Spark architecture (drivers, executors, partitions, shuffles, stages) is machinery for doing that split, that parallel work, and that combine, without you writing the networking yourself.

Spark and PySpark

3. So What Exactly Is Apache Spark?

Apache Spark is a distributed data processing engine. Its job is to answer: you tell me WHAT you want to do with the data. I will figure out HOW to distribute and execute that work across multiple machines.

You stay in the language of tables: filter these rows, group by this column, sum this measure. Spark turns that description into tasks, puts those tasks on machines that have capacity, and brings the result back. You are not writing a cluster program. You are writing a query that a cluster can run.

PythonWhat you write
df.groupBy("country").sum("sales")

You don't manually say which machine processes India, which processes USA, which processes UK, and who combines the results. Spark handles that. The moment you catch yourself thinking in machine numbers, you have dropped below the API Spark wants you to use. Think in tables and keys. Let Spark place the work.

4. Then What is PySpark?

Spark itself is primarily built using Scala and Java. The engine, the schedulers, the memory managers, and the query optimizer live on the JVM. PySpark is the Python API for Apache Spark. It is how Python programmers talk to that engine without writing Scala.

So when you write a group-by in Python, PySpark translates your instructions so that Spark's underlying distributed engine can execute them. Your Python process is not looping over every row. It is sending a description of the computation into Spark, then waiting for an action to demand a result.

PythonPython that Spark will execute
df.groupBy("department").count()
How Python reaches the cluster
You / Python codePySpark APIApache Spark engineCluster of machines

You type Python. PySpark talks to Spark. Spark runs on the machines.

PySpark is not a second engine

PySpark is NOT a separate distributed computing engine. PySpark is simply a way to interact with Apache Spark using Python.

That distinction matters the first time someone asks whether they should 'learn Spark or learn PySpark.' You learn Spark's architecture. You type PySpark. Same engine, Python spelling.

Driver, workers, and executors

5. Before Architecture: Understand "Cluster"

A cluster simply means multiple computers working together. Each one is still an ordinary computer. The cluster is the agreement that they share jobs, storage, and a resource manager. Nothing magical sits between them except the network and the software that assigns work.

A cluster is several machines
ClusterMachine 1Machine 2Machine 3Machine 4

Four heaps, four disks, one network. Not one giant computer.

Each machine can have CPU, RAM, and storage. Together they process large datasets. A four-machine cluster with 32 GB RAM each is not '128 GB of one computer.' It is four separate heaps, four sets of disks, and a network in between. Spark's whole job is to use those four heaps as if they were one processing system, while respecting that they are not one memory space.

6. The Most Important Spark Architecture Diagram

Here is the architecture you should permanently remember. If you can redraw this from memory, the rest of this lesson is names for the boxes you already know.

Spark architecture to memorize
User / PySpark appDriverCluster managerWorkers + executors

User talks to the Driver. Driver asks the cluster manager. Workers run executors and tasks. Result comes back.

There are four major things you need to understand: the Driver, the Cluster Manager, Worker Nodes, and Executors. Let's understand each deeply. If you mix these four up, every later topic (shuffle, skew, broadcast, cache) will feel like unrelated trivia. They are all consequences of this picture.

7. Driver: The Brain of Your Spark Application

The Driver is the component to understand first. Think of Spark as a company. The Driver is the CEO, the brain. It does not necessarily process all the data itself. It decides what should happen, who should do it, and what 'done' looks like.

Driver coordinating workers
Driver (brain)Worker 1Worker 2Worker 3

The Driver is the brain. Workers do the row work.

Instead, it receives your code, creates the execution plan, divides work into tasks, sends tasks to executors, monitors execution, and collects results. When a task fails, the Driver is the one that notices and asks for a retry. When you call show(), the small sample comes back to the Driver so it can print it.

Example. You write:

PythonWhat the Driver receives
df = spark.read.csv("data.csv")

result = (
    df.filter(df.salary > 50000)
      .groupBy("department")
      .sum("salary")
)

result.show()

The Driver sees your instructions. Conceptually it lines them up as a sequence of steps, then creates a plan for executing that sequence across the cluster. You did not write those steps as a workflow engine. Spark inferred them from the DataFrame chain.

How the Driver sees your chain
STEP 1
Read data

STEP 2
Filter salary > 50000

STEP 3
Group by department

STEP 4
Calculate sum

STEP 5
Show result

Driver

Driver = The coordinator and brain of your Spark application.

8. SparkSession: Your Entry Point

When working with PySpark, you usually start by building a SparkSession. That object is the handle your Python program holds for the rest of the application. Without it, you have no catalog, no readers, no SQL, and no connection to a cluster.

PythonCreating a SparkSession
from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("My Application") \
    .getOrCreate()

What is spark? It is your SparkSession. Think of it as your main entry point into Spark functionality. Architecture-wise, your Python code talks to SparkSession, SparkSession belongs to the Spark application on the Driver, and the Driver talks to the cluster.

Through SparkSession, you can read data, run SQL, and create DataFrames:

PythonWhat SparkSession unlocks
spark.read
spark.sql()
spark.createDataFrame()
SparkSession in the stack
Your Python codeSparkSessionSpark application / DriverSpark cluster

Your Python talks to SparkSession. SparkSession lives on the Driver. The Driver talks to the cluster.

For beginners: SparkSession is how your PySpark application connects to Spark. getOrCreate() means 'reuse a session if this process already has one.' That is why notebooks often call it at the top of every cell without spawning a new application each time.

9. Cluster Manager: The Resource Manager

Imagine you have 100 machines. The Driver says 'I need 10 machines and 40 CPU cores.' Who decides which machines should be assigned? That is the job of the Cluster Manager. The Driver knows what it wants. The Cluster Manager knows what the cluster can spare.

Spark supports cluster managers such as Standalone, YARN, and Kubernetes. You will meet those names in job configs. The programming model does not change. Only the thing that hands out CPUs and memory changes.

Driver asking the cluster manager
Cluster managerMachine 1Machine 2Machine 3

The Driver asks for resources. The cluster manager picks the machines.

The Cluster Manager says: okay. You can use Machine 1 with 4 cores, Machine 2 with 4 cores, Machine 3 with 4 cores. Then Spark starts executors on those machines. If the cluster is busy, you might get fewer cores than you asked for, or wait in a queue. That wait is not a Spark bug. It is the resource manager doing its job for every application on the cluster, not only yours.

Important distinction:

Planning vs allocating. Do not mix these two.

RoleResponsible for
DriverPlanning and coordinating work.
Cluster ManagerManaging and allocating resources.

10. Worker Nodes

A Worker Node is simply a machine available in the cluster to perform Spark work. It is hardware (or a VM, or a pod) that can run executor processes. When people say 'the cluster has twenty nodes,' they mean twenty of these.

Worker nodes in a cluster
CLUSTER                   ├── Machine 1├── Machine 2├── Machine 3└── Machine 4

These are worker nodes. Inside worker nodes, Spark launches Executors. A worker can host one executor or several, depending on how the application is configured. Beginners can treat 'worker' as the machine and 'executor' as the Spark process on that machine. That split will keep you out of trouble when you read Spark UI later.

11. Executors: The Actual Workers

This is where actual computation happens. The Driver can be a brilliant planner. If no executor is running tasks, no data is processed. Executors are JVM processes that live on worker nodes for the life of your application (or until they are lost and replaced).

What executors do
Driver                           │                             │ sends work                  ▼                          Executors                        │                             ├── Process data              ├── Run tasks                 ├── Store intermediate data   └── Return results         

Example: the Driver fans work out to three executors. Each one processes a different slice of the data. They do not share Python lists in memory. They share a plan, and they exchange data only when a shuffle or a result collection requires it.

Executors processing slices
                DRIVER                                                                     │                                                           ┌───────────┼───────────┐                                               ▼           ▼           ▼                                            Executor    Executor    Executor       1           2           3                                               │           │           │                                           Process      Process      Process   Data A       Data B       Data C 

Executors perform the actual computation. They have CPU cores, memory, and tasks. A simple analogy: the Driver is the project manager, and the Executors are the employees doing the actual work. If you cache a DataFrame, it is executor memory that holds the cached partitions, not the Driver's laptop RAM.

12. Driver vs Executor: Extremely Important

If you remember only one comparison from this lesson, make it this table. Almost every production outage beginners cause is a Driver/Executor mix-up: collecting too much to the Driver, or assuming 'my Python variable' holds the lake.

The Driver decides what should happen. Executors make it happen.

DriverExecutor
BrainWorker
Creates execution planExecutes tasks
Coordinates clusterProcesses data
Schedules tasksRuns tasks
Usually one per applicationMultiple
Maintains application stateProcesses partitions

Who does what

The Driver decides what should happen. Executors make it happen.

Partitions and tasks

13. Now We Need to Understand the Core Concept: Partitions

Imagine you have a 1 TB file. Spark cannot efficiently give this entire dataset to one executor. Even if one machine had the RAM, you would throw away every other CPU in the cluster. The dataset has to be split into pieces that can be processed independently.

Instead, 1 TB of data gets divided into Partition 1, Partition 2, Partition 3, Partition 4, and so on. Each partition can be processed independently. Then Executor 1 can take Partition 1, Executor 2 takes Partition 2, Executor 3 takes Partition 3. This is how Spark achieves parallelism.

A dataset split into partitions
              DATASET                                                     1 TB                                                       │                                                ┌────────┼────────┐                                       ▼        ▼        ▼                                      P1       P2       P3                                     250GB    250GB    250GB

Those 250 GB numbers are only a sketch. Real partitions are often much smaller (128 MB or 256 MB is a common default from filesystems and Spark). The idea is the same: many piles, many tasks, many CPUs busy at once.

14. Partition = Unit of Parallelism

This sentence is extremely important: partitions determine how much parallel processing Spark can perform. Not the number of machines you rented. Not the number of lines of Python you wrote. The number of partitions that still have work left.

Imagine 100 GB of data in 10 partitions. Spark can potentially process multiple partitions simultaneously, one per task, spread across executors. But if you only have 1 partition, then one executor processes everything, even if you have 100 machines available. The extra machines sit idle because there is no extra pile to hand them.

How parallelism is born
DATA
 ↓
PARTITIONS
 ↓
TASKS
 ↓
EXECUTORS

So remember that chain. Data becomes partitions. Partitions become tasks. Tasks run on executors. If a job is slow and the cluster looks idle, the first question is not 'do we need more nodes?' It is 'how many partitions does this stage actually have?'

15. What is a Task?

A task is the smallest unit of work sent to an executor. Let's say a dataset has four partitions. Spark may create four tasks: Task 1 processes Partition 1, Task 2 processes Partition 2, and so on. Then those tasks are scheduled onto executors that have a free core.

Tasks mapped onto partitions
Dataset                                                 ├── Partition 1             ├── Partition 2             ├── Partition 3             └── Partition 4                                         Spark may create:                                       Task 1 → Process Partition 1Task 2 → Process Partition 2Task 3 → Process Partition 3Task 4 → Process Partition 4                            Then:                                                   Executor 1 → Task 1         Executor 2 → Task 2         Executor 3 → Task 3         Executor 4 → Task 4         

So a useful mental model is: Partition is the unit of data. Task is the unit of work. Typically, one task processes one partition for a particular stage. If a stage has 200 partitions, you will see about 200 tasks in Spark UI for that stage. That is not a coincidence. It is the same idea counted two ways.

Lazy evaluation

16. Now Let's Understand What Happens When You Run PySpark Code

Suppose you write a small pipeline: read a CSV, filter high salaries, count by department, then show the result. Let's trace the complete journey. This is the same story you will replay, with different column names, for the rest of your Spark career.

PythonA small PySpark pipeline
df = spark.read.csv("employees.csv")

filtered_df = df.filter("salary > 50000")

result = filtered_df.groupBy("department").count()

result.show()

PHASE 1: You Write PySpark Code

At this point, something surprising happens. Spark usually does NOT immediately execute the code. The assignments succeed. The Python variables exist. No cluster CPU has started scanning employees.csv yet. Why? Because Spark uses lazy evaluation.

17. Lazy Evaluation: One of Spark's Superpowers

When you write df.filter(...), Spark doesn't necessarily execute immediately. Instead, Spark says: okay, I remember that you eventually want to filter. Then groupBy(...): okay, I also remember that. Then .count() as a transformation chained on a grouped DataFrame: okay, noted. Spark builds a plan.

Only when you call an action, such as show(), collect(), count(), or write(), does Spark actually execute the work. Until that call, you have a recipe. You do not have cooked food. That is why printing a DataFrame object in Python shows a schema and a plan, not your rows.

Beginners fight this the first week. They add a filter, then stare at the cluster hoping to see CPU. Nothing happens. They think Spark is broken. Spark is waiting for an action. Once you expect that wait, lazy evaluation stops being a quirk and starts being the reason Spark can optimize.

18. Why Lazy Evaluation Exists

Imagine you write three transformations in a row: filter salary, filter age, select a few columns. If Spark executed every line immediately, each operation would scan or materialize on its own. That could be inefficient. You would read columns you are about to throw away, and you would run two filter passes where one would do.

PythonSeveral transformations, one later action
df.filter("salary > 50000") \
  .filter("age > 25") \
  .select("name", "salary")

Instead, Spark waits. Then it sees the complete logic: read data, filter salary, filter age, select columns. Now Spark can optimize it. For example: I don't actually need all columns. Let me read only name, salary, and age. Or: let me combine these filters. This is one of the major reasons Spark can be efficient. Eager Python cannot easily do that, because each line already ran before the next line was written.

19. Transformations vs Actions

This leads to two extremely important categories. If you can sort a method into one of these two buckets, you can predict whether a line will hit the cluster.

Transformations

Transformations describe what you want to do. They generally build the plan. filter, select, groupBy, join, withColumn, drop are the usual examples. When you write df2 = df.filter("salary > 50000"), no immediate execution necessarily happens. df2 is a new DataFrame object that knows how it would be computed.

PythonA transformation: plan only
df2 = df.filter("salary > 50000")

Actions

Actions say: okay Spark, now execute. show, collect, count, write, take are the usual examples. When you write df2.show(), now Spark starts execution. The plan is compiled, stages are scheduled, tasks run, and a result (or a written dataset) appears.

PythonAn action: execute the plan
df2.show()

A practical habit: count the actions in a notebook. Each one can re-run lineage from the start unless you cache. That is why a notebook with twelve show() calls on the same giant frame can cost twelve scans.

DAG, shuffle, and stages

20. The Complete Execution Flow

This is the architecture flow you should visualize. Keep it in your head as a vertical slide. Every later section is a zoom-in on one of these boxes.

End-to-end execution flow
PYTHON CODE                        │              ▼                        PYSPARK API                        │              ▼                        TRANSFORMATIONS                    │              ▼                        LAZY EVALUATION                    │              ▼                        LOGICAL PLAN                       │              ▼                        OPTIMIZATION                       │              ▼                        PHYSICAL PLAN                      │              ▼                        DAG                                │              ▼                        STAGES                             │              ▼                        TASKS                              │              ▼                        EXECUTORS                          │              ▼                        PROCESS DATA                       │              ▼                        RESULT         

Now let's understand every layer. You do not need to memorize JVM class names. You do need to know that Python is not the thing scanning Parquet, and that 'the DAG' is not a synonym for 'the job.'

21. What is a DAG?

DAG means Directed Acyclic Graph. Don't get scared by the name. Break it down, because each word is doing a job.

Directed: operations move in a direction, A to B to C. Acyclic: no circular loop. You never go A to B to C and back to A. Graph: a representation of connected operations. Spark internally builds this computational graph so it can see the whole recipe at once.

A computation as a DAG
Read Data      │          ▼      Filter         │          ▼      Select         │          ▼      Group By       │          ▼      Aggregation

This execution relationship can be represented as a DAG. When Spark UI shows you a job, the picture you are looking at is this graph, sliced into stages. If you can sketch the DAG for a query, you can predict where it will shuffle before you run it.

22. Why Does Spark Need a DAG?

Because Spark needs to understand what operations depend on other operations. You cannot aggregate by department until the filter has produced the rows that belong in the aggregate. You cannot filter until the read has produced rows. The DAG is that dependency list, drawn.

PythonA chain with clear dependencies
df.filter(...)
  .select(...)
  .groupBy(...)
  .count()
Dependencies in order
Read

 ↓

Filter

 ↓

Select

 ↓

GroupBy

 ↓

Count

Spark creates this dependency graph. Then it can decide what can run together, what can run in parallel, where data movement is required, and how to recover if something fails. Without a DAG, retrying a failed task would mean guessing. With a DAG, Spark knows exactly which parent output a lost partition was built from.

23. Narrow Transformations vs Wide Transformations

This concept shows up in almost every slow Spark job. Let's start with partitions. Imagine three input partitions holding a few numbers each. A transformation either can finish inside each pile, or it cannot.

Three independent partitions
Input

Partition 1 → [1, 2, 3]
Partition 2 → [4, 5, 6]
Partition 3 → [7, 8, 9]

Narrow Transformation

A narrow transformation means data from one input partition is used by only one output partition. filter, select, and map are the usual examples. Suppose you write df.filter("age > 25"). Spark can do Input P1 to Output P1, Input P2 to Output P2, Input P3 to Output P3. No need for partitions to communicate. This is efficient.

Narrow: no cross-partition traffic
P1 ─────────▶ P1                P2 ─────────▶ P2                P3 ─────────▶ P3

Narrow transformations pipeline well. Spark can run a read, a filter, and a select in one task, on one partition, without stopping to wait for any other executor. That is why a long chain of filters is usually cheap compared with a single groupBy.

24. Wide Transformation

A wide transformation means data needs to move across partitions. groupBy, join, orderBy, and distinct are the usual examples. Suppose we have India and USA rows split across two partitions, and we want groupBy("country"). Spark needs all India records together. But currently India lives in Partition 1 and also in Partition 2. So data must move.

Wide transformation: shuffle by key
Before:

P1                    P2

India                 USA
USA                   India
India                 UK

          SHUFFLE

After:

P1                    P2

India                 USA
India                 USA

                      UK

This movement of data is called shuffle. Once you can see a shuffle in your head, Spark UI's 'Exchange' operator stops being a mystery. It is this picture, implemented with files and network transfer.

25. Shuffle: The Expensive Operation

Shuffle means redistributing data across partitions. Example: df.groupBy("department"). Spark might need every executor to send some of its rows to other executors so that each department lives in one place. That is not a Python groupby on a local dict. It is a distributed regrouping.

Shuffle across executors
Executor 1                                     │                                           ├──────────────┐                            │              │                         Executor 2        │                            │              │                            └──────────────┼──────► Redistribute Data                  │                         Executor 3        │                                           │                                           ▼                                      New Partitions                 

Shuffle is expensive because it can involve network transfer, disk I/O, serialization, and sorting. This is why operations like groupBy, join, and orderBy can be expensive. They are not 'slow Python.' They are 'move a lot of data between machines.' When you optimize Spark, you are usually trying to do less of this, or to shuffle less bytes when you must.

26. Stages: Where Spark Breaks the Job

Next comes how Spark cuts the DAG into stages. Spark sees the DAG and divides it into stages. Why? Because some operations can happen continuously without moving data. Others require a shuffle. Spark groups the first kind together, then draws a line at the shuffle, then starts a new group.

PythonNarrow work, then a wide boundary
df
.filter(...)
.select(...)
.groupBy(...)
.agg(...)

Spark might see Stage 1 as read, filter, select, and maybe a partial aggregation, then a shuffle, then Stage 2 as the final groupBy and aggregation. The shuffle creates a boundary. So: wide transformations usually create stage boundaries.

Stages split at shuffle
Stage 1

Read
 ↓
Filter
 ↓
Select

------ SHUFFLE ------

Stage 2

GroupBy
 ↓
Aggregation

When Spark UI shows 'Stage 0' and 'Stage 1,' you are looking at this split. Stage 0 is not 'less important.' It is 'everything we can do before we have to move data.'

27. Stage vs Task vs Job

Beginners confuse these three constantly. Let's fix that forever. Suppose you filter, groupBy, count, and show. One Python statement chain. Several Spark nouns.

PythonOne action, several Spark nouns
df.filter(...)
  .groupBy(...)
  .count()
  .show()

Job

An action triggers a Job. show() creates a job. If you later call count() on the same DataFrame, that is a second job. Jobs are not 'your whole application.' They are 'one demand for a result.'

Stages

Spark divides the job into stages. Often two, separated by a shuffle boundary: Stage 1, then Stage 2. A job with two shuffles can have three stages. Count the wide boundaries, then add one.

Tasks

Each stage gets divided into tasks. Stage 1 might have Task 1 through Task 4, one per partition. Each task processes a partition. The hierarchy is the thing to memorize:

The essential hierarchy
ACTION

  ↓

JOB

  ↓

STAGES

  ↓

TASKS

  ↓

PARTITIONS

This hierarchy is absolutely essential. Interviewers ask it because it is how you read Spark UI without drowning. Application is the whole program. Jobs are per action. Stages are per shuffle-free chunk. Tasks are per partition.

Tracing a query

28. The Complete Example: Let's Trace Everything

Let's use this code: read sales Parquet, filter amount over 100, group by country, sum amount, then show. Let's see what actually happens, from application start to the printed table.

PythonThe example we will trace
result = (
    spark.read.parquet("sales")
    .filter("amount > 100")
    .groupBy("country")
    .sum("amount")
)

result.show()

Step 1: Application Starts

You submit a PySpark Application. Spark creates a Driver. That is the process that will hold SparkSession, the DAG, and the schedulers. In local mode, that Driver is on your laptop. In cluster mode, it may be a process inside the cluster. Either way, one Driver per application.

Step 2: Driver Connects to Cluster Manager

The Driver says it needs resources. The Cluster Manager agrees (or queues you). Then executors start. You now have a Driver plus Executor 1, Executor 2, and Executor 3. Until executors register, there is nobody to run tasks. A job submitted in that window waits.

Executors after resource grant
Driver                           │                          Cluster Manager                  │                          ├── Executor 1 ├── Executor 2 └── Executor 3 

Step 3: Spark Reads Data

Suppose data is divided into four partitions under sales/. Spark knows the data can be processed in parallel. For Parquet, Spark can also see footer metadata: min/max values, column names, row-group sizes. That information feeds later optimizations. A CSV would give Spark much less to work with.

Input already split
sales/

Partition 1
Partition 2
Partition 3
Partition 4

Step 4: Transformations Build the Plan

You wrote filter, groupBy, and sum. Spark builds something conceptually like Read, then Filter amount > 100, then Group by country, then Sum amount. Still no execution until show(). The Python variables result and the chain that produced it are plan objects.

Step 5: Action Triggers a Job

show() is the action. Spark says: now I need to actually execute this. A Job is created. From this moment, Spark UI will show a job you can click into. Before this moment, Spark UI had an application and maybe some executors, but no job for this query.

Step 6: Catalyst Optimizer Optimizes the Query

For DataFrames and Spark SQL, Spark has something called the Catalyst Optimizer. Think of it as Spark's query optimizer. You tell Spark WHAT you want. Spark tries to decide HOW to execute it efficiently. That is the same split you already know from SQL databases: a query parser, then an optimizer, then an engine.

Example: you select two columns and filter a third. Spark may optimize with column pruning, predicate pushdown, filter reordering, and join strategies. Conceptually: your code becomes a logical plan, Catalyst rewrites it, you get an optimized logical plan, then a physical plan.

PythonA chain Catalyst can rewrite
df.select("name", "salary") \
  .filter("salary > 50000")
Catalyst's path
YOUR CODE

      ↓

LOGICAL PLAN

      ↓

CATALYST OPTIMIZER

      ↓

OPTIMIZED LOGICAL PLAN

      ↓

PHYSICAL PLAN

Catalyst and Tungsten

29. Logical Plan vs Physical Plan

This is easier than it sounds. The Logical Plan describes what needs to happen: read sales, filter amount > 100, group by country, calculate sum. No detailed decision yet about exactly how. It is the query in operator form.

The Physical Plan describes how Spark will actually execute it. For example: Executor 1 reads partitions 1 and 2, Executor 2 reads partitions 3 and 4, they perform local filtering, shuffle by country, redistribute records, and perform aggregation. So: Logical Plan is WHAT. Physical Plan is HOW.

When you call explain(), you are asking Spark to print this HOW. You will practice that in a later lesson. For now, know that the physical plan is where join algorithms and shuffle operators show up.

30. Catalyst Optimizer

The Catalyst Optimizer is mainly responsible for optimizing DataFrame and Spark SQL operations. Think: the developer says 'I want this result.' Catalyst looks for a better way to execute it. Then you get an optimized plan. You still wrote the same PySpark. The engine chose a cheaper path.

Common optimizations include the four below. You do not implement them. You write DataFrame code that Catalyst can see. Python UDFs, by contrast, are opaque, which is why they often disable this good luck.

Predicate Pushdown

Instead of reading everything and filtering later, Spark may push filtering closer to the data source when supported: read only relevant data. On Parquet and many lake tables, that can mean skipping entire row groups whose max amount is already below 100. On CSV, there is much less to push into. The filter still runs. It just runs after a full parse.

Column Pruning

Suppose your dataset has 100 columns, but you only use name and salary. Spark can avoid unnecessarily reading all columns when the source format supports it. That is one of the practical reasons people store lakes as Parquet or Delta, not as giant CSVs.

Constant Folding

Simplifies expressions where possible. Example conceptually: salary > (1000 + 500) can become salary > 1500. Tiny on its own. Together with the other rewrites, it keeps the physical plan from doing busywork.

Join Optimization

Spark can choose different strategies for joining datasets. For example: Broadcast Join, Sort Merge Join, Shuffle Hash Join. We'll revisit joins later. The beginner takeaway is that Spark, not your Python, picks the join algorithm, using sizes and hints Catalyst can see.

31. Tungsten: Spark's Execution Engine Optimization

Another term you'll hear: Tungsten. Catalyst focuses heavily on optimizing the query plan. Tungsten focuses heavily on making execution efficient. Historically, key areas include memory management, CPU efficiency, binary processing, and code generation.

You don't need to memorize all internal implementation details initially. Just remember: Catalyst optimizes the plan. Tungsten helps optimize execution. When someone says 'whole-stage codegen,' they are talking in the Tungsten neighborhood: generate tight code for a stage instead of interpreting one row at a time.

Catalyst vs Tungsten
CATALYST

"What's the best plan?"

TUNGSTEN

"Let's execute that plan efficiently."

How tasks run

32. Back to Stages

Our example filters sales then groups by country. Spark might create a job with Stage 1 as read data, filter amount > 100, and a partial aggregation, then a shuffle, then Stage 2 as the final aggregation. Why? Because groupBy("country") requires data to be reorganized. All records belonging to the same country need to come together.

Partial agg, shuffle, final agg
JOB                                                                       │                                    ├── Stage 1                          │                                    │    Read Data                       │                                    │    Filter amount > 100             │                                    │    Partial aggregation             │                                    │                                    └──────── SHUFFLE ────────                                    │                                    ▼                               Stage 2                                                                   Final aggregation                                                              │                                                                         ▼                                                                       Result        

That 'partial aggregation' line is worth a pause. Before shuffling, each partition can sum amounts per country locally. Then the shuffle moves much smaller partial totals, not every raw sale. Spark does this when it can. It is one reason a groupBy is expensive but not as expensive as shipping every row.

33. Tasks Are Sent to Executors

Suppose Stage 1 has 8 partitions. Spark creates approximately 8 tasks. Each task is 'run Stage 1's pipeline on this partition.' Executors process them. Depending on available CPU cores, multiple tasks can run simultaneously. An executor with one core runs one task at a time. An executor with four cores can run four.

Tasks scheduled onto executors
Stage 1

Task 1 → Partition 1
Task 2 → Partition 2
Task 3 → Partition 3
Task 4 → Partition 4
Task 5 → Partition 5
Task 6 → Partition 6
Task 7 → Partition 7
Task 8 → Partition 8

Executor 1

Task 1
Task 4
Task 7

Executor 2

Task 2
Task 5
Task 8

Executor 3

Task 3
Task 6

34. Executor Cores and Parallelism

Suppose Executor 1 has 4 CPU cores. Then it may run approximately 4 tasks simultaneously. Visualize four cores, four tasks. After one finishes, that core takes the next pending task. This continues until all tasks are complete.

Cores run tasks in parallel
Executor

Core 1 → Task 1

Core 2 → Task 2

Core 3 → Task 3

Core 4 → Task 4

This is why spark.executor.cores and the number of partitions interact. 10,000 tiny partitions on 4 cores means a lot of scheduling overhead. 2 giant partitions on 40 cores means 38 cores nap. Parallelism is a matching problem between piles of data and cores that can lift them.

35. Complete Architecture: The Full Picture

Now combine everything. The user writes PySpark. SparkSession sits on the Driver. The Driver builds the DAG, optimizes, and schedules. The Cluster Manager grants workers. Executors run tasks in memory. Results flow back. This is the same diagram as section 6, now with the words you have earned.

The full architecture picture
                     USER                                                                                    │                                                                              PySpark Code                                                                                 │                                                                                     ▼                                                                            ┌─────────────────┐                        │  SparkSession   │                        └────────┬────────┘                                 │                                          ▼                                                                            ┌─────────────────┐                        │     DRIVER      │                        │                 │                        │ • Application   │                        │ • DAG Creation  │                        │ • Optimization  │                        │ • Scheduling    │                        └────────┬────────┘                                 │                                          │ Request Resources                        ▼                                                                            ┌─────────────────┐                        │ CLUSTER MANAGER │                        │                 │                        │ YARN / K8s /    │                        │ Standalone      │                        └────────┬────────┘                                 │                             ┌────────────┼────────────┐                │            │            │                ▼            ▼            ▼                                                      ┌──────────┐ ┌──────────┐ ┌──────────┐     │ WORKER 1 │ │ WORKER 2 │ │ WORKER 3 │     └────┬─────┘ └────┬─────┘ └────┬─────┘          │            │            │           ┌────▼─────┐ ┌────▼─────┐ ┌────▼─────┐     │ EXECUTOR │ │ EXECUTOR │ │ EXECUTOR │     │          │ │          │ │          │     │ Tasks    │ │ Tasks    │ │ Tasks    │     │          │ │          │ │          │     │ Memory   │ │ Memory   │ │ Memory   │     └──────────┘ └──────────┘ └──────────┘                                                     │            │            │                                                           └────────────┼────────────┘                             │                                          ▼                                                                                  RESULTS                

Python and the JVM

36. But Wait: How Does Python Actually Communicate With Spark?

This is particularly important for PySpark. Your code is Python: df.filter(...). But Spark's core engine runs primarily on the JVM. So conceptually Python talks to PySpark, PySpark communicates with the JVM Spark engine, and that engine talks to executors.

Python to JVM to executors
PYTHON                                    Your PySpark Code                               │                                         ▼                                   PYSPARK                                         │ Communication                           ▼                                   JVM                                       Spark Engine                                    │                                         ▼                                   EXECUTORS            

A simplified mental model: Python API, then PySpark, then Spark JVM, then distributed execution. For DataFrame operations, your Python code generally describes operations that Spark can plan and execute in its optimized engine rather than shipping ordinary Python logic to every record.

This distinction becomes very important when comparing df.filter(...) with Python UDFs. Native functions stay inside the plan. A Python function you wrote by hand often cannot. That is the next section.

37. Why Are Python UDFs Sometimes Slow?

Suppose you write a UDF that doubles a value. Now data may need to cross between the JVM and a Python process and back again. This can introduce serialization and communication overhead. Each row (or batch, depending on the UDF type) pays a toll to leave the optimized engine, run your Python, and return.

PythonA Python UDF
@udf
def my_function(x):
    return x * 2
UDF data crossing
JVM             │              │ Convert Data ▼             Python          │              │ Run Function ▼             JVM            

That's why native Spark functions are generally preferred. Instead of a UDF that doubles a column, prefer a column expression Spark understands and can optimize.

PythonNative expression Spark can optimize
from pyspark.sql import functions as F

df.withColumn(
    "double",
    F.col("value") * 2
)

Because Spark understands native operations and can optimize them. A later production lesson goes deeper on UDFs. The architecture point is already here: UDFs punch a hole in the JVM plan and make the Python/JVM boundary a per-row cost.

38. Spark DataFrame vs Pandas DataFrame

This is a useful comparison, because the method names will trick you into thinking they are the same object. They are not. Pandas generally runs on one machine, with the DataFrame in RAM. Spark generally distributes across multiple machines, with the DataFrame as a plan over partitions.

Pandas: one machine
Your Computer                         ┌─────────────────┐│                 ││ Pandas DataFrame││                 ││      RAM        ││                 │└─────────────────┘
Spark: many machines
                  CLUSTER                                                      ┌───────────┼───────────┐                                               ▼           ▼           ▼                                           Machine      Machine      Machine                                           │           │           │                                               ▼           ▼           ▼                                           Partition   Partition   Partition

39. Why Spark DataFrames Feel Similar to Pandas

You can write a salary filter in pandas with boolean indexing, and in PySpark with filter. They look conceptually similar. But internally they are very different. Pandas executes immediately on the local machine. Spark builds an execution plan, optimizes, distributes, and executes across the cluster.

PythonPandas
df[df["salary"] > 50000]
PythonPySpark
df.filter("salary > 50000")

The similarity is a gift for learning and a trap for debugging. If you bring pandas habits (iterate rows, mutate in place, assume len(df) is free) into Spark, you will either wait forever or crash the Driver. Use the similar verbs. Drop the similar mental model of 'this variable is the data.'

Memory, cache, and lineage

40. The Spark Concept That Explains Speed: Data Locality

Suppose your data is stored on Machine 1, but your computation runs on Machine 5. Then data needs to travel over the network. Network transfer is expensive. So Spark tries to achieve data locality: try to process data close to where it already exists.

Bad: move data to compute
Machine 1        │            │ Network    ▼        Machine 5    

Instead of move data to computation, Spark prefers move computation to data. Conceptually: where is this data? Machine 2. Run the task on Machine 2. This reduces network transfer. You will not always get locality (caches, busy CPUs, and cloud object storage change the story), but the preference is baked into scheduling. It is the same instinct as 'don't copy the warehouse to the office; send the crew to the warehouse.'

41. Spark Memory: What Happens Inside an Executor?

Executors have memory. A simplified split is execution memory versus storage memory. Execution memory is used while performing operations: shuffle, join, sort, aggregation. Storage memory is used for cache and persisted data. The two pools can borrow from each other under pressure, which is why a heavy shuffle can evict a cache, and a huge cache can starve a shuffle.

Executor memory pools
EXECUTOR MEMORY                             ├── Execution Memory  │                     │   Used for:         │   ├── Shuffle       │   ├── Join          │   ├── Sort          │   └── Aggregation   │                     └── Storage Memory        │                     Used for:             ├── Cache             └── Persisted Data

A simplified way to remember: Execution Memory is used while performing operations. Storage Memory is used for cached data. When a job dies with 'out of memory' on an executor, you are usually looking at one of these two pools overflowing, not at Python's laptop RAM. Driver OOM is a different failure: too much result collected home.

42. Cache and Persist

Suppose you read a large Parquet dataset, then you use it multiple times: a filter, a groupBy, a join. Spark might repeatedly recompute or reread parts of the lineage. You can say df.cache(). Now Spark can keep data in memory when appropriate. The first action after cache pays the cost. Later actions try to reuse the pin.

PythonReuse without a cache: lineage reruns
df = spark.read.parquet("large_data")

df.filter(...)
df.groupBy(...)
df.join(...)
Recompute every action
Without Cache

Read Data
   ↓
Process

Read Data Again
   ↓
Process

Read Data Again
   ↓
Process
Cache: read once, reuse
With cache:

Read Data

   ↓

CACHE

   ↓

Use Again
   ↓

Use Again
   ↓

Use Again

But important: don't cache everything. Caching consumes memory. Cache when the dataset is reused multiple times, recomputation is expensive, and memory resources are sufficient. Do not cache the raw lake you touch once. A later lesson practices the cache() call. The architecture point is that cache lives in executor storage memory, and it is a trade of RAM for less recomputation.

43. Lineage: How Spark Knows What Happened

Suppose you build a chain: read CSV into df1, filter into df2, select into df3, groupBy count into df4. Spark remembers the chain. This is called lineage. Lineage means the history of transformations used to produce a dataset. It is the DAG viewed from one DataFrame backward to its sources.

PythonA lineage chain
df1 = spark.read.csv("data.csv")

df2 = df1.filter(...)

df3 = df2.select(...)

df4 = df3.groupBy(...).count()
Lineage of df4
data.csv

   ↓

filter

   ↓

select

   ↓

groupBy

   ↓

count

44. Why Lineage Matters: Fault Tolerance

Suppose Executor 3 crashes. It was processing Partition 7. Does Spark need to restart the entire job? Not necessarily. Spark knows the lineage: original data, then filter, then transformation, then Partition 7. Spark can recompute the lost partition.

Recompute from lineage
Partition Lost                              │                                     ▼                               Check Lineage                               │                                     ▼                               Recompute Partition

This is one of Spark's fault tolerance mechanisms. You do not write the retry logic. The lineage graph is the retry logic. That is why Spark can lose an executor and still finish, and why a shuffle file going missing is a bigger deal than a lost in-memory partition of a narrow stage.

45. Fault Tolerance in Simple Terms

Traditional thinking: something failed, start everything again. Spark: something failed, what exactly was lost, can I recompute only that, yes, recompute. This makes distributed processing more reliable. It is also why Spark wants your sources to still be there. If the original CSV was a temp file you deleted, lineage cannot help.

46. Shuffle Files and Fault Tolerance

There is an important distinction. For narrow transformations, Spark can often recompute lost partitions from lineage relatively straightforwardly. But after shuffle boundaries, intermediate data and execution dependencies become more involved. Shuffle outputs are often written to local disk so downstream stages can fetch them. If those files vanish with the executor, more work has to be redone.

Still, at a beginner level, remember: Spark uses lineage to understand how data was produced and can recompute lost work when failures occur. You will not debug shuffle-service internals in this lesson. You only need the instinct: narrow loss is cheaper to replay than shuffle loss.

Joins, skew, and Spark UI

47. What Happens During a Join?

Suppose employees.join(departments, "department_id"). Imagine employees split across partitions by person, and departments split across partitions by department name. Matching keys are not automatically in the same place. For a normal distributed join, Spark may need matching keys to meet. That can require a shuffle. Data gets reorganized so all Engineering rows meet, all Sales rows meet, all HR rows meet. Then joining happens.

Join inputs, keys not aligned
Employees

Partition 1

101 → Engineering
102 → Sales

Partition 2

103 → Engineering
104 → HR

Departments

Partition 1

Engineering → Building A
HR → Building B

Partition 2

Sales → Building C

After the shuffle, a partition is no longer 'the first chunk of the file.' It is 'the chunk of this join key.' That is why joins change parallelism and why a skewed key (one department with most employees) becomes one fat task.

48. Broadcast Join

Now imagine Employees is 1 TB and Departments is 10 MB. Does it make sense to shuffle the entire 1 TB dataset? No. Instead, Spark can send the small table to all executors. Then each executor joins locally with its portion of the large dataset. This is called Broadcast Join.

Broadcast the small side
                Small Table                        10 MB                                                        ┌────────┼────────┐                                                 ▼        ▼        ▼                                             Executor  Executor  Executor
Broadcast join idea
Small Dataset
      ↓
Broadcast Everywhere
      ↓
Join Locally

Conceptually: small dataset, broadcast everywhere, join locally. This can avoid a large shuffle. Spark may pick this automatically when it knows the table is small. You can also hint it. A later joins lesson practices the call. The architecture point is that broadcasting copies the small table into executor memory so the large side never has to move by key.

49. The Different Levels of Spark Execution

Here's a very useful hierarchy. Let's understand this with one example: df.groupBy("country").count().show().

Levels of Spark execution
APPLICATION

    ↓

JOB

    ↓

STAGE

    ↓

TASK

    ↓

PARTITION

Application: your entire Spark program, My PySpark Program. Job: created when show() is called. Stages: maybe Stage 1 is read and prepare, then shuffle, then Stage 2 aggregates. Tasks: if there are 100 partitions, Stage 1 has 100 tasks. Partition: each task processes a partition. This hierarchy is one of the most useful things to remember.

50. Let's Talk About Spark UI

When you run Spark, there is a UI that helps you understand execution. It can show Jobs, Stages, Tasks, Executors, Storage, SQL, and Environment. Conceptually you click a job, see its stages, then see the tasks inside a stage. Executors have their own tab with memory and disk.

Spark UI shape
SPARK UI                        Jobs             └── Stage 1          ├── Task 1      ├── Task 2      └── Task 3                Executors        ├── Executor 1  ├── Executor 2  └── Executor 3 

When debugging performance, Spark UI is extremely useful. You can see slow tasks, shuffle size, memory usage, skew, and failed tasks. For a data engineer, learning to read Spark UI is a major skill. This tab does not host Spark UI. On a real cluster, it is usually a port on the Driver (or a history server). The nouns in the UI are the nouns in this lesson. That is why we spent so long naming them.

51. What is Data Skew?

Imagine 10 partitions. Normally each is about 10 GB and everything is balanced. But imagine P1 is 90 GB and the others are 1 GB. Then the executor processing P1 runs for a long time, and the other executors finish and wait. This is called data skew. Spark jobs are only as fast as their slowest critical tasks.

Balanced vs skewed partitions
Normally:

P1 → 10 GB
P2 → 10 GB
P3 → 10 GB
P4 → 10 GB
P5 → 10 GB

Skewed:

P1 → 90 GB
P2 → 1 GB
P3 → 1 GB
P4 → 1 GB
P5 → 1 GB

Executor processing P1

████████████████████████

Other Executors

██ Done
██ Done
██ Done
██ Done

A simplified mental model: parallel processing only works if everyone carries similar weight. If one person carries 100kg and everyone else carries 1kg, the entire team waits. A later lesson shows a country histogram and salting. The architecture point is already visible: skew is a partition (often a key after a shuffle) that destroys parallelism.

Run modes

52. Spark Architecture Through a Restaurant Analogy

Sometimes the easiest way to remember everything is an analogy. Imagine a large restaurant. You are the customer: 'I want 10 pizzas.' The Driver is the head chef / manager, and receives the order: what needs to be prepared? The DAG is the cooking plan: prepare dough, add toppings, bake, pack. The Cluster Manager is the resource manager: which kitchens are available, which chefs have capacity? Executors are the chefs. They do actual work.

Partitions are pieces of the order: Pizza 1, Pizza 2, Pizza 3, Pizza 4. Tasks are instructions assigned to chefs: Chef 1 makes Pizza 1, Chef 2 makes Pizza 2. Shuffle is ingredients needing to move between kitchens. This is expensive. Kitchen A sends cheese to Kitchen B, vegetables to Kitchen C. A Stage is a group of operations that can happen before waiting for a major dependency: Stage 1 prepares ingredients, then wait / move ingredients, then Stage 2 bakes pizzas.

If you can retell a Spark job as that restaurant, you understand the architecture. The rest is vocabulary for the same kitchen.

53. Spark Execution Modes

Spark applications can run in different modes. The code you write is largely the same. The processes live in different places.

Local Mode

Everything runs on one machine: your laptop is both Driver and Executor. Example: SparkSession.builder.master("local[*]").getOrCreate(). Great for learning, development, and testing. local[*] means use as many threads as cores. It is still Spark's scheduler and lazy plans. It is not a cluster. It will not prove a terabyte job.

PythonLocal mode
SparkSession.builder \
    .master("local[*]") \
    .getOrCreate()

Cluster Mode

Distributed across multiple machines: a Driver and many Executors. Used for production and large-scale processing. This is the picture the rest of the lesson assumed. When someone says 'we ran it on the cluster,' they mean this, not local[*] on a laptop.

54. Client Mode vs Cluster Mode

You may hear these terms. They are about where the Driver process lives, not about whether executors are distributed. Executors are in the cluster in both cases (unless you are in local mode).

Client Mode

The Driver runs close to where you submit the application. Conceptually: your machine runs the Driver, and the cluster runs the Executors. If your laptop sleeps, the Driver can die and the job dies with it. Notebooks often feel like this.

Client mode
Your Machine            DRIVER                      │                       ▼                   CLUSTER                 Executors   

Cluster Mode

The Driver itself runs inside the cluster. The cluster holds the Driver and the Executors. For production, cluster mode is often preferred because the application is less dependent on your local machine remaining connected. You submit, you go away, the Driver lives next to the workers.

Cluster mode
CLUSTER                 ├── Driver  ├── Executor├── Executor└── Executor

55. Spark Standalone vs YARN vs Kubernetes

These are different cluster managers. The key point for beginners: Spark can run on different infrastructure managers. The Spark programming model remains conceptually similar. You still have a Driver, executors, partitions, and stages. Only the thing that launches those processes changes.

Spark Standalone

Spark's built-in cluster manager. Spark Master, then Workers. Simple. Fine for a dedicated Spark cluster. Less common as the only manager in a company that already has YARN or Kubernetes.

Standalone
Spark         Master          │           Workers

YARN

Common in Hadoop ecosystems. Your Spark application asks YARN, and YARN hands out cluster resources. If your company grew up on HDFS and Hive, you will meet Spark-on-YARN.

YARN
Spark Application

↓

YARN

↓

Cluster Resources

Kubernetes

Modern container orchestration. Spark asks Kubernetes, and Kubernetes runs pods: a Driver pod, then executor pods. Same Spark. Different launcher. You will see this more on cloud platforms that already run everything else on Kubernetes.

Kubernetes
Spark                           ↓                               Kubernetes                      ↓                               Pods                            ├── Driver Pod  ├── Executor Pod├── Executor Pod

Where data lives

56. One Very Important Question: Where Is My Data?

Spark does not necessarily store your data permanently itself. Spark can read from HDFS, S3, Azure Data Lake, Google Cloud Storage, databases, Kafka, local storage, Parquet files, CSV files, and Delta tables. Spark is primarily a processing engine. The lake, the warehouse, and the stream live elsewhere. Spark is the engine that scans them and writes results back.

Spark reads, transforms, writes
                   DATA SOURCES                                                           ┌──────────┬───────────┬───────────┐                                                ▼          ▼           ▼                                                           S3        Kafka       Database                                                       │          │           │                                                            └──────────┼───────────┘                             │                                         ▼                                                                                 SPARK                                                                                 │                                                                                   ▼                                                                              TRANSFORMATION                                                                           │                                                                                   ▼                                                                               OUTPUT DATA                  

57. Why Parquet Is Better Than CSV for Spark

Consider CSV: name, age, salary, department as a row of text. If you need only name and salary, Spark may still have to parse through row-oriented text. Parquet is columnar. Conceptually: Column A is Name, Column B is Age, Column C is Salary, Column D is Department. If Spark only needs Name and Salary, it can often read only the relevant columns.

This works beautifully with column pruning. That's one reason columnar formats are widely used in big-data systems. Combined with predicate pushdown, a Parquet lake can skip both columns you do not select and row groups that cannot match your filter. CSV gives Catalyst much less to work with. You will still meet CSV (vendors love it). You should not choose it for a lake you control.

The picture to memorize

58. The Entire Architecture as a Story

Let's put everything together. You write a Parquet read, a filter, a groupBy sum, then show. Now imagine Spark saying the following out loud. This is the same pipeline as section 28, told as a checklist you can reuse on any job.

PythonThe story's source code
result = (
    spark.read.parquet("sales")
    .filter("amount > 100")
    .groupBy("country")
    .sum("amount")
)

result.show()

Step 1

The user wrote a PySpark application. Spark creates the Driver. That process will own the session, the plan, and the schedulers for the rest of this story.

Step 1
Create Driver

Step 2

The user wants Spark resources. The Driver asks the Cluster Manager. Nothing has scanned sales yet. This is only 'may I have machines?'

Step 2
Driver → Cluster Manager

Step 3

Start executors. Until these processes register, there is nobody to run tasks.

Step 3
Executor 1
Executor 2
Executor 3

Step 4

The user described transformations. Spark records Read, then Filter, then GroupBy, then Sum. It does not execute yet. The plan is sitting on the Driver.

Step 4
Read

↓

Filter

↓

GroupBy

↓

Sum

Step 5

The user called show(). That is the ACTION. Now execution begins. A job appears because something finally demanded a result.

Step 5
ACTION

Step 6

Create the logical plan. WHAT needs to happen? Still no executor cores chewing Parquet. This is the query in operator form.

Step 7

Optimize it. Catalyst looks for better ways: push the amount filter toward the scan, read only the columns the sum needs, maybe partial-aggregate before the shuffle.

Step 8

Create the physical execution plan. HOW should this happen? Which scan operator, which aggregate, which shuffle.

Step 9

Build the DAG so Spark can see dependencies, not only a flat list of steps.

Step 9
Read
 ↓
Filter
 ↓
GroupBy
 ↓
Aggregate

Step 10

Where are shuffle boundaries? Read and Filter can stay together. GroupBy needs keys to meet. That is the line.

Step 10
Read
 ↓
Filter

-------- SHUFFLE --------

GroupBy
 ↓
Aggregate

Step 11

Create stages from that line. Stage 1 is Read + Filter. Stage 2 is GroupBy + Aggregate. Two stages, one job.

Step 11
Stage 1

Read + Filter

Stage 2

GroupBy + Aggregate

Step 12

Create tasks based on partitions. If Stage 1 has four partitions, Spark creates about four tasks. Partition 1 becomes Task 1, Partition 2 becomes Task 2, and so on.

Step 12
Stage 1

Partition 1 → Task 1

Partition 2 → Task 2

Partition 3 → Task 3

Partition 4 → Task 4

Step 13

Send tasks to executors. An executor with free cores takes the next pending task. The Driver is coordinating, not scanning sales itself.

Step 13
Executor 1 → Task 1

Executor 2 → Task 2

Executor 3 → Task 3

Executor 1 → Task 4

Step 14

Process partitions in parallel. Each executor works its slice. Filters run locally. No country totals exist yet, because the keys have not met.

Step 14
Executor 1 ██████████

Executor 2 ██████████

Executor 3 ██████████

Step 15

Shuffle data by country. India rows from every partition need to sit together. Same for USA and UK. This is the expensive part of the job.

Step 15
India → Same group

USA → Same group

UK → Same group

Step 16

Perform the final aggregation on those grouped partitions. Each country becomes one total. Stage 2's tasks are fewer or at least differently shaped than Stage 1, because they follow the shuffled keys, not the original files.

Step 16
India → Total

USA → Total

UK → Total

Step 17

Return the result. show() prints a small sample on the Driver. Done. That is the entire architecture, one pipeline, start to finish.

PythonStep 17
result.show()

59. The Architecture Diagram to Memorize

If you remember only one thing, remember this flow, plus the two small hierarchies under it. Everything in Spark UI maps onto these three pictures.

The one flow to memorize
                    PYSPARK CODE                             │                                   ▼                            SPARK SESSION                              │                                   ▼                                DRIVER                                 │                                   │                            Creates DAG                                │                                   ▼                           OPTIMIZES PLAN                              │                                   ▼                             CREATES JOB                               │                                   ▼                             SPLITS INTO                            STAGES                                 │                            (Shuffle Boundary)                         │                                   ▼                             CREATES TASKS                             │                                   ▼                            CLUSTER MANAGER                            │                                   ▼                              EXECUTORS                                │                                   ▼                             PROCESS DATA                         IN PARALLEL                              │                                   ▼                                 RESULT       

And alongside it:

Data to work
DATASET

   ↓

PARTITIONS

   ↓

TASKS

   ↓

EXECUTORS

And:

Demand to work
ACTION

   ↓

JOB

   ↓

STAGES

   ↓

TASKS

60. The 15 Concepts You Must Know

If you're beginning PySpark, master these in this order. They are the same 15 ideas this lesson already walked. Treat the list as a self-check, not as new material.

Fifteen names. One engine.

#ConceptIn one sentence
1Why Spark?Because data becomes too large for one machine.
2Distributed ComputingMultiple machines process data together.
3DriverThe brain that coordinates execution.
4Cluster ManagerAllocates resources.
5ExecutorsPerform actual computation.
6PartitionsPieces of distributed data.
7TasksSmallest units of execution, typically working on partitions.
8TransformationsDescribe what should happen: filter, select, groupBy.
9ActionsTrigger execution: show, count, write.
10Lazy EvaluationSpark waits before executing so it can optimize.
11DAGThe dependency graph of computations.
12ShuffleRedistribution of data across partitions.
13StagesGroups of tasks separated largely by shuffle boundaries.
14Catalyst OptimizerOptimizes DataFrame and SQL execution plans.
15LineageSpark remembers how data was created, helping with fault tolerance.

Final Mental Model

Whenever you see PySpark code, think: I am not directly processing data. I am describing a distributed computation. You write df.filter(...).groupBy(...).agg(...). But internally, you describe computation, Spark waits, an action is called, the Driver creates a plan, Catalyst optimizes it, a DAG is created, the DAG is split into stages, stages split into tasks, tasks are assigned to executors, executors process partitions in parallel, shuffle happens if required, and results are produced.

What that one-liner actually does
You describe computation
        ↓
Spark waits
        ↓
Action is called
        ↓
Driver creates plan
        ↓
Catalyst optimizes it
        ↓
DAG is created
        ↓
DAG split into stages
        ↓
Stages split into tasks
        ↓
Tasks assigned to executors
        ↓
Executors process partitions in parallel
        ↓
Shuffle happens if required
        ↓
Results are produced

The one-line definition I want you to remember

PySpark in one sentence

PySpark lets you describe data processing using Python, while Apache Spark's distributed engine plans, optimizes, and executes that processing across partitions and executors in a cluster.

If you genuinely understand the flow above, you have the foundation needed before moving into DataFrames, joins, partitioning, performance tuning, Spark UI, AQE, caching, skew, and production-level Spark engineering. The next lessons in this track are that next step, with a small orders table you can actually run.

What comes next

The next lesson is the short hands-on version of 'when do you need Spark?' You will load spark.table("orders"), count, and take a sample. The architecture stays the same. The dataset gets small on purpose so you can practice the verbs.

Practice

Run Sample to print the four architecture names and a lazy filter plus groupBy. Then complete Exercise: fill the pieces dict (driver, executors, partition, task), filter orders where order_total > 50, print the count action, and assign a groupBy("order_status").count() to result.

You are proving two things: you can name who does the work, and you can tell a transformation from an action on a real DataFrame in this tab.

Practicals · load into the editor

After you read the theory, run these in the pane on the right. They execute in this tab, no cluster.

Practice this

Same ideas as interview drills. These challenges open in the studio with a dataset and tests already set up.

  • Add Derived ColumnInterview-style drill: Add full_name by joining first_name and last_name with a space.Studiobeginnerpyspark12 min
  • Explode Array ColumnInterview-style drill: Explode tags arrays into one row per tag.Studiobeginnerpyspark12 min
  • Filter and Select ColumnsInterview-style drill: Keep purchase rows with amount > 100 and project user_id, amount.Studiobeginnerpyspark12 min
Rate:
Was this useful?
PySpark for Distributed ProcessingWhen do you need Spark?