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"]

핵심 흐름

text
Application
└─ Job
   └─ Stage
      └─ Task
         └─ Partition 처리

다만 이 계층은 단순히 크기순으로 포함되는 관계만은 아님.

  • 하나의 Application에는 여러 Job이 존재할 수 있음
  • 하나의 Job은 여러 Stage를 필요로 할 수 있음
  • 하나의 Stage에는 여러 Task가 존재함
  • 각 Task는 일반적으로 해당 Stage의 Partition 하나를 계산함
  • 이미 계산된 Shuffle Stage는 여러 Job에서 재사용될 수도 있음

4.1 Spark Application

Spark Application은 사용자가 제출한 Spark 프로그램의 전체 실행 단위임.

bash
spark-submit order_analysis.py

일반적으로 한 번의 spark-submit으로 다음과 같은 하나의 Application이 실행됨.

text
Spark Application
├─ Driver
├─ Executor 1
├─ Executor 2
└─ Executor 3

Application 안에서는 여러 개의 Action을 호출할 수 있음.

python
orders = spark.read.parquet("s3://orders/")

paid = orders.filter("status = 'PAID'")

paid.count()
paid.write.parquet("s3://results/paid-orders/")

개념적으로:

text
Spark Application

├─ Job 1: paid.count()
└─ Job 2: paid.write.parquet(...)

처럼 여러 Job이 만들어질 수 있음.

Application과 Job의 차이

text
Application
= Spark 프로그램 전체의 생명주기

Job
= Application 안에서 특정 결과를 계산하기 위한 실행

예를 들어 Airflow가 Spark 작업을 실행한다면:

text
Airflow Task
   ↓ spark-submit
Spark Application
   ├─ Job 1
   ├─ Job 2
   └─ Job 3

가 될 수 있음.

Airflow Task 하나가 Spark Job 하나와 반드시 일치하는 것은 아님.

text
Airflow Task
→ Spark Application 제출

Spark Application
→ 내부적으로 여러 Spark Job 실행 가능

4.2 From Actions to Jobs

RDD나 DataFrame의 Transformation은 바로 실행되지 않음.

python
result = (
    orders
    .filter(lambda x: x.status == "PAID")
    .map(lambda x: (x.region, x.amount))
    .reduceByKey(lambda a, b: a + b)
)

여기까지는:

text
RDD Lineage 생성
실제 Job은 아직 없음

Action이 호출되면 실제 결과가 필요해짐.

python
result.collect()
text
Action
   ↓
Job Submission
   ↓
DAGScheduler

RDD API의 대표적인 Action:

  • count()
  • collect()
  • take()
  • reduce()
  • first()
  • saveAsTextFile()

DataFrame API의 대표적인 Action:

  • show()
  • count()
  • collect()
  • write
  • foreach()

Action 하나가 항상 Job 하나인가?

기본적인 RDD 예제에서는 다음처럼 이해해도 괜찮음.

text
Action 하나
≈ Job 하나

하지만 현대 Spark SQL에서는 완전히 엄밀한 1:1 관계라고 보기는 어려움.

예를 들어:

  • Broadcast Join을 위한 작은 테이블 수집
  • Subquery 실행
  • Adaptive Query Execution
  • 파일 메타데이터 처리

등으로 하나의 사용자 Action을 처리하는 과정에서 추가 Spark Job이 관찰될 수 있음.


4.3 From RDD Lineage to an Execution DAG

다음 예제를 사용해 전체 과정을 살펴보자.

python
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:

text
TextFileRDD
     ↓ map(parse_order)
MapPartitionsRDD
     ↓ filter(status == PAID)
MapPartitionsRDD
     ↓ map(region, amount)
MapPartitionsRDD
     ↓ reduceByKey
ShuffledRDD

Action이 호출되면 Driver의 DAGScheduler가 최종 RDD부터 부모 방향으로 Dependency를 추적함.

text
collect(result)
       ↓
result의 부모는 무엇인가?
       ↓
그 부모의 부모는 무엇인가?
       ↓
원본 데이터까지 추적

이 과정에서 DAGScheduler는 다음을 확인함.

text
각 RDD의 Partition은 몇 개인가?
어떤 부모 RDD가 필요한가?
어떤 Dependency가 Narrow인가?
어떤 Dependency가 Shuffle Dependency인가?
이미 캐시된 Partition이 있는가?
이미 계산된 Shuffle Output이 있는가?

RDD Lineage와 실행 DAG는 관련 있지만 같은 것은 아님.

