1463 words
7 minutes
[Spark] 실행모델 및 DAG
2025-10-21
2026-07-30

Spark 실행 모델 개관과 DAG#

1-1. Cluster#

Cluster : 네트워크로 연결된 여러 대의 머신(node)을 하나의 논리적 컴퓨팅 자원으로 묶은 것.

  • 단일 머신의 한계: CPU 코어 수, RAM 용량, 디스크 I/O 대역폭이 물리적으로 고정됨
  • Scale-up(더 큰 장비)은 비용이 지수적으로 증가하고 상한이 존재
  • Scale-out(장비 추가)이 대안이나, 여러 대에 걸쳐 하나의 계산을 나눠 실행하는 소프트웨어가 필요
  • Spark = 그 소프트웨어. Cluster를 “하나의 거대한 컴퓨터”처럼 보이게 하는 추상화 계층

Cluster에서 Spark가 해결하는 문제 4가지:

  1. 분할(Partitioning) — 데이터를 어떻게 쪼갤 것인가
  2. 배치(Scheduling) — 어느 노드에서 어떤 작업을 실행할 것인가
  3. 이동(Shuffle) — 노드 간 데이터 재배치를 어떻게 할 것인가
  4. 복구(Fault tolerance) — 노드 하나가 죽었을 때 어떻게 되살릴 것인가

DAG#

DAG (Directed Acyclic Graph) : 방향이 있고 순환이 없는 그래프. Spark가 계산을 표현하는 방식

  • Directed(방향성): A → B는 “B는 A의 결과에 의존”을 의미. 데이터 흐름의 방향
  • Acyclic(비순환): A → B → A 같은 되돌아오는 경로가 없음
  • Graph(그래프): 노드(데이터셋 = RDD/DataFrame) + 엣지(연산)

왜 DAG인가#

1. 의존성 표현

분산 실행에서 가장 중요한 정보는 “무엇을 먼저 계산해야 하는가”임. DAG는 이 의존 관계를 직접 표현하는 자료구조.

read(A) read(B)
│ │
filter filter
│ │
└──── join ──┘
groupBy
write
  • join은 양쪽 filter가 끝나야 실행 가능 → 엣지가 그 제약을 표현
  • filter는 서로 의존하지 않음 → 병렬 실행 가능하다는 사실이 그래프 구조에서 자동 도출

2. 비순환성이 종료를 보장

순환이 있으면 “언제 끝나는가”를 정적으로 알 수 없음. 비순환이면 **위상 정렬(topological sort)**이 항상 가능 → 실행 순서를 확정적으로 결정 가능. 스케줄러가 무한 루프에 빠지지 않음.

3. 최적화대상

DAG는 “실행 기록”이 아니라 실행 전에 만들어지는 계획임. 계획 상태이므로 자유롭게 재작성 가능.

  • 연속된 두 filter → 하나로 합치기
  • filterjoin 아래로 밀어내려 데이터량 줄이기 (→ 챕터 5 Filter Pushdown)
  • 불필요한 컬럼 읽기 제거 (→ Projection Pushdown)

즉시 실행(eager) 방식이었다면 이미 실행된 연산을 되돌릴 수 없어 이런 최적화가 불가능. DAG + Lazy Evaluation은 세트로 동작함 (→ 챕터 4).

4. 장애 복구(Lineage)

DAG는 “이 데이터가 어떻게 만들어졌는지”를 담은 완전한 계보(lineage). 노드 하나가 죽어 파티션 일부가 소실되면:

  • 전체 재실행 불필요
  • DAG를 역추적해 소실된 파티션을 만드는 경로만 재계산
  • 데이터 복제(replication) 없이 fault tolerance 확보 → 저장 비용 절감

이것이 Spark와 MapReduce를 가르는 핵심 설계 차이 중 하나.

코드로 보는 DAG#

df = spark.read.parquet("/data/events") # 노드 생성
f1 = df.filter("country = 'KR'") # 엣지 추가
f2 = f1.select("user_id", "event_time") # 엣지 추가
f2.count() # 여기서 비로소 실행
  • 1~3행: DAG만 쌓임. 실제 데이터는 단 1바이트도 읽지 않음
  • 4행 count(): Action. 이 시점에 DAG가 완성되고 실행 계획으로 변환됨

확인 방법:

f2.explain(True) # DAG가 어떤 계획으로 번역됐는지 출력

1-3. 전체 개념 지도#

spark 코드 한 줄이 클러스터에서 실행되기까지의 경로:

[사용자 코드]
│ (Lazy — 실행 안 됨)
[DAG / Unresolved Logical Plan]
│ Catalyst Optimizer
[Logical Plan → Optimized → Physical Plan]
│ Action 호출 시점에 실행 트리거
[Job] ──분할──▶ [Stage] ──분할──▶ [Task]
│ (Shuffle 경계) (Partition 단위)
[DAGScheduler → TaskScheduler]
[Driver] ──자원 요청──▶ [Resource Manager]
│ │
│ ▼
└──Task 전송──────▶ [Executor 여러 개]
[Executor 메모리 영역]
Execution / Storage / User ...

Concept

  • Cluster : 네트워크로 연결된 다수의 노드를 하나의 논리적 컴퓨팅 자원으로 묶은 집합. Spark는 이를 단일 컴퓨터처럼 추상화
  • Scale-out : 장비를 추가해 처리 능력을 늘리는 방식. Scale-up(단일 장비 증설)과 대비
  • DAG (Directed Acyclic Graph) : 방향성 있고 순환 없는 그래프. Spark가 연산 의존 관계를 표현하는 자료구조
  • Directed : 엣지에 방향이 존재. “B는 A에 의존”이라는 실행 순서 제약을 표현
  • Acyclic : 순환 경로 없음. 위상 정렬이 항상 가능해 실행 순서 확정과 종료를 보장
  • 위상 정렬 (Topological Sort) : DAG의 노드를 의존성 순서대로 나열하는 알고리즘. Spark 스케줄러의 실행 순서 결정 기반
  • Lineage : DAG가 담고 있는 데이터 생성 계보. 파티션 소실 시 해당 경로만 재계산해 복구
  • Fault Tolerance : 장애 발생 시 복구 능력. Spark는 데이터 복제 대신 Lineage 재계산으로 확보
  • Partition : 분산 데이터셋의 물리적 분할 단위. Task 실행의 최소 단위이자 병렬도의 상한
  • Eager Evaluation : 연산 호출 즉시 실행하는 방식. 실행 후 되돌릴 수 없어 전역 최적화 불가
  • Lazy Evaluation : 실행을 지연시키고 계획만 축적하는 방식. DAG 재작성 최적화의 전제 조건
  • explain() : DataFrame의 실행 계획을 문자열로 출력하는 메서드. DAG가 어떤 계획으로 번역됐는지 확인하는 도구

[Spark] 실행모델 및 DAG
https://yjinheon.netlify.app/posts/02de/03-spark/spark-01-execution-models/
Author
Datamind
Published at
2025-10-21
License
CC BY-NC-SA 4.0