Every Spark program is a graph. Transformations such as map, filter and join only record how each dataset derives from its parents; nothing runs until an action such as count or a write asks for a result. At that moment a component in the driver, the DAGScheduler, takes the graph of datasets, cuts it into stages at every shuffle, decides which stages already have their outputs, submits the rest in dependency order, and recovers when a shuffle file goes missing. The name refers to this graph: directed because data flows from parents to children, acyclic because a dataset can never depend on itself.

This article follows a job through the DAGScheduler using the names in Spark's source, so that what you see in the Spark UI and in logs stops being mysterious. Task-level scheduling, locality and speculation are covered in Spark stages and tasks; how lineage enables recomputation is covered in Spark RDD lineage.

Advertisement

Two graphs, not one

It helps to separate two graphs. The lineage graph has one node per RDD and one edge per dependency. Dependencies are either narrow, where each child partition reads a bounded set of parent partitions (map, filter, union, a join of co-partitioned data), or shuffle dependencies, where each child partition may read from every parent partition (reduceByKey, groupByKey, a join of differently partitioned data, repartition).

The stage graph is what the DAGScheduler derives from it. A stage is a maximal chain of narrow dependencies that can run as one pipelined task per partition without moving data between machines. Each shuffle dependency ends one stage and starts another. Stages that produce shuffle output are ShuffleMapStages; the final stage that computes the action's result is the ResultStage. Each job has exactly one ResultStage and any number of map stages. Most actions run a single job, though some, such as take, zipWithIndex or the sampling behind sortByKey, run extra ones.

The event loop

The DAGScheduler is single-threaded by design. Everything that changes its state arrives as an event on a queue processed by DAGSchedulerEventProcessLoop: JobSubmitted when an action runs, CompletionEvent when a task finishes or fails, executor added and lost events, cancellations, and ResubmitFailedStages after a fetch failure. Because one thread handles every event in order, the scheduler's maps of jobs, stages and shuffles need no locks, and a slow handler delays everything behind it. A driver whose event queue backs up shows as jobs that appear to hang between stages while executors sit idle.

The flow for one action is: SparkContext.runJob calls the DAGScheduler, which posts JobSubmitted; the handler builds the ResultStage and its ancestors, then submits the ResultStage; submitting a stage first submits any parent stage whose output is missing; a stage with no missing parents becomes a TaskSet handed to the TaskScheduler, which places tasks on executors. As each task completes, a CompletionEvent comes back, and when a map stage finishes, its waiting children are submitted.

Advertisement

Building the stage graph

One action, one job: the DAGScheduler cuts the RDD graph at shufflestextFileordersmapto (customer, amount)reduceByKeyshuffletextFilecustomersmapto (customer, region)joincogroup on customernarrowshufflecollectresultStage 0: map stageorders, writes shuffle 0Stage 1: map stagecustomers, writes shuffle 1Stage 2: ResultStagereduceByKey reads shuffle 0, join reads shuffle 1Stages 0 and 1 have no parents and run in parallel;stage 2 waits for both. The join reuses the reduceByKeypartitioner, so that side needs no extra shuffle.Narrow dependencies are pipelined inside a stage; every ShuffleDependency becomes a stage boundary.
A job that aggregates orders per customer and joins them with customer regions. Two shuffle dependencies produce two map stages and one result stage.

Stage construction starts at the final RDD and walks backwards. createResultStage asks getOrCreateParentStages for the stages that feed it, which finds the nearest shuffle dependencies by traversing narrow edges and calls getOrCreateShuffleMapStage for each. That function first looks the shuffle up in shuffleIdToMapStage; only if no stage exists for that shuffle does it create one, recursively building its ancestors. The core of the algorithm fits in a few lines:

def parent_stages(rdd, first_job):
    stages, seen, todo = [], set(), [rdd]
    while todo:                                   # walk narrow edges until a shuffle
        r = todo.pop()
        if r in seen:
            continue
        seen.add(r)
        for dep in r.dependencies:
            if isinstance(dep, ShuffleDependency):
                stages.append(get_or_create_map_stage(dep, first_job))
            else:
                todo.append(dep.rdd)              # narrow: same stage, keep walking
    return stages