text
RDD Lineage

TextFileRDD
   ↓ map
RDD A
   ↓ filter
RDD B
   ↓ map
RDD C
   ↓ reduceByKey
RDD D
text
Execution DAG

Stage 0
   ↓ Shuffle
Stage 1

RDD 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에만 의존하는 관계임.

text
Parent Partition 0 → Child Partition 0
Parent Partition 1 → Child Partition 1
Parent Partition 2 → Child Partition 2

대표적인 연산:

  • map
  • filter
  • flatMap
  • mapPartitions
  • union

예시:

python
rdd2 = rdd1.map(f)
rdd3 = rdd2.filter(g)
text
rdd1 Partition 0
        ↓ map
rdd2 Partition 0
        ↓ filter
rdd3 Partition 0

부모 Partition 0만 있으면 자식 Partition 0을 계산할 수 있음.

따라서 하나의 Task 안에서 연산을 연결할 수 있음.

text
Task for Partition 0

Read Partition 0
→ map
→ filter
→ map
→ output

RDD별 중간 결과 전체를 별도로 저장하지 않고 Iterator를 따라 한 레코드씩 처리함.

text
Record 1 → map → filter → map
Record 2 → map → filter → map
Record 3 → map → filter → map

이것이 Pipelining임.


Wide Dependency

자식 RDD의 Partition 하나를 계산하려면 여러 부모 Partition의 데이터가 필요한 관계임.

text
Parent P0 ─┬→ Child P0
Parent P1 ─┼→ Child P0
Parent P2 ─┘

대표적인 연산:

  • reduceByKey
  • groupByKey
  • join
  • distinct
  • repartition
  • 정렬 연산

예를 들어 지역별 합계를 계산한다고 하자.

text
Parent Partition 0
Seoul, Busan, Seoul

Parent Partition 1
Busan, Daejeon, Seoul

Parent Partition 2
Daejeon, Busan

같은 Key를 한곳에 모아야 함.

text
Seoul   → Child Partition 0
Busan   → Child Partition 1
Daejeon → Child Partition 2

이를 위해 여러 부모 Partition의 데이터가 네트워크를 통해 재분배됨.

text
Parent Partitions
       ↓
     Shuffle
       ↓
Child Partitions

Dependency 비교

구분Narrow DependencyWide Dependency
부모 관계일부 부모 Partition여러 부모 Partition
네트워크 이동일반적으로 없음Shuffle 발생
Stage같은 Stage 가능Stage 경계 생성
실패 복구해당 부모 Partition 중심여러 Shuffle Output 필요
대표 연산map, filtergroupByKey, join

4.5 Shuffle and Stage Boundaries

Shuffle은 Stage를 나누는 가장 중요한 경계임.

text
Stage 0
   ↓ Shuffle
Stage 1

왜 Shuffle 앞뒤를 같은 Stage에서 바로 실행할 수 없을까?

다음과 같은 reduceByKey가 있다고 하자.

python
pairs.reduceByKey(lambda a, b: a + b)

지역별 합계를 계산하려면 모든 Map Task가 각 Key에 해당하는 데이터를 내보내야 함.

text
Map Task 0 ─┬→ Seoul Partition
            ├→ Busan Partition
            └→ Daejeon Partition

Map Task 1 ─┬→ Seoul Partition
            ├→ Busan Partition
            └→ Daejeon Partition

그 뒤 Reduce 쪽 Task가 여러 Map Task의 결과를 가져감.

text
Seoul Task
├─ Map Task 0의 Seoul Block
├─ Map Task 1의 Seoul Block
└─ Map Task 2의 Seoul Block

따라서 일반적으로 앞 Stage의 Shuffle Output이 준비돼야 뒤 Stage가 이를 안정적으로 읽을 수 있음.

text
Stage 0 완료
→ Shuffle Output 준비
→ Stage 1 실행

Shuffle Write

Shuffle 이전 Stage의 Task가 수행함.

text
Input Partition
      ↓
Record 처리
      ↓
Key별 Target Partition 결정
      ↓
Sort / Aggregate
      ↓
Shuffle File Write

예시:

text
Map Task 0의 입력

(Seoul, 100)
(Busan, 200)
(Seoul, 300)
text
Shuffle Write 결과

Reduce Partition 0용 Block
→ Seoul 데이터

Reduce Partition 1용 Block
→ Busan 데이터

Shuffle Read

다음 Stage의 Task가 수행함.

