Spark looks like a library you call: read a table, join it, aggregate, write. Underneath, one call to write or count sets off a chain of translations. Your query becomes a logical plan, the optimiser rewrites it, a physical plan picks concrete algorithms, the scheduler cuts that plan into stages, each stage becomes a set of tasks, and those tasks run on executor JVMs scattered across a cluster, moving data between each other through shuffle files on local disk.

Almost every performance problem and failure in Spark lives at one of those translations. This article walks the whole path from first principles, with a worked example sized like a real job, and ends with a checklist you can apply to the next job you tune. Version notes refer to Spark 3.x and 4.x; Mesos support was deprecated in 3.2 and removed in 4.0, so the cluster managers that matter are YARN, Kubernetes and standalone.

From one action to thousands of tasks: the Spark execution pipelineYour codeDataFrame / SQLCatalystanalyse + optimisePhysical plancodegen, exchangesAQEre-plan at stage boundariesDAGSchedulerjobs -> stages at shufflesTaskSchedulerTaskSetManager, localitySchedulerBackendYARN / K8s / standaloneaction triggers a jobdriver process (one JVM): planning, scheduling, MapOutputTracker, result collectionExecutor 1Executor 2Executor NtasktaskBlockManager + shuffle filestasktaskBlockManager + shuffle filestask slots = coresspills to local disklaunch tasksfetchfetchStorage: Parquet / Delta / Iceberg on HDFS, S3, GCSinput splits feed stage 1; final stage writes files
The execution pipeline: planning and scheduling happen in the driver, work happens in executor task slots, and shuffle files on executor disks connect one stage to the next.

The execution model in three facts

Three facts explain most of Spark's behaviour. First, transformations are lazy. filter, join and groupBy only build a plan; nothing runs until an action such as count, collect or a write. Each action creates one job, which is why a notebook that calls count three times scans the input three times unless you cache.

Second, data is split into partitions and one task processes one partition. Parallelism is therefore the number of partitions in a stage, capped by the number of task slots, which is executors times cores per executor. A stage with 16,000 tasks on 200 slots runs in 80 waves; a stage with 12 tasks on 200 slots leaves 188 cores idle no matter how big the cluster is.

Third, dependencies come in two kinds. A narrow dependency means each output partition needs one input partition: filters, projections, map-side work. Narrow operators are fused into the same task and never touch the network. A wide dependency means an output partition needs data from many input partitions: a sort-merge join, a group-by, a repartition. Wide dependencies require a shuffle, and every shuffle is a stage boundary.

The execution model in three facts

Three facts explain most of Spark's behaviour. First, transformations are lazy. filter, join and groupBy only build a plan; nothing runs until an action such as count, collect or a write. Each action creates one job, which is why a notebook that calls count three times scans the input three times unless you cache.

Second, data is split into partitions and one task processes one partition. Parallelism is therefore the number of partitions in a stage, capped by the number of task slots, which is executors times cores per executor. A stage with 16,000 tasks on 200 slots runs in 80 waves; a stage with 12 tasks on 200 slots leaves 188 cores idle no matter how big the cluster is.

Third, dependencies come in two kinds. A narrow dependency means each output partition needs one input partition: filters, projections, map-side work. Narrow operators are fused into the same task and never touch the network. A wide dependency means an output partition needs data from many input partitions: a sort-merge join, a group-by, a repartition. Wide dependencies require a shuffle, and every shuffle is a stage boundary.

From query to physical plan

When the action fires, the driver hands the query to Catalyst. The analyser resolves table and column names against the catalog and checks types. The optimiser applies rules: it pushes filters down towards the scan, prunes columns you never use, folds constants and simplifies expressions. With table statistics collected, cost-based optimisation can also reorder joins. The planner then picks physical operators: a broadcast hash join when one side is estimated below spark.sql.autoBroadcastJoinThreshold (10 MB by default), otherwise usually a sort-merge join; a hash aggregate split into partial and final halves; and an Exchange node wherever data has to be redistributed.

Whole-stage code generation then collapses chains of narrow operators into one generated Java function per fused region, so a filter, projection and partial aggregate run as a tight loop over rows instead of a chain of virtual calls. In explain() output these fused regions are the operators prefixed with *(n). Reading this plan is the single most useful Spark skill:

df = (spark.read.parquet("s3://lake/orders")
        .filter("order_date >= '2026-09-01'")
        .join(spark.read.parquet("s3://lake/customers"), "customer_id")
        .groupBy("country")
        .agg({"amount": "sum"}))

df.explain("formatted")
# Look for:
#   PushedFilters: [GreaterThanOrEqual(order_date,2026-09-01)]  -> filter reached the scan
#   BroadcastHashJoin vs SortMergeJoin                         -> join algorithm chosen
#   Exchange hashpartitioning(country, 200)                    -> a shuffle, i.e. a stage cut
#   AdaptiveSparkPlan isFinalPlan=false                        -> AQE may still change it

