Jobs, Stages and Shuffles

Jobs, Stages and Shuffle Boundaries

An action triggers jobs; each job is cut at shuffle boundaries into stages, and each stage runs one task per partition. Adaptive execution (Adaptive Query Execution) runs one shuffle stage at a time, so a single query can launch several jobs.

Running the query and listing its jobs, stages and tasks
sc.setJobGroup("genre", "BookNest revenue by genre")
genre_revenue.show()
st = sc.statusTracker()
for job_id in sorted(st.getJobIdsForGroup("genre")):
    stages = st.getJobInfo(job_id).stageIds
    tasks = [st.getStageInfo(s).numTasks if st.getStageInfo(s) else "skipped" for s in stages]
    print(f"job {job_id}: stages {list(stages)} tasks {tasks}")
Output
+---------------+----------+
|          genre|   revenue|
+---------------+----------+
|     Technology|8878336.00|
|        Cooking|6791664.00|
...
job 2: stages [2] tasks [1]
job 3: stages [3] tasks [4]
job 4: stages [4, 5] tasks [4, 1]

Job 2 reads the six-row catalog and broadcasts it (one task). Job 3 scans the four Parquet 129 files, filters, joins and pre-aggregates in four tasks, then writes shuffle files keyed by genre. Job 4 needs stage 4's output, finds it already on disk and skips it, then runs stage 5 as a single task. The Spark 129 UI draws the same thing.

The Spark UI's DAG visualization for job 4: stage 4 skipped, stage 5 reading its shuffle output
The Spark UI's DAG visualization for job 4: stage 4 skipped, stage 5 reading its shuffle output