text
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에 있으면 네트워크로 가져옴.

text
Executor A ─┐
Executor B ─┼→ Executor C의 Reduce Task
Executor D ─┘

Shuffle 비용

text
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 집합임.

text
Job
├─ Stage 0
├─ Stage 1
└─ Stage 2

Spark Core에는 대표적으로 두 종류의 Stage가 있음.

ShuffleMapStage

Shuffle 데이터를 생성하는 Stage임.

text
Input RDD
   ↓
map / filter / map
   ↓
ShuffleMapStage
   ↓
Shuffle Output Files

이 Stage의 최종 결과는 Driver에 직접 반환되는 값이 아니라, 다음 Stage가 읽을 Shuffle Output임.

text
ShuffleMapStage
→ 다른 Stage가 읽을 중간 데이터 생성

ResultStage

Job의 최종 Action을 수행하는 Stage임.

text
Final RDD
   ↓
ResultStage
   ↓
Driver Result 또는 External Storage

예시:

python
rdd.collect()
text
ResultStage
→ 각 Partition의 결과를 Driver에 반환
python
rdd.saveAsTextFile(path)
text
ResultStage
→ 각 Partition의 결과를 외부 저장소에 기록

Spark의 DAGScheduler는 최종 Action을 수행하는 ResultStage와 Shuffle Map Output을 기록하는 ShuffleMapStage를 구분함. (GitHub)


예시: Shuffle 한 번

python
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()
text
Job 0

Stage 0: ShuffleMapStage
├─ Read
├─ map(parse_order)
├─ filter
├─ map(region, amount)
└─ Shuffle Write

             ↓ Shuffle

Stage 1: ResultStage
├─ Shuffle Read
├─ reduceByKey
└─ collect 결과 반환

예시: Shuffle 두 번

python
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()
text
Stage 0
→ 첫 번째 Shuffle Write

Stage 1
→ 첫 번째 Shuffle Read
→ 두 번째 Shuffle Write

Stage 2
→ 두 번째 Shuffle Read
→ 최종 결과
text
Job
├─ Stage 0
│    ↓ Shuffle
├─ Stage 1
│    ↓ Shuffle
└─ Stage 2

4.7 Partitions and Tasks

Stage가 만들어지면 각 Partition을 계산하기 위한 Task가 생성됨.

text
Stage
├─ Partition 0 → Task 0
├─ Partition 1 → Task 1
├─ Partition 2 → Task 2
└─ Partition 3 → Task 3

일반적으로 한 Stage 안에서:

Partition 하나당 Task 하나가 생성된다.

Partition과 Task의 차이

text
Partition
= 계산해야 할 데이터 조각

Task
= 해당 데이터 조각을 계산하라는 실행 명령
text
Partition 17
→ 논리적인 데이터 조각

Task Attempt 17.0
→ Partition 17을 계산하기 위한 첫 번째 실행 시도

Task가 실패해 재시도되면:

text
Partition 17
├─ Task Attempt 17.0 → 실패
└─ Task Attempt 17.1 → 성공

처럼 같은 Partition에 대해 여러 Task Attempt가 생길 수 있음.


ShuffleMapTask

ShuffleMapStage 안에서 실행되는 Task임.

text
ShuffleMapTask
→ 부모 RDD의 Partition 하나 계산
→ 결과를 Shuffle Partition별로 분리
→ Shuffle Output 기록
text
Input Partition 0
      ↓
ShuffleMapTask 0
      ↓
├─ Reduce Partition 0용 Block
├─ Reduce Partition 1용 Block
└─ Reduce Partition 2용 Block

ResultTask

ResultStage 안에서 실행되는 Task임.

text
ResultTask
→ 최종 RDD Partition 계산
→ Action 함수 실행
→ 결과 반환 또는 저장

예시:

text
collect()
→ Partition 결과를 Driver로 반환

save()
→ Partition 결과를 Storage에 기록

Partition 수와 동시 Task 수

Stage에 Partition이 100개라고 하자.

text
Stage
→ 100 Partitions
→ 100 Tasks

전체 Executor Core 수가 20개라면 일반적으로 한 번에 최대 약 20개의 Task가 실행됨.

text
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라고 부르기도 함.

text
Task 개수
= 총 처리해야 할 병렬 작업 수

Executor Core 수
= 동시에 처리할 수 있는 작업 수

Task의 동시 실행 가능 수는 Executor의 사용 가능한 Core와 Task별 자원 요구량에 따라 정해짐.