Adaptive Query Execution, on by default since Spark 3.2, adds a loop to this pipeline. The plan is run stage by stage, and at each Exchange AQE reads the real map output sizes before planning the next stage. It can coalesce many small shuffle partitions into fewer reasonably sized ones, switch a sort-merge join to a broadcast join once it sees that one side is actually small, and split skewed partitions of a sort-merge join into several tasks. More on its limits in the AQE deep dive.

Scheduling: jobs, stages and task sets

The physical plan reaches the scheduler as an RDD graph. The DAGScheduler walks it backwards from the final RDD and cuts a new stage at every shuffle dependency. Stages that end in a shuffle are shuffle map stages; the last stage of a job is the result stage. It submits a stage only when all its parent stages have finished, and it records where each map task's output lives in the MapOutputTracker so reducers know whom to fetch from. In simplified pseudocode:

def submit_job(final_rdd, action):
    final_stage = create_result_stage(final_rdd)        # recursively creates parent stages
    submit_stage(final_stage)

def create_stage(rdd):
    parents = []
    for dep in walk_narrow_ancestors(rdd):              # narrow deps stay in this stage
        if isinstance(dep, ShuffleDependency):
            parents.append(get_or_create_shuffle_map_stage(dep))   # reused if already computed
    return Stage(rdd, parents)

def submit_stage(stage):
    missing = [p for p in stage.parents if not p.is_available()]
    if missing:
        for p in missing:
            submit_stage(p)
        waiting.add(stage)                              # resumes when parents finish
    else:
        tasks = [make_task(stage, part) for part in stage.missing_partitions()]
        task_scheduler.submit(TaskSet(stage, tasks))

Notice get_or_create: if an earlier job already produced a shuffle whose files still exist, the stage is skipped, which is why the Spark UI shows grey 'skipped' stages.

The TaskScheduler receives one TaskSet per stage. A TaskSetManager tracks each task's attempts and tries to place it near its data, falling back from process-local to node-local to rack-local to any after spark.locality.wait (3 seconds by default) at each level. A failed task is retried up to spark.task.maxFailures times (4 by default) before the whole job fails. With speculation enabled, slow tasks get a duplicate attempt and the first to finish wins. The SchedulerBackend is the adapter to the cluster manager: it asks YARN, Kubernetes or the standalone master for executors and sends serialised tasks to them. Stages and tasks covers the UI view of this in detail.

Inside an executor

An executor is a JVM with a fixed number of cores and a fixed heap. Each core runs one task at a time. A task deserialises its closure, reads its partition (an input split, a cached block, or shuffle blocks fetched from other executors), runs the generated code, and writes its output either as shuffle files or, in the final stage, as result data or output files.

Executor memory is split by the unified memory manager. After a fixed 300 MB reservation, spark.memory.fraction (0.6) of the remaining heap is shared between execution memory (sort buffers, hash tables for joins and aggregations) and storage memory (cached blocks, broadcast variables). spark.memory.storageFraction (0.5) marks the part of that region where cached blocks are protected from eviction; outside it, execution can evict cache and either side can borrow from the other. The remaining 40 percent holds user data structures and Spark internals. Native buffers, thread stacks and Python workers without their own limit live outside the heap, in spark.executor.memoryOverhead (the larger of 384 MB and 10 percent of executor memory). The container request is heap plus overhead plus any explicit spark.memory.offHeap.size and spark.executor.pyspark.memory; exceeding it gets the container killed, which looks very different from a Java OutOfMemoryError.

When execution memory runs out, sorts and aggregations spill to local disk rather than failing. Spill is the most common silent cost in Spark: the job finishes, but each spilled task writes and rereads its data. The task table in the UI reports spill per task. Detail is in Spark memory management.

Shuffle: the boundary between stages

Where the DAGScheduler cuts stages: every wide dependency is an Exchangescan orders16,000 splitsfilter + projectnarrow, same taskscan customers50 MBbroadcastbuilt once, shippedbroadcast hash join + partial aggstage 1: map side, writes shuffleExchangehashpartitioningfinal agg + writestage 2: reduce sideNarrow operators pipeline inside one task. An Exchange forces map output to disk and a new stage.AQE inspects the Exchange's real sizes before stage 2 starts and may coalesce or split partitions.
Stage boundaries in the worked example: the broadcast join and partial aggregate pipeline in one stage; the group-by forces an Exchange.

A shuffle has a write side and a read side. On the write side, each map task partitions its output by the target reducer, sorts it by partition id (and by key when needed), and writes one data file plus an index file to local disk. When the number of reducers is small (at most spark.shuffle.sort.bypassMergeThreshold, 200 by default) and there is no map-side combine, the writer skips the sort and writes per-reducer files that it concatenates.

On the read side, each reduce task asks the MapOutputTracker for the locations of its blocks and fetches them over the network from every map task's executor, in parallel and with a limit on bytes in flight. If the executor that wrote a block has died, the fetch fails with a FetchFailedException. The scheduler does not just retry the reduce task: it marks the map output as lost, resubmits the parent stage to recompute the missing map partitions, then reruns the reducers. A stage that keeps failing this way is aborted after spark.stage.maxConsecutiveAttempts (4).

