2. Spark Architecture: Then and Now
RDD 중심 구조에서 DataFrame·Spark Connect까지
- Haram Lee
- 2026-08-16
- studies / Topics / Spark
- 초기의 Driver–Executor 기반 분산 실행 구조는 유지
- 그 위에 DataFrame, SQL, Catalyst Optimizer, Spark Connect 등의 계층이 추가됨
Early Spark
flowchart TD rdd["RDD API"] --> core["Spark Core"] core --> runtime["Driver / Executor"]
Modern Spark
flowchart TD api["DataFrame · SQL · Spark Connect"] --> optimization["Query Optimization"] optimization --> core["Spark Core"] core --> runtime["Driver / Executor"]
변화한 부분
- 사용자가 사용하는 API
- 데이터와 연산을 표현하는 방식
- 실행 계획을 생성하고 최적화하는 과정
- Driver
유지된 부분
- Executor
- Partition
- Task
- Shuffle
- Cluster Manager
2.1 The Early RDD-Centric Architecture
- 초기 Spark에서는 사용자가
SparkContext를 생성하고 RDD API를 직접 사용하는 방식이 중심
from pyspark import SparkContext
sc = SparkContext()
result = (
sc.textFile("hdfs://logs")
.filter(lambda line: "ERROR" in line)
.map(parse_log)
.count()
)SparkContext
- Spark Application과 클러스터를 연결하는 진입점
- RDD 생성
- Broadcast Variable, Accumulator 생성
- 클러스터 매니저에 Executor 자원 요청
- Job 실행을 Spark Core에 전달
User Application
↓
SparkContext
↓
RDD Transformations
↓
RDD Lineage / DAG
↓
Spark Core Scheduler
↓
Driver ───────── Executors- 사용자가
map,filter,reduceByKey등의 Transformation을 호출 - 각 연산은 기존 RDD를 수정하는 것이 아니라 새로운 RDD를 생성
- RDD 사이의 의존 관계가 Lineage와 DAG를 구성
HadoopRDD
↓ map
MapPartitionsRDD
↓ filter
MapPartitionsRDD
↓ reduceByKey
ShuffledRDD- Spark Core는 RDD DAG를 보고 다음을 결정함.
- 각 RDD가 어떤 파티션으로 구성되는가?
- 어떤 부모 RDD를 계산해야 하는가?
- 어느 지점에서 Shuffle이 발생하는가?
- 실패한 파티션을 어떻게 재계산할 것인가?
RDD API의 한계
Spark가 사용자 함수 내부의 의미를 알기 어려움.
rdd.map(lambda row: custom_function(row))Spark가 알 수 있는 것
- 각 레코드에 함수를 적용해야 한다는 사실
Spark가 알기 어려운 것
- 어떤 컬럼을 사용하는지
- 필터를 더 앞에서 실행할 수 있는지
- 조인 순서를 변경할 수 있는지
- 일부 연산을 데이터 소스에 내려보낼 수 있는지
2.2 The Modern DataFrame- and SQL-Centric Architecture
- 현재 Spark에서는 일반적으로
SparkContext보다SparkSession을 진입점으로 사용
from pyspark.sql import SparkSession
spark = (
SparkSession.builder
.appName("order-analysis")
.getOrCreate()
)SparkSession
- DataFrame API의 진입점
- Spark SQL 실행
- 데이터 소스 읽기와 쓰기
- Catalog, Table, View 관리
- 내부적으로 SparkContext를 포함
SparkSession
├─ DataFrame API
├─ Spark SQL
├─ Catalog
└─ SparkContext사용자는 RDD의 각 레코드에 어떤 함수를 실행할지보다, 원하는 결과를 선언적으로 표현함.
result = (
spark.read.parquet("s3://orders/")
.filter("status = 'PAID'")
.groupBy("region")
.sum("amount")
)SELECT region, SUM(amount)
FROM orders
WHERE status = 'PAID'
GROUP BY region;Spark는 DataFrame과 SQL을 바로 실행하지 않고 실행 계획으로 변환함.
SQL / DataFrame API
↓
Unresolved Logical Plan
↓
Analyzed Logical Plan
↓
Optimized Logical Plan
↓
Physical Plan
↓
Spark Core
↓
Job / Stage / Task
↓
ExecutorsDataFrame은 스키마와 연산 구조를 Spark에 제공함.
orders
user_id : LONG
status : STRING
region : STRING
amount : DECIMALSpark가 알 수 있는 정보:
- 어떤 컬럼이 필요한가?
- 어떤 조건으로 필터링하는가?
- 어떤 Key로 조인하는가?
- 어떤 컬럼을 집계하는가?
- 데이터의 예상 크기는 어느 정도인가?
이를 통해 Spark가 실행 방식을 최적화할 수 있음.
사용자가 표현한 흐름
모든 컬럼 읽기
→ status 필터
→ region, amount 선택최적화된 흐름
status, region, amount만 읽기
→ 가능한 경우 파일 읽기 단계에서 필터2.3 Spark Core and High-Level APIs
현재 Spark는 여러 고수준 API를 제공함.
Spark SQL
DataFrame
Dataset
Structured Streaming
MLlib
GraphX이들은 서로 완전히 별개의 실행 엔진이 아니라, 대부분 Spark Core 위에서 동작함.
High-Level APIs
├─ Spark SQL
├─ DataFrame
├─ Structured Streaming
└─ MLlib
↓
Query Planning / Library Logic
↓
Spark Core
↓
Driver / ExecutorSpark Core가 제공하는 핵심 기능
- 분산 데이터 파티셔닝
- Task 스케줄링
- Executor 관리
- Shuffle
- Cache와 Persist
- Broadcast
- 장애 복구
DataFrame과 SQL이 RDD를 완전히 제거한 것은 아님.
DataFrame / SQL
↓
Logical Plan
↓
Physical Plan
↓
RDD[InternalRow]
또는
RDD[ColumnarBatch]
↓
Task Execution정리하면:
RDD
→ 저수준 분산 데이터 API이자 Spark Core 실행 모델
DataFrame / SQL
→ Spark Core 위에서 실행 계획 최적화를 제공하는 고수준 API2.4 Spark Components at a Glance
flowchart TD user["User Code"] --> driver["Driver"] driver -- "① 자원 요청" --> manager["Cluster Manager"] manager -- "② Executor 실행" --> executors["Executors on Worker Nodes"] driver -- "③ Task 할당" --> executors executors --> storage["S3 / GCS / HDFS / Database"]
| 컴포넌트 | 정체 | 핵심 책임 | 개수 | 죽으면? |
|---|---|---|---|---|
| SparkSession / SparkContext | Application 코드의 진입점 객체 | API 제공, 클러스터 연결 | Application당 1개 | Driver와 운명 공동체 |
| Driver | 조율자 프로세스 | 실행 계획 생성, Job→Stage→Task 분해, Task 할당과 추적 | Application당 1개 | Application 전체 실패 |
| Cluster Manager | 클러스터 자원 관리자 | Executor를 어디에 얼마나 띄울지 결정 | 클러스터당 1개 | 신규 자원 할당 불가 |
| Executor | 일꾼 프로세스 | Task 실행, Shuffle Write·Read, Cache 보관 | Application당 N개 | Task 재시도, Cache·Shuffle Output 손실 |
| Worker Node | 물리·가상 머신 | Executor에 CPU와 Memory 제공 | 클러스터당 N대 | 그 위의 Executor 전부 손실 |
| Task | 실행 최소 단위 | Partition 하나 계산 | Stage당 Partition 수만큼 | 다른 Executor에서 재시도 |
- 수량 관계
Cluster 1개
├─ Cluster Manager 1개
├─ Worker Node N대
└─ Application M개 (동시 실행 가능)
├─ Driver 1개
└─ Executor N개
└─ 동시에 Core 수만큼 Task 실행- 헷갈리기 쉬운 경계
- Cluster Manager는 Executor를 어디에 띄울지까지만 관여, 어떤 Task를 어느 Executor에 보낼지는 Driver가 결정 (4.8 TaskScheduler)
- Worker Node ≠ Executor: 머신과 프로세스의 관계, 한 Node에 여러 Executor 가능 (2.7)
- Driver는 데이터를 직접 처리하지 않음 —
collect()로 결과를 모으는 순간만 예외이고, 이때 Driver OOM 위험 (2.6)
2.5 Spark Application Architecture
- 하나의 Spark Application은 기본적으로 하나의 Driver와 여러 Executor로 구성
Spark Application
Driver
SparkSession / SparkContext
Planning · Scheduling
│
│ Task Assignment
┌────────────┼────────────┐
↓ ↓ ↓
Executor 1 Executor 2 Executor 3
├─ Task ├─ Task ├─ Task
├─ Task ├─ Task ├─ Task
└─ Cache └─ Cache └─ Cache하나의 spark-submit 실행이 일반적으로 하나의 Spark Application이 됨.
spark-submit application.pySpark Application
├─ Driver Process
├─ Executor Process 1
├─ Executor Process 2
└─ Executor Process 3서로 다른 Spark Application은 일반적으로 독립적인 Driver와 Executor를 사용함.
Application A
├─ Driver A
├─ Executor A1
└─ Executor A2
Application B
├─ Driver B
├─ Executor B1
└─ Executor B2- Cluster Manager가 전체 클러스터 자원을 Application 사이에 배분
- 각 Executor는 특정 Spark Application에 속해 해당 Application의 Task만 실행
2.6 The Role of the Driver
Driver는 Spark Application의 전체 실행을 조율하는 프로세스임.
Driver
├─ 사용자 main 코드 실행
├─ SparkSession / SparkContext 생성
├─ 실행 계획 생성
├─ Job 생성
├─ Stage 분리
├─ Task 생성
├─ Executor에 Task 할당
├─ 실행 상태 추적
└─ 실패한 작업 재시도예시:
result = (
orders
.filter("status = 'PAID'")
.groupBy("region")
.count()
)
result.write.parquet(output_path)Driver가 수행하는 과정:
- DataFrame 연산 관계 생성
- Logical Plan 분석
- Optimized Plan 생성
- Physical Plan 선택
- Action 호출 시 Job 생성
- Job을 Stage와 Task로 분리
- Executor에 Task 전달
- 실행 결과와 상태 관리
Driver는 전체 데이터를 직접 처리하지 않음.
Driver
→ 실행 계획과 메타데이터 관리
Executor
→ 실제 데이터 파티션 처리Driver가 관리하는 정보
- RDD Lineage
- Logical Plan과 Physical Plan
- Job, Stage, Task 상태
- Executor 목록과 상태
- Broadcast metadata
- Task 재시도 정보
collect()를 호출하면 결과 데이터가 Driver로 이동함.
rows = dataframe.collect()Executor 1 ─┐
Executor 2 ─┼→ Driver Memory
Executor 3 ─┘결과 데이터가 너무 크면 Driver OOM이 발생할 수 있음.
2.7 The Role of Executors
- Executor는 Worker Node에서 실행되는 Spark 프로세스
- Driver가 전달한 Task를 실제로 수행
Worker Node
└─ Executor Process
├─ Core 1 → Task
├─ Core 2 → Task
├─ Core 3 → Task
├─ Core 4 → Task
└─ Memory
├─ Execution Memory
├─ Cached Data
└─ Shuffle DataExecutor의 역할
- 데이터 Partition 읽기
map,filter,join,aggregate실행- Shuffle Write
- Shuffle Read
- Cache와 Persist 데이터 저장
- 결과와 실행 상태를 Driver에 보고
Driver
↓ Task
Executor
↓
Read Partition
↓
Transform
↓
Shuffle / Cache / OutputExecutor의 CPU Core와 동시 Task 수
Executor with 4 Cores
Core 1 → Task A
Core 2 → Task B
Core 3 → Task C
Core 4 → Task D- 일반적으로 Executor Core 하나가 한 번에 Task 하나를 실행
- Executor에 Core가 4개라면 최대 4개의 Task를 동시에 실행 가능
용어 구분
Worker Node
= 물리 머신 또는 가상 머신
Executor
= Worker Node에서 실행되는 Spark 프로세스
Task
= Executor 안에서 실행되는 실제 작업 단위하나의 Worker Node에 여러 Executor가 실행될 수도 있음.
Worker Node
├─ Executor A
└─ Executor BKubernetes에서는 보통 Executor 하나가 Pod 하나로 실행됨.
Kubernetes Node
├─ Executor Pod 1
├─ Executor Pod 2
└─ Other Pods2.8 The Role of the Cluster Manager
- Cluster Manager는 Spark 연산을 직접 수행하지 않음
- 클러스터의 CPU와 Memory를 Spark Application에 할당하는 역할
Spark Driver
“Executor 5개가 필요합니다.”
↓
Cluster Manager
“사용 가능한 노드에 자원을 할당합니다.”
↓
Executor Processes역할 구분
Cluster Manager
→ Executor를 어디에 얼마나 실행할지 결정
Spark Driver
→ 각 Executor가 어떤 Task를 실행할지 결정예시:
Kubernetes
→ Executor Pod를 Node A에 배치
Spark Driver
→ 해당 Executor에 Partition 17 처리 Task 할당현재 주요 Cluster Manager:
- Spark Standalone
- Hadoop YARN
- Kubernetes
2.9 Standalone, YARN, and Kubernetes
1. Spark Standalone
- Spark가 자체적으로 제공하는 Cluster Manager
- 별도의 YARN이나 Kubernetes 없이 Spark 클러스터 구성 가능
Spark Standalone Cluster
Master
├─ Worker 1
├─ Worker 2
└─ Worker 3Master의 역할:
- 클러스터 자원 관리
- Worker 상태 추적
- Application에 Executor 자원 할당
Worker의 역할:
- Executor를 실행할 수 있는 노드
적합한 상황:
- Spark 전용 클러스터
- 비교적 단순한 온프레미스 환경
- 학습과 테스트 환경
2. Hadoop YARN
- Hadoop 생태계의 범용 자원 관리자
YARN Cluster
├─ Spark Application
├─ MapReduce Application
└─ Other Hadoop Applications- 기존에 HDFS와 Hadoop 클러스터를 운영 중인 환경에서 많이 사용
- 여러 종류의 분산 Application이 동일한 자원을 공유
적합한 상황:
- 기존 Hadoop 인프라 존재
- HDFS 중심 데이터 처리
- YARN Queue와 보안 정책 사용
3. Kubernetes
- 컨테이너 기반의 범용 오케스트레이터
- Spark Driver와 Executor를 Pod로 실행
Kubernetes Cluster
├─ Driver Pod
├─ Executor Pod 1
├─ Executor Pod 2
└─ Executor Pod 3역할:
Kubernetes
→ Pod가 어느 Node에서 실행될지 결정
Spark Driver
→ Executor Pod 생성 요청
→ Executor에 Task 할당적합한 상황:
- 기존 Kubernetes 인프라 존재
- Spark Application을 컨테이너 이미지로 관리
- Spark와 다른 서비스가 동일한 클러스터를 공유
- 클라우드 네이티브 운영 환경
2.10 Client Mode and Cluster Mode
Cluster Manager 선택과 Driver 실행 위치 선택은 서로 다른 개념임.
--master
→ 어떤 Cluster Manager를 사용할 것인가?
--deploy-mode
→ Driver를 어디에서 실행할 것인가?Client Mode
Driver가 spark-submit을 실행한 머신에서 동작함.
Client Machine
├─ spark-submit
└─ Driver
│
│ Network
↓
Cluster
├─ Executor 1
├─ Executor 2
└─ Executor 3spark-submit \
--master yarn \
--deploy-mode client \
application.py장점:
- Driver 로그와 출력에 바로 접근 가능
- 디버깅과 대화형 실행에 편리
- Notebook, Spark Shell 등에 적합
단점:
- Client Machine이 종료되면 Application에 영향
- Driver와 Executor 사이의 네트워크 거리가 멀 수 있음
Cluster Mode
Driver도 클러스터 내부에서 실행됨.
Client Machine
└─ spark-submit
↓ submit
Cluster
├─ Driver
├─ Executor 1
├─ Executor 2
└─ Executor 3spark-submit \
--master k8s://https://kubernetes-api-server \
--deploy-mode cluster \
application.py장점:
- 제출 머신이 종료돼도 Application 실행 가능
- Driver와 Executor가 같은 클러스터 네트워크에 위치
- 운영 배치 작업에 적합
2.11 Spark Connect: The Client–Server Architecture
기존 Spark는 사용자 Application과 Driver 환경이 강하게 결합돼 있었음.
Classic PySpark
Python Application
↓ Py4J
Driver JVM
↓
ExecutorsSpark Connect는 Client와 Spark Driver를 분리함.
flowchart TD client["Client Application<br/>Python · Scala · Notebook · IDE"] client -- "Logical Plan<br/>gRPC + Protocol Buffers" --> server["Spark Connect Server"] server --> optimization["Query Optimization"] optimization --> core["Spark Core"] core --> executors["Executors"]
사용자가 DataFrame 연산을 호출하면 Client에서 실제 데이터를 계산하지 않음.
Client는 Logical Plan을 생성해 Spark Connect Server로 전달함.
spark = (
SparkSession.builder
.remote("sc://spark-server:15002")
.getOrCreate()
)Client
df.filter(...)
.groupBy(...)
.count()
↓ Logical Plan 전송
Spark Server
Analyze
→ Optimize
→ Execute- 실행 결과는 주로 Arrow 형식으로 Client에 반환
- Client는 Driver JVM과 직접 결합되지 않음
- 서버의
SparkContext나 내부 JVM 객체에 직접 접근할 수 없음 - 선언적인 DataFrame과 SQL API 중심
- RDD API는 Spark Connect에서 지원하지 않음
의미:
Classic Spark
→ Spark를 Application 내부 라이브러리처럼 사용
Spark Connect
→ 원격 Spark 실행 서버를 여러 Client가 사용2.12 The Complete Modern Spark Architecture
Client Layer
Python / Scala / SQL / Notebook / Application
│
Classic Spark API or Spark Connect
↓
Driver Layer
SparkSession / SparkContext
↓
Logical Plan Analysis
↓
Catalyst Optimization
↓
Physical Planning
↓
Spark Core Scheduler
│
│ Resource Request
↓
Cluster Manager
Standalone / YARN / Kubernetes
│
│ Executor Allocation
↓
Execution Layer
┌────────────────┼────────────────┐
↓ ↓ ↓
Executor 1 Executor 2 Executor 3
├─ Task ├─ Task ├─ Task
├─ Cache ├─ Cache ├─ Cache
└─ Shuffle └─ Shuffle └─ Shuffle
│ │ │
└────────────────┼────────────────┘
↓
Data Layer
S3 / GCS / HDFS / Kafka
Parquet / Iceberg / Databases구성요소 정리
SparkSession
→ DataFrame과 SQL 작업의 진입점
SparkContext
→ Spark Application과 클러스터의 핵심 연결
Driver
→ 실행 계획을 만들고 Application 전체를 조율
Cluster Manager
→ CPU와 Memory 자원을 Application에 할당
Executor
→ 실제 데이터 Partition을 Task 단위로 계산
Worker Node
→ Executor가 실행되는 물리적 또는 가상 머신
Spark Connect Client
→ Logical Plan을 원격 Spark Server에 전달핵심 정리
Modern Spark
= DataFrame · SQL · Spark Connect
+ Query Planning and Optimization
+ Spark Core
+ Driver–Executor Architecture