4.8 DAGScheduler and TaskScheduler

Spark Driver 안에는 실행을 서로 다른 수준에서 담당하는 Scheduler들이 있음.

text
Driver

DAGScheduler
      ↓
TaskScheduler
      ↓
SchedulerBackend
      ↓
Executors

DAGScheduler

DAGScheduler는 무엇을 어떤 순서로 실행해야 하는가를 결정함.

text
DAGScheduler
├─ Action을 Job으로 등록
├─ RDD Dependency 분석
├─ Shuffle 경계 탐색
├─ Stage 생성
├─ Stage Dependency 관리
├─ Missing Parent Stage 확인
├─ TaskSet 생성
└─ Shuffle Output 손실 시 Stage 재실행

예를 들어:

text
Stage 0 → Stage 1 → Stage 2

일 때 Stage 2를 바로 실행하지 않고 먼저 필요한 부모 Stage를 확인함.

text
Stage 2 요청
   ↓
Stage 1 결과 없음
   ↓
Stage 1 요청
   ↓
Stage 0 결과 없음
   ↓
Stage 0부터 실행

TaskScheduler

TaskScheduler는 Stage 안의 Task를 어느 Executor에서 실행할 것인가를 결정함.

text
TaskScheduler
├─ TaskSet 등록
├─ 사용 가능한 Executor 확인
├─ Executor 자원과 Task 매칭
├─ Data Locality 고려
├─ Task 실행 요청
├─ Task 상태 추적
└─ 일반적인 Task 실패 재시도

DAGScheduler가 Stage를 Task 집합으로 만들면:

text
DAGScheduler
→ TaskSet(Stage 0의 모든 Task)

TaskScheduler가 이를 받음.

text
TaskSet
├─ Task 0
├─ Task 1
├─ Task 2
└─ Task 3

SchedulerBackend

TaskScheduler와 실제 Cluster Manager·Executor 환경 사이를 연결함.

text
TaskScheduler
      ↓
SchedulerBackend
      ↓
Standalone / YARN / Kubernetes
      ↓
Executor

다만 Cluster Manager가 각 Spark Task를 직접 선택하는 것은 아님.

text
Cluster Manager
→ Executor 자원 제공 및 프로세스 배치

Spark TaskScheduler
→ 사용 가능한 Executor에 Task 할당

4.9 Task Scheduling and Executor Assignment

TaskScheduler는 단순히 빈 Executor에 아무 Task나 보내지 않음.

가능하면 데이터에 가까운 Executor를 선택함.

text
Partition 데이터가 Node A에 있음
          ↓
Node A의 Executor에서 Task 실행 선호

이를 Data Locality라고 함.

대표적인 Locality 단계

text
PROCESS_LOCAL
→ 같은 Executor Process에 데이터가 있음

NODE_LOCAL
→ 같은 Worker Node에 데이터가 있음

RACK_LOCAL
→ 같은 Rack에 데이터가 있음

ANY
→ 어느 Executor에서든 실행

Cache된 RDD Partition이 특정 Executor에 있다면:

text
Executor A
└─ Cached Partition 3

Partition 3을 처리하는 Task를 Executor A에 보내는 것이 유리함.

text
Task for Partition 3
→ Executor A 선호

하지만 항상 가장 가까운 Executor만 기다리면 자원이 놀 수 있음.

text
Executor A
→ 데이터가 가까움
→ 현재 Busy

Executor B
→ 데이터가 멀리 있음
→ 현재 Idle

Spark는 일정 조건에서 Locality보다 병렬 실행을 우선해 다른 Executor에서 Task를 실행할 수 있음.

text
Data Locality
vs.
Resource Utilization

사이의 균형을 잡는 것임.


Task Closure 전달

다음 코드가 있다고 하자.

python
threshold = 100

result = rdd.filter(lambda x: x.amount > threshold)

Executor가 Task를 실행하려면 다음 정보가 필요함.

text
Task
├─ 실행할 Partition ID
├─ RDD 계산 정보
├─ 사용자 함수
├─ 캡처된 외부 변수
└─ 기타 실행 메타데이터

이 정보는 직렬화되어 Executor로 전달됨.

text
Driver
   ↓ Serialized Task
Executor
   ↓ Deserialize
Run Task

그래서 함수가 불필요하게 큰 객체를 캡처하면 Task 전달 비용이 커질 수 있음.

python
large_lookup = load_large_lookup()

rdd.map(lambda x: process(x, large_lookup))