The external shuffle service lets shuffle files outlive their executor, which makes dynamic allocation safe; see the shuffle deep dive.

Worked example: a 2 TB join and aggregation

Take the query above against realistic sizes: orders is 2 TB of Parquet, the date filter keeps about a fifth of it, customers is 50 MB, and the result is roughly 200 countries. The cluster has 50 executors with 4 cores and 16 GB heap each, so 200 task slots.

Stage 1. Scan splits are sized by spark.sql.files.maxPartitionBytes (128 MB), so 2 TB becomes about 16,000 tasks, or 80 waves on 200 slots. Partition pruning on order_date, if the table is partitioned by date, cuts that to about 3,200 tasks before anything runs. At 50 MB, customers is above the 10 MB default threshold, so without help Spark plans a sort-merge join and shuffles 400 GB of orders (AQE may rescue it at runtime, but only after paying for the shuffle's map side). Raising the broadcast threshold to 64 MB, or adding broadcast(customers), turns the join into a broadcast hash join inside stage 1. The partial aggregate then reduces each task's output to at most 200 rows.

Stage 2. The Exchange now moves a few hundred thousand rows instead of 400 GB. With spark.sql.shuffle.partitions at its default of 200, AQE sees tiny partitions and coalesces them into a handful of tasks. The whole job is dominated by the scan.

Now suppose the join really is large to large, and 400 GB must cross the Exchange. Two hundred partitions means 2 GB per reducer, far more than one task's share of execution memory, so every reducer spills. The fix is counter-intuitive: AQE can only coalesce partitions downwards; it does not raise the count. Set the initial number high, for example 4,000 for about 100 MB each, and let AQE merge them towards spark.sql.adaptive.advisoryPartitionSizeInBytes (64 MB):

spark.conf.set("spark.sql.adaptive.enabled", "true")                     # default since 3.2
spark.conf.set("spark.sql.shuffle.partitions", "4000")                   # start high
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128m")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")            # split hot keys
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", str(64 * 1024 * 1024))

Then verify in the Spark UI: check the final plan's partition counts and compare median and maximum task duration per stage. A maximum ten times the median is skew, and AQE only splits skew in sort-merge joins, not in a plain group-by.

Failure modes and their signatures

SymptomWhere it livesWhat usually fixes it
Driver OutOfMemoryErrorcollect, toPandas, huge broadcast, very wide plansWrite results instead of collecting; cap spark.driver.maxResultSize; smaller broadcasts
Executor killed by YARN or Kubernetesoff-heap, Python workers, native buffers exceed overheadRaise spark.executor.memoryOverhead; fewer cores per executor
Executor Java OOM in a join or aggregateone partition far bigger than the restMore shuffle partitions, AQE skew join, salting hot keys
FetchFailedException and repeated stage retrieslost executor took its shuffle files with itExternal shuffle service or decommissioning; stop preempting workers mid-shuffle
One task runs for an hour while others take a minuteskewed key or an unsplittable input fileInspect key distribution; split files; salt; isolate the hot key
Thousands of tasks taking milliseconds eachtoo many small files or partitionsCompact input; let AQE coalesce; avoid repartition before write
Disk full on workersshuffle and spill on small local volumesBigger or more local disks; reduce shuffle volume with earlier filters

Trade-offs

  • Few large executors versus many small ones. Large executors share broadcast variables and cache across more cores but suffer longer garbage-collection pauses; four to five cores per executor is a common middle ground.
  • Broadcast versus shuffle joins. Broadcasting removes a shuffle but copies the table to every executor and builds it on the driver first; a 1 GB broadcast to 200 executors is 200 GB of memory.
  • Caching. It saves recomputation but takes memory from execution; cache only what several actions reuse.
  • Dynamic allocation. It saves money on bursty jobs, but executors released mid-job take cached blocks with them, and without shuffle tracking they take shuffle files too.

What to do next

  1. Run explain("formatted") on your slowest job and mark every Exchange, join type and pushed filter.
  2. Open its stage page and record task count, median and max duration, shuffle read and spill for each stage.
  3. Check that no stage has fewer tasks than your cluster has slots, unless the data is genuinely tiny.
  4. Set a high initial shuffle partition count and let AQE coalesce; compare spill before and after.
  5. Broadcast every dimension table that comfortably fits in executor memory, and confirm the plan changed.
  6. Look for max-to-median task ratios above about five and treat them as skew to investigate.
  7. Size executor overhead explicitly for PySpark and native libraries before raising heap.
  8. Make writes idempotent so retries, speculation and stage resubmission cannot duplicate output.
Key takeaway: Spark runs your query as a chain of translations: plan, stages cut at shuffles, tasks bound to partitions, executors with a fixed memory budget. Tune by reading the plan, counting tasks per stage, measuring shuffle and spill, and fixing the one stage that dominates, rather than by turning knobs at random.