4. From Application to Task: How Spark Executes Work
Job·Stage·Task로 이어지는 Spark 실행 모델
- Haram Lee
- 2026-08-16
- studies / Topics / Spark
- 이번 장에서는 그 설계도가 실제 클러스터에서 실행 가능한 단위로 어떻게 변환되는지 설명
flowchart TD code["User Code"] --> plan["RDD / DataFrame Plan"] plan --> action["Action"] action --> job["Job"] job --> stage["Stage"] stage --> task["Task"] task --> executor["Executor"]
핵심 흐름
Application
└─ Job
└─ Stage
└─ Task
└─ Partition 처리다만 이 계층은 단순히 크기순으로 포함되는 관계만은 아님.
- 하나의 Application에는 여러 Job이 존재할 수 있음
- 하나의 Job은 여러 Stage를 필요로 할 수 있음
- 하나의 Stage에는 여러 Task가 존재함
- 각 Task는 일반적으로 해당 Stage의 Partition 하나를 계산함
- 이미 계산된 Shuffle Stage는 여러 Job에서 재사용될 수도 있음
4.1 Spark Application
Spark Application은 사용자가 제출한 Spark 프로그램의 전체 실행 단위임.
spark-submit order_analysis.py일반적으로 한 번의 spark-submit으로 다음과 같은 하나의 Application이 실행됨.
Spark Application
├─ Driver
├─ Executor 1
├─ Executor 2
└─ Executor 3Application 안에서는 여러 개의 Action을 호출할 수 있음.
orders = spark.read.parquet("s3://orders/")
paid = orders.filter("status = 'PAID'")
paid.count()
paid.write.parquet("s3://results/paid-orders/")개념적으로:
Spark Application
├─ Job 1: paid.count()
└─ Job 2: paid.write.parquet(...)처럼 여러 Job이 만들어질 수 있음.
Application과 Job의 차이
Application
= Spark 프로그램 전체의 생명주기
Job
= Application 안에서 특정 결과를 계산하기 위한 실행예를 들어 Airflow가 Spark 작업을 실행한다면:
Airflow Task
↓ spark-submit
Spark Application
├─ Job 1
├─ Job 2
└─ Job 3가 될 수 있음.
Airflow Task 하나가 Spark Job 하나와 반드시 일치하는 것은 아님.
Airflow Task
→ Spark Application 제출
Spark Application
→ 내부적으로 여러 Spark Job 실행 가능4.2 From Actions to Jobs
RDD나 DataFrame의 Transformation은 바로 실행되지 않음.
result = (
orders
.filter(lambda x: x.status == "PAID")
.map(lambda x: (x.region, x.amount))
.reduceByKey(lambda a, b: a + b)
)여기까지는:
RDD Lineage 생성
실제 Job은 아직 없음Action이 호출되면 실제 결과가 필요해짐.
result.collect()Action
↓
Job Submission
↓
DAGSchedulerRDD API의 대표적인 Action:
count()collect()take()reduce()first()saveAsTextFile()
DataFrame API의 대표적인 Action:
show()count()collect()writeforeach()
Action 하나가 항상 Job 하나인가?
기본적인 RDD 예제에서는 다음처럼 이해해도 괜찮음.
Action 하나
≈ Job 하나하지만 현대 Spark SQL에서는 완전히 엄밀한 1:1 관계라고 보기는 어려움.
예를 들어:
- Broadcast Join을 위한 작은 테이블 수집
- Subquery 실행
- Adaptive Query Execution
- 파일 메타데이터 처리
등으로 하나의 사용자 Action을 처리하는 과정에서 추가 Spark Job이 관찰될 수 있음.
4.3 From RDD Lineage to an Execution DAG
다음 예제를 사용해 전체 과정을 살펴보자.
orders = sc.textFile("orders.csv")
result = (
orders
.map(parse_order)
.filter(lambda x: x.status == "PAID")
.map(lambda x: (x.region, x.amount))
.reduceByKey(lambda a, b: a + b)
)
result.collect()Transformation이 구성하는 RDD Lineage:
TextFileRDD
↓ map(parse_order)
MapPartitionsRDD
↓ filter(status == PAID)
MapPartitionsRDD
↓ map(region, amount)
MapPartitionsRDD
↓ reduceByKey
ShuffledRDDAction이 호출되면 Driver의 DAGScheduler가 최종 RDD부터 부모 방향으로 Dependency를 추적함.
collect(result)
↓
result의 부모는 무엇인가?
↓
그 부모의 부모는 무엇인가?
↓
원본 데이터까지 추적이 과정에서 DAGScheduler는 다음을 확인함.
각 RDD의 Partition은 몇 개인가?
어떤 부모 RDD가 필요한가?
어떤 Dependency가 Narrow인가?
어떤 Dependency가 Shuffle Dependency인가?
이미 캐시된 Partition이 있는가?
이미 계산된 Shuffle Output이 있는가?RDD Lineage와 실행 DAG는 관련 있지만 같은 것은 아님.
RDD Lineage
TextFileRDD
↓ map
RDD A
↓ filter
RDD B
↓ map
RDD C
↓ reduceByKey
RDD DExecution DAG
Stage 0
↓ Shuffle
Stage 1RDD DAG
- 데이터가 어떤 Transformation을 거쳐 생성되는지 표현
- 노드는 RDD
- 간선은 Dependency
- Lazy Evaluation과 장애 복구의 기반
Stage DAG
- 실제 실행할 Stage들의 순서 표현
- Shuffle을 기준으로 나뉨
- DAGScheduler가 RDD Dependency를 분석해 생성
4.4 Narrow and Wide Dependencies
Stage가 어떻게 나뉘는지 이해하려면 RDD Dependency를 구분해야 함.
Narrow Dependency
자식 RDD의 Partition 하나가 부모 RDD의 일부 Partition에만 의존하는 관계임.
Parent Partition 0 → Child Partition 0
Parent Partition 1 → Child Partition 1
Parent Partition 2 → Child Partition 2대표적인 연산:
mapfilterflatMapmapPartitionsunion
예시:
rdd2 = rdd1.map(f)
rdd3 = rdd2.filter(g)rdd1 Partition 0
↓ map
rdd2 Partition 0
↓ filter
rdd3 Partition 0부모 Partition 0만 있으면 자식 Partition 0을 계산할 수 있음.
따라서 하나의 Task 안에서 연산을 연결할 수 있음.
Task for Partition 0
Read Partition 0
→ map
→ filter
→ map
→ outputRDD별 중간 결과 전체를 별도로 저장하지 않고 Iterator를 따라 한 레코드씩 처리함.
Record 1 → map → filter → map
Record 2 → map → filter → map
Record 3 → map → filter → map이것이 Pipelining임.
Wide Dependency
자식 RDD의 Partition 하나를 계산하려면 여러 부모 Partition의 데이터가 필요한 관계임.
Parent P0 ─┬→ Child P0
Parent P1 ─┼→ Child P0
Parent P2 ─┘대표적인 연산:
reduceByKeygroupByKeyjoindistinctrepartition- 정렬 연산
예를 들어 지역별 합계를 계산한다고 하자.
Parent Partition 0
Seoul, Busan, Seoul
Parent Partition 1
Busan, Daejeon, Seoul
Parent Partition 2
Daejeon, Busan같은 Key를 한곳에 모아야 함.
Seoul → Child Partition 0
Busan → Child Partition 1
Daejeon → Child Partition 2이를 위해 여러 부모 Partition의 데이터가 네트워크를 통해 재분배됨.
Parent Partitions
↓
Shuffle
↓
Child PartitionsDependency 비교
| 구분 | Narrow Dependency | Wide Dependency |
|---|---|---|
| 부모 관계 | 일부 부모 Partition | 여러 부모 Partition |
| 네트워크 이동 | 일반적으로 없음 | Shuffle 발생 |
| Stage | 같은 Stage 가능 | Stage 경계 생성 |
| 실패 복구 | 해당 부모 Partition 중심 | 여러 Shuffle Output 필요 |
| 대표 연산 | map, filter | groupByKey, join |
4.5 Shuffle and Stage Boundaries
Shuffle은 Stage를 나누는 가장 중요한 경계임.
Stage 0
↓ Shuffle
Stage 1왜 Shuffle 앞뒤를 같은 Stage에서 바로 실행할 수 없을까?
다음과 같은 reduceByKey가 있다고 하자.
pairs.reduceByKey(lambda a, b: a + b)지역별 합계를 계산하려면 모든 Map Task가 각 Key에 해당하는 데이터를 내보내야 함.
Map Task 0 ─┬→ Seoul Partition
├→ Busan Partition
└→ Daejeon Partition
Map Task 1 ─┬→ Seoul Partition
├→ Busan Partition
└→ Daejeon Partition그 뒤 Reduce 쪽 Task가 여러 Map Task의 결과를 가져감.
Seoul Task
├─ Map Task 0의 Seoul Block
├─ Map Task 1의 Seoul Block
└─ Map Task 2의 Seoul Block따라서 일반적으로 앞 Stage의 Shuffle Output이 준비돼야 뒤 Stage가 이를 안정적으로 읽을 수 있음.
Stage 0 완료
→ Shuffle Output 준비
→ Stage 1 실행Shuffle Write
Shuffle 이전 Stage의 Task가 수행함.
Input Partition
↓
Record 처리
↓
Key별 Target Partition 결정
↓
Sort / Aggregate
↓
Shuffle File Write예시:
Map Task 0의 입력
(Seoul, 100)
(Busan, 200)
(Seoul, 300)Shuffle Write 결과
Reduce Partition 0용 Block
→ Seoul 데이터
Reduce Partition 1용 Block
→ Busan 데이터Shuffle Read
다음 Stage의 Task가 수행함.
Reduce Task 0
├─ Map Task 0의 Block 0
├─ Map Task 1의 Block 0
└─ Map Task 2의 Block 0
↓
Merge / Sort
↓
Aggregate각 Task는 필요한 Shuffle Block의 위치를 확인하고, 다른 Executor에 있으면 네트워크로 가져옴.
Executor A ─┐
Executor B ─┼→ Executor C의 Reduce Task
Executor D ─┘Shuffle 비용
Shuffle
├─ Serialization
├─ Local Disk Write
├─ Network Transfer
├─ Remote Disk Read
├─ Sorting
├─ Aggregation
└─ Memory Spill따라서 Spark 성능 문제는 다음과 연결되는 경우가 많음.
- Shuffle 데이터 양이 지나치게 큼
- Partition 수가 부적절함
- 특정 Key에 데이터가 몰림
- Shuffle Fetch가 느림
- 메모리 부족으로 디스크 Spill 발생
4.6 Jobs and Stages
Job은 Action으로 시작된 전체 계산이고, Stage는 Job 안에서 Shuffle 없이 연속 실행 가능한 Task 집합임.
Job
├─ Stage 0
├─ Stage 1
└─ Stage 2Spark Core에는 대표적으로 두 종류의 Stage가 있음.
ShuffleMapStage
Shuffle 데이터를 생성하는 Stage임.
Input RDD
↓
map / filter / map
↓
ShuffleMapStage
↓
Shuffle Output Files이 Stage의 최종 결과는 Driver에 직접 반환되는 값이 아니라, 다음 Stage가 읽을 Shuffle Output임.
ShuffleMapStage
→ 다른 Stage가 읽을 중간 데이터 생성ResultStage
Job의 최종 Action을 수행하는 Stage임.
Final RDD
↓
ResultStage
↓
Driver Result 또는 External Storage예시:
rdd.collect()ResultStage
→ 각 Partition의 결과를 Driver에 반환rdd.saveAsTextFile(path)ResultStage
→ 각 Partition의 결과를 외부 저장소에 기록Spark의 DAGScheduler는 최종 Action을 수행하는 ResultStage와 Shuffle Map Output을 기록하는 ShuffleMapStage를 구분함. (GitHub)
예시: Shuffle 한 번
result = (
orders
.map(parse_order)
.filter(lambda x: x.status == "PAID")
.map(lambda x: (x.region, x.amount))
.reduceByKey(lambda a, b: a + b)
)
result.collect()Job 0
Stage 0: ShuffleMapStage
├─ Read
├─ map(parse_order)
├─ filter
├─ map(region, amount)
└─ Shuffle Write
↓ Shuffle
Stage 1: ResultStage
├─ Shuffle Read
├─ reduceByKey
└─ collect 결과 반환예시: Shuffle 두 번
result = (
orders
.map(lambda x: (x.user_id, x.amount))
.reduceByKey(lambda a, b: a + b)
.map(lambda x: (bucket(x[1]), 1))
.reduceByKey(lambda a, b: a + b)
)
result.collect()Stage 0
→ 첫 번째 Shuffle Write
Stage 1
→ 첫 번째 Shuffle Read
→ 두 번째 Shuffle Write
Stage 2
→ 두 번째 Shuffle Read
→ 최종 결과Job
├─ Stage 0
│ ↓ Shuffle
├─ Stage 1
│ ↓ Shuffle
└─ Stage 24.7 Partitions and Tasks
Stage가 만들어지면 각 Partition을 계산하기 위한 Task가 생성됨.
Stage
├─ Partition 0 → Task 0
├─ Partition 1 → Task 1
├─ Partition 2 → Task 2
└─ Partition 3 → Task 3일반적으로 한 Stage 안에서:
Partition 하나당 Task 하나가 생성된다.
Partition과 Task의 차이
Partition
= 계산해야 할 데이터 조각
Task
= 해당 데이터 조각을 계산하라는 실행 명령Partition 17
→ 논리적인 데이터 조각
Task Attempt 17.0
→ Partition 17을 계산하기 위한 첫 번째 실행 시도Task가 실패해 재시도되면:
Partition 17
├─ Task Attempt 17.0 → 실패
└─ Task Attempt 17.1 → 성공처럼 같은 Partition에 대해 여러 Task Attempt가 생길 수 있음.
ShuffleMapTask
ShuffleMapStage 안에서 실행되는 Task임.
ShuffleMapTask
→ 부모 RDD의 Partition 하나 계산
→ 결과를 Shuffle Partition별로 분리
→ Shuffle Output 기록Input Partition 0
↓
ShuffleMapTask 0
↓
├─ Reduce Partition 0용 Block
├─ Reduce Partition 1용 Block
└─ Reduce Partition 2용 BlockResultTask
ResultStage 안에서 실행되는 Task임.
ResultTask
→ 최종 RDD Partition 계산
→ Action 함수 실행
→ 결과 반환 또는 저장예시:
collect()
→ Partition 결과를 Driver로 반환
save()
→ Partition 결과를 Storage에 기록Partition 수와 동시 Task 수
Stage에 Partition이 100개라고 하자.
Stage
→ 100 Partitions
→ 100 Tasks전체 Executor Core 수가 20개라면 일반적으로 한 번에 최대 약 20개의 Task가 실행됨.
Wave 1: Task 0–19
Wave 2: Task 20–39
Wave 3: Task 40–59
Wave 4: Task 60–79
Wave 5: Task 80–99이를 Task Wave라고 부르기도 함.
Task 개수
= 총 처리해야 할 병렬 작업 수
Executor Core 수
= 동시에 처리할 수 있는 작업 수Task의 동시 실행 가능 수는 Executor의 사용 가능한 Core와 Task별 자원 요구량에 따라 정해짐.
4.8 DAGScheduler and TaskScheduler
Spark Driver 안에는 실행을 서로 다른 수준에서 담당하는 Scheduler들이 있음.
Driver
DAGScheduler
↓
TaskScheduler
↓
SchedulerBackend
↓
ExecutorsDAGScheduler
DAGScheduler는 무엇을 어떤 순서로 실행해야 하는가를 결정함.
DAGScheduler
├─ Action을 Job으로 등록
├─ RDD Dependency 분석
├─ Shuffle 경계 탐색
├─ Stage 생성
├─ Stage Dependency 관리
├─ Missing Parent Stage 확인
├─ TaskSet 생성
└─ Shuffle Output 손실 시 Stage 재실행예를 들어:
Stage 0 → Stage 1 → Stage 2일 때 Stage 2를 바로 실행하지 않고 먼저 필요한 부모 Stage를 확인함.
Stage 2 요청
↓
Stage 1 결과 없음
↓
Stage 1 요청
↓
Stage 0 결과 없음
↓
Stage 0부터 실행TaskScheduler
TaskScheduler는 Stage 안의 Task를 어느 Executor에서 실행할 것인가를 결정함.
TaskScheduler
├─ TaskSet 등록
├─ 사용 가능한 Executor 확인
├─ Executor 자원과 Task 매칭
├─ Data Locality 고려
├─ Task 실행 요청
├─ Task 상태 추적
└─ 일반적인 Task 실패 재시도DAGScheduler가 Stage를 Task 집합으로 만들면:
DAGScheduler
→ TaskSet(Stage 0의 모든 Task)TaskScheduler가 이를 받음.
TaskSet
├─ Task 0
├─ Task 1
├─ Task 2
└─ Task 3SchedulerBackend
TaskScheduler와 실제 Cluster Manager·Executor 환경 사이를 연결함.
TaskScheduler
↓
SchedulerBackend
↓
Standalone / YARN / Kubernetes
↓
Executor다만 Cluster Manager가 각 Spark Task를 직접 선택하는 것은 아님.
Cluster Manager
→ Executor 자원 제공 및 프로세스 배치
Spark TaskScheduler
→ 사용 가능한 Executor에 Task 할당4.9 Task Scheduling and Executor Assignment
TaskScheduler는 단순히 빈 Executor에 아무 Task나 보내지 않음.
가능하면 데이터에 가까운 Executor를 선택함.
Partition 데이터가 Node A에 있음
↓
Node A의 Executor에서 Task 실행 선호이를 Data Locality라고 함.
대표적인 Locality 단계
PROCESS_LOCAL
→ 같은 Executor Process에 데이터가 있음
NODE_LOCAL
→ 같은 Worker Node에 데이터가 있음
RACK_LOCAL
→ 같은 Rack에 데이터가 있음
ANY
→ 어느 Executor에서든 실행Cache된 RDD Partition이 특정 Executor에 있다면:
Executor A
└─ Cached Partition 3Partition 3을 처리하는 Task를 Executor A에 보내는 것이 유리함.
Task for Partition 3
→ Executor A 선호하지만 항상 가장 가까운 Executor만 기다리면 자원이 놀 수 있음.
Executor A
→ 데이터가 가까움
→ 현재 Busy
Executor B
→ 데이터가 멀리 있음
→ 현재 IdleSpark는 일정 조건에서 Locality보다 병렬 실행을 우선해 다른 Executor에서 Task를 실행할 수 있음.
Data Locality
vs.
Resource Utilization사이의 균형을 잡는 것임.
Task Closure 전달
다음 코드가 있다고 하자.
threshold = 100
result = rdd.filter(lambda x: x.amount > threshold)Executor가 Task를 실행하려면 다음 정보가 필요함.
Task
├─ 실행할 Partition ID
├─ RDD 계산 정보
├─ 사용자 함수
├─ 캡처된 외부 변수
└─ 기타 실행 메타데이터이 정보는 직렬화되어 Executor로 전달됨.
Driver
↓ Serialized Task
Executor
↓ Deserialize
Run Task그래서 함수가 불필요하게 큰 객체를 캡처하면 Task 전달 비용이 커질 수 있음.
large_lookup = load_large_lookup()
rdd.map(lambda x: process(x, large_lookup))이런 경우 Broadcast Variable을 사용하는 것이 적합할 수 있음.
lookup_broadcast = sc.broadcast(large_lookup)
rdd.map(
lambda x: process(x, lookup_broadcast.value)
)4.10 Task Execution Inside an Executor
Task가 Executor에 도착하면 대략 다음 순서로 실행됨.
Serialized Task 수신
↓
Task Deserialize
↓
TaskContext 생성
↓
RDD Partition 계산
↓
Iterator Pipeline 실행
↓
Shuffle / Cache / Output
↓
Task Result 보고예시:
result = (
orders
.map(parse_order)
.filter(is_paid)
.map(to_region_amount)
)Executor에서 Partition 0의 Task를 실행하면:
Task 0
↓
orders.iterator(Partition 0)
↓
read record
↓ parse_order
↓ is_paid
↓ to_region_amount
↓ outputRDD의 compute()가 부모 RDD의 iterator()를 호출하면서 연산이 연결됨.
Final RDD.compute(P0)
↓
Parent RDD.iterator(P0)
↓
Parent RDD.compute(P0)
↓
Source RDD.iterator(P0)
↓
Read Input4.11 Failure and Retry
분산 환경에서는 Task와 Executor 실패를 정상적인 가능성으로 다룸.
일반적인 Task 실패
Task Attempt 5.0
→ 일시적 오류
→ 실패
Task Attempt 5.1
→ 다른 Executor에서 재시도
→ 성공일반적인 Stage 내부 Task 실패는 TaskScheduler가 일정 횟수 재시도함.
TaskScheduler
→ Task 재시도
→ 계속 실패하면 Stage 실패 보고Shuffle Fetch 실패
뒤 Stage의 Task가 앞 Stage의 Shuffle Block을 가져오지 못할 수 있음.
Stage 0
→ Executor A에 Shuffle Output 저장
Executor A 장애
→ Shuffle Output 손실
Stage 1
→ Fetch Failed이 경우 Stage 1의 Task만 다시 실행해도 필요한 데이터가 없음.
Stage 0 재실행
→ 손실된 Shuffle Output 재생성
→ Stage 1 재실행DAGScheduler는 Fetch Failure나 Executor Loss로 Shuffle Output이 손실되면 해당 Shuffle Map Stage를 다시 제출할 수 있음. 반면 Stage 내부의 일반 Task 실패 재시도는 하위 TaskScheduler가 담당함. (GitHub)
Stage Attempt
같은 Stage가 장애로 다시 실행되면 새로운 Stage Attempt가 생김.
Stage 0.0
→ 첫 번째 실행
→ 실패
Stage 0.1
→ 두 번째 실행
→ 성공Spark UI에서 Stage ID 0, Attempt 1처럼 확인할 수 있음.
4.12 Complete Execution Example
다음 코드를 처음부터 끝까지 추적해 보자.
orders = sc.textFile("orders.csv")
result = (
orders
.map(parse_order)
.filter(lambda x: x.status == "PAID")
.map(lambda x: (x.region, x.amount))
.reduceByKey(lambda a, b: a + b)
)
result.collect()① Application 시작
spark-submit
↓
Driver 시작
↓
SparkContext 생성
↓
Executor 확보② RDD Lineage 생성
TextFileRDD
↓ map
MapPartitionsRDD
↓ filter
MapPartitionsRDD
↓ map
MapPartitionsRDD
↓ reduceByKey
ShuffledRDD아직 실제 데이터 계산은 시작되지 않음.
③ Action 호출
result.collect()collect()
→ Job 생성④ DAGScheduler가 Dependency 분석
map → Narrow
filter → Narrow
map → Narrow
reduceByKey → Shuffle Dependency⑤ Stage 분리
Stage 0: ShuffleMapStage
Read orders.csv
→ parse_order
→ filter PAID
→ map(region, amount)
→ Shuffle WriteStage 1: ResultStage
Shuffle Read
→ region별 reduce
→ collect 결과 반환⑥ Task 생성
원본 RDD가 4개 Partition이고 Shuffle 이후 Partition이 2개라면:
Stage 0
├─ Task 0 → Input Partition 0
├─ Task 1 → Input Partition 1
├─ Task 2 → Input Partition 2
└─ Task 3 → Input Partition 3Stage 1
├─ Task 0 → Shuffle Partition 0
└─ Task 1 → Shuffle Partition 1⑦ Executor 할당
Executor A
├─ Stage 0 Task 0
└─ Stage 0 Task 2
Executor B
├─ Stage 0 Task 1
└─ Stage 0 Task 3⑧ Shuffle 실행
Stage 0 Tasks
→ Key별 Shuffle Block 작성
Stage 1 Tasks
→ 필요한 Block Fetch
→ Key별 집계⑨ 결과 반환
Stage 1 Task 0 ─┐
├→ Driver
Stage 1 Task 1 ─┘[
("Seoul", 150000),
("Busan", 80000)
]4.13 The Complete Spark Execution Model
flowchart TD app["Spark Application"] --> transform["Transformation<br/>RDD / DataFrame Plan 생성"] transform --> action["Action"] action --> job["Job"] job --> dagScheduler["DAGScheduler<br/>RDD Dependency와 Shuffle 분석"] dagScheduler --> stageDag["Stage DAG"] stageDag --> shuffleStage["ShuffleMapStage"] stageDag --> resultStage["ResultStage"] shuffleStage --> shuffleTaskSet["TaskSet"] resultStage --> resultTaskSet["TaskSet"] shuffleTaskSet --> taskScheduler["TaskScheduler<br/>사용 가능한 Executor와 Task 매칭"] resultTaskSet --> taskScheduler taskScheduler --> executors["Executors"] executors --> task0["Task 0"] executors --> task1["Task 1"] executors --> task2["Task 2"] task0 --> partition0["Partition 0"] task1 --> partition1["Partition 1"] task2 --> partition2["Partition 2"] partition0 --> output["Shuffle / Cache / Output"] partition1 --> output partition2 --> output output --> complete["Job Complete"]