이런 경우 Broadcast Variable을 사용하는 것이 적합할 수 있음.

python
lookup_broadcast = sc.broadcast(large_lookup)

rdd.map(
    lambda x: process(x, lookup_broadcast.value)
)

4.10 Task Execution Inside an Executor

Task가 Executor에 도착하면 대략 다음 순서로 실행됨.

text
Serialized Task 수신
       ↓
Task Deserialize
       ↓
TaskContext 생성
       ↓
RDD Partition 계산
       ↓
Iterator Pipeline 실행
       ↓
Shuffle / Cache / Output
       ↓
Task Result 보고

예시:

python
result = (
    orders
    .map(parse_order)
    .filter(is_paid)
    .map(to_region_amount)
)

Executor에서 Partition 0의 Task를 실행하면:

text
Task 0
  ↓
orders.iterator(Partition 0)
  ↓
read record
  ↓ parse_order
  ↓ is_paid
  ↓ to_region_amount
  ↓ output

RDD의 compute()가 부모 RDD의 iterator()를 호출하면서 연산이 연결됨.

text
Final RDD.compute(P0)
       ↓
Parent RDD.iterator(P0)
       ↓
Parent RDD.compute(P0)
       ↓
Source RDD.iterator(P0)
       ↓
Read Input

4.11 Failure and Retry

분산 환경에서는 Task와 Executor 실패를 정상적인 가능성으로 다룸.

일반적인 Task 실패

text
Task Attempt 5.0
→ 일시적 오류
→ 실패

Task Attempt 5.1
→ 다른 Executor에서 재시도
→ 성공

일반적인 Stage 내부 Task 실패는 TaskScheduler가 일정 횟수 재시도함.

text
TaskScheduler
→ Task 재시도
→ 계속 실패하면 Stage 실패 보고

Shuffle Fetch 실패

뒤 Stage의 Task가 앞 Stage의 Shuffle Block을 가져오지 못할 수 있음.

text
Stage 0
→ Executor A에 Shuffle Output 저장

Executor A 장애
→ Shuffle Output 손실

Stage 1
→ Fetch Failed

이 경우 Stage 1의 Task만 다시 실행해도 필요한 데이터가 없음.

text
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가 생김.

text
Stage 0.0
→ 첫 번째 실행
→ 실패

Stage 0.1
→ 두 번째 실행
→ 성공

Spark UI에서 Stage ID 0, Attempt 1처럼 확인할 수 있음.


4.12 Complete Execution Example

다음 코드를 처음부터 끝까지 추적해 보자.

python
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 시작

text
spark-submit
      ↓
Driver 시작
      ↓
SparkContext 생성
      ↓
Executor 확보

② RDD Lineage 생성

text
TextFileRDD
    ↓ map
MapPartitionsRDD
    ↓ filter
MapPartitionsRDD
    ↓ map
MapPartitionsRDD
    ↓ reduceByKey
ShuffledRDD

아직 실제 데이터 계산은 시작되지 않음.

③ Action 호출

python
result.collect()
text
collect()
→ Job 생성

④ DAGScheduler가 Dependency 분석

text
map       → Narrow
filter    → Narrow
map       → Narrow
reduceByKey → Shuffle Dependency

⑤ Stage 분리

text
Stage 0: ShuffleMapStage

Read orders.csv
→ parse_order
→ filter PAID
→ map(region, amount)
→ Shuffle Write
text
Stage 1: ResultStage

Shuffle Read
→ region별 reduce
→ collect 결과 반환

⑥ Task 생성

원본 RDD가 4개 Partition이고 Shuffle 이후 Partition이 2개라면:

text
Stage 0
├─ Task 0 → Input Partition 0
├─ Task 1 → Input Partition 1
├─ Task 2 → Input Partition 2
└─ Task 3 → Input Partition 3
text
Stage 1
├─ Task 0 → Shuffle Partition 0
└─ Task 1 → Shuffle Partition 1

⑦ Executor 할당

text
Executor A
├─ Stage 0 Task 0
└─ Stage 0 Task 2

Executor B
├─ Stage 0 Task 1
└─ Stage 0 Task 3

⑧ Shuffle 실행

text
Stage 0 Tasks
→ Key별 Shuffle Block 작성

Stage 1 Tasks
→ 필요한 Block Fetch
→ Key별 집계

⑨ 결과 반환

text
Stage 1 Task 0 ─┐
                ├→ Driver
Stage 1 Task 1 ─┘
python
[
    ("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"]

Discussion