def get_or_create_map_stage(dep, first_job):
    if dep.shuffle_id in shuffle_id_to_map_stage:  # reuse across jobs
        return shuffle_id_to_map_stage[dep.shuffle_id]
    parents = parent_stages(dep.rdd, first_job)
    stage = ShuffleMapStage(dep.rdd, dep, parents, first_job)
    shuffle_id_to_map_stage[dep.shuffle_id] = stage
    return stage

def submit_stage(stage):
    missing = [p for p in stage.parents if not p.is_available()]
    if not missing:
        task_scheduler.submit(TaskSet(stage, stage.missing_partitions()))
    else:
        waiting.add(stage)
        for p in missing:
            submit_stage(p)

Two consequences follow. A stage is keyed by its shuffle, so two jobs that share a shuffle share a stage object. And missing_partitions is computed per stage, so a stage resubmitted after a failure runs tasks only for the partitions whose output is gone.

Map outputs and skipped stages

When a shuffle map task finishes, it reports a MapStatus: which executor holds its output and how large each reduce partition's slice is. The driver's MapOutputTracker stores these per shuffle. A ShuffleMapStage is available when every map partition has a registered output, and reduce tasks ask the tracker where to fetch from.

This is why the Spark UI shows skipped stages. Run a second action on an RDD that ends in an existing shuffle and the new job's graph includes the old map stage, but its outputs are still registered, so the scheduler never submits it. The UI draws the stage in grey and labels it skipped. Shuffle files act as an implicit cache. They live on executor disks, or with an external shuffle service, until the shuffle is cleaned up when its RDD is garbage-collected on the driver.

Worked example: reading the DAG

Here is the job from the diagram in the Scala shell. Run it and open the job in the UI alongside.

val orders = sc.textFile("orders.csv", 8)
  .map(_.split(",")).map(f => (f(1), f(2).toDouble))     // (customer, amount)
val totals = orders.reduceByKey(_ + _)                    // shuffle 0, HashPartitioner(8)
val regions = sc.textFile("customers.csv", 4)
  .map(_.split(",")).map(f => (f(0), f(3)))               // (customer, region)
val joined = totals.join(regions)                         // reuses totals' partitioner
println(joined.toDebugString)
joined.collect()    // job 0: stages 0, 1 and 2
joined.count()      // job 1: stages 0 and 1 skipped, only the result stage runs

The interesting line is the join. totals already has a hash partitioner with 8 partitions, so the join adopts it: the totals side becomes a narrow dependency and only regions is shuffled. The job has two map stages and a result stage, not three map stages. Give the join a different partition count and a third shuffle appears, which is the cheapest optimisation in this example to lose by accident.

To read toDebugString, start with a smaller job. A word count in the Scala shell prints a shape like this:

(2) ShuffledRDD[4] at reduceByKey at <console>:25 []
 +-(2) MapPartitionsRDD[3] at map at <console>:25 []
    |  MapPartitionsRDD[2] at flatMap at <console>:25 []
    |  README.md MapPartitionsRDD[1] at textFile at <console>:25 []
    |  README.md HadoopRDD[0] at textFile at <console>:25 []

The number in parentheses is the partition count. Each +- marks a shuffle boundary: everything indented below it belongs to the parent stage, here the pipelined read, flatMap and map that write the shuffle. Ids and call sites vary by version. In the join job the printout has two such boundaries, one per map stage. In the UI the first action shows three stages, with the two map stages running in parallel; the second action shows the same two map stages in grey as skipped, because their shuffle outputs are still registered, and only the result stage runs.

When a shuffle file goes missing

A reduce task that cannot fetch a map output, because the executor died, the node was preempted or the file was deleted, fails with FetchFailed. This is not an ordinary task failure. The scheduler marks the reduce stage as failed, unregisters the lost map outputs from the MapOutputTracker, adds both stages to failedStages, and posts ResubmitFailedStages after a short delay so that several fetch failures from the same loss are handled together. On resubmission the map stage runs only its missing partitions, then the reduce stage runs again.

Repeated fetch failures are bounded by spark.stage.maxConsecutiveAttempts (default 4); past that the job aborts. Recovery depends on lineage reaching back to durable data, and an external shuffle service or decommissioning that migrates shuffle blocks reduces how often it is needed. One subtle case: if a map stage's output is not deterministic, for example repartition distributing rows round-robin over unsorted input, recomputing some partitions can route rows differently than the first run did. Spark detects such indeterminate stages and reruns their downstream stages fully, or fails the job if a result was already partially committed. Sorting before repartition or repartitioning by key avoids the problem.

DataFrames, exchanges and AQE

DataFrame and SQL queries reach the same scheduler. Catalyst produces a physical plan, and every Exchange node in it becomes a shuffle dependency in the RDDs the plan executes. With adaptive query execution enabled, Spark does not submit the whole graph at once. It cuts the plan into query stages at exchanges, materialises the leaves first, looks at the real shuffle statistics, and re-optimises the rest of the plan, coalescing partitions, switching join strategies or splitting skewed partitions, before submitting the next stages. That is why one SQL query can show several jobs in the UI. The mechanics are covered in Spark AQE architecture, and the shuffle itself in Spark shuffle architecture.

Long lineages

Iterative programs that reassign a dataset in a loop, such as graph algorithms, iterative refinements or a streaming state built by repeated unions, grow the lineage on every iteration. Planning time on the driver grows with it, recursive traversal of a very deep graph can overflow the stack, and recovery after a failure may recompute many iterations. Caching does not help here: cache keeps data but keeps the lineage too.

spark.sparkContext.setCheckpointDir("hdfs:///tmp/ckpt")
for i in range(50):
    ranks = update(ranks, links)          # ranks is a DataFrame
    if i % 10 == 9:
        # eager by default: materialises now and returns a DataFrame with a short lineage
        ranks = ranks.localCheckpoint() if fast_and_fragile else ranks.checkpoint()

# for an RDD, checkpoint() returns nothing: mark it, then run an action
rdd.checkpoint()
rdd.count()

checkpoint writes to reliable storage and replaces the lineage with a read of those files; localCheckpoint keeps blocks on executors, which is faster but loses the data if an executor dies. Checkpoint every few iterations and check that the plan or toDebugString stays short.

Failure modes

SymptomLikely causeFix
Job aborts after repeated FetchFailedExecutors lost during shuffle, often spot preemption or memory killsExternal shuffle service or decommissioning; fix executor memory
Same stage reruns several timesLost map outputs resubmitted repeatedlyFind the executor loss in the driver log; stabilise nodes
Driver busy, executors idle between stagesEvent loop backed up by huge task counts or listenersFewer, larger partitions; remove heavy listeners
Slow planning grows each iterationUnbounded lineage in a loopCheckpoint every N iterations
Wrong or duplicated rows after a retryNondeterministic map outputSort before repartition, or repartition by key
Expensive stage recomputed per actionNo shuffle or cache to reusePersist the shared dataset or reuse the shuffle

Trade-offs

Recomputation from lineage makes Spark cheap in the normal case and expensive when failures cluster: no data is replicated on the way, but a lost shuffle costs a stage rerun. Wide transformations give the scheduler stage boundaries where it can recover and re-optimise, while each one also costs a full shuffle. Checkpointing trades write I/O for bounded recovery and planning. The Spark UI's DAG view, described in the Spark UI guide, is the fastest way to see which side of these trade-offs a job is on.

What to do next

  1. Print toDebugString for one real job and mark every shuffle boundary.
  2. Open the same job in the UI and match each stage to a boundary; note which are skipped and why.
  3. Search the driver log of your flakiest job for FetchFailed and count stage attempts.
  4. Enable an external shuffle service or decommissioning wherever executors are preemptible.
  5. Add a checkpoint every few iterations to any loop that reassigns an RDD or DataFrame.
  6. Replace round-robin repartitions on unsorted data with key-based ones.
Key takeaway: The DAGScheduler turns each action into a job, cuts the lineage graph into stages at every shuffle, reuses stages whose shuffle outputs are still registered, and submits the rest parent first as task sets. Read toDebugString and the UI with that model in mind: shuffle boundaries are stages, skipped stages are reused outputs, FetchFailed means a parent stage reruns its missing partitions, and long or nondeterministic lineages are the cases where you must step in with checkpoints and deterministic partitioning.