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를 직접 사용하는 방식이 중심
python
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에 전달
text
User Application
       ↓
SparkContext
       ↓
RDD Transformations
       ↓
RDD Lineage / DAG
       ↓
Spark Core Scheduler
       ↓
Driver ───────── Executors
  • 사용자가 map, filter, reduceByKey 등의 Transformation을 호출
  • 각 연산은 기존 RDD를 수정하는 것이 아니라 새로운 RDD를 생성
  • RDD 사이의 의존 관계가 Lineage와 DAG를 구성
text
HadoopRDD
    ↓ map
MapPartitionsRDD
    ↓ filter
MapPartitionsRDD
    ↓ reduceByKey
ShuffledRDD
  • Spark Core는 RDD DAG를 보고 다음을 결정함.
    • 각 RDD가 어떤 파티션으로 구성되는가?
    • 어떤 부모 RDD를 계산해야 하는가?
    • 어느 지점에서 Shuffle이 발생하는가?
    • 실패한 파티션을 어떻게 재계산할 것인가?

RDD API의 한계

Spark가 사용자 함수 내부의 의미를 알기 어려움.

python
rdd.map(lambda row: custom_function(row))

Spark가 알 수 있는 것

  • 각 레코드에 함수를 적용해야 한다는 사실

Spark가 알기 어려운 것

  • 어떤 컬럼을 사용하는지
  • 필터를 더 앞에서 실행할 수 있는지
  • 조인 순서를 변경할 수 있는지
  • 일부 연산을 데이터 소스에 내려보낼 수 있는지

2.2 The Modern DataFrame- and SQL-Centric Architecture

  • 현재 Spark에서는 일반적으로 SparkContext보다 SparkSession을 진입점으로 사용
python
from pyspark.sql import SparkSession

spark = (
    SparkSession.builder
    .appName("order-analysis")
    .getOrCreate()
)

SparkSession

  • DataFrame API의 진입점
  • Spark SQL 실행
  • 데이터 소스 읽기와 쓰기
  • Catalog, Table, View 관리
  • 내부적으로 SparkContext를 포함
text
SparkSession
├─ DataFrame API
├─ Spark SQL
├─ Catalog
└─ SparkContext

사용자는 RDD의 각 레코드에 어떤 함수를 실행할지보다, 원하는 결과를 선언적으로 표현함.

python
result = (
    spark.read.parquet("s3://orders/")
    .filter("status = 'PAID'")
    .groupBy("region")
    .sum("amount")
)
sql
SELECT region, SUM(amount)
FROM orders
WHERE status = 'PAID'
GROUP BY region;

Spark는 DataFrame과 SQL을 바로 실행하지 않고 실행 계획으로 변환함.

text
SQL / DataFrame API
         ↓
Unresolved Logical Plan
         ↓
Analyzed Logical Plan
         ↓
Optimized Logical Plan
         ↓
Physical Plan
         ↓
Spark Core
         ↓
Job / Stage / Task
         ↓
Executors

DataFrame은 스키마와 연산 구조를 Spark에 제공함.

text
orders

user_id : LONG
status  : STRING
region  : STRING
amount  : DECIMAL

Spark가 알 수 있는 정보:

  • 어떤 컬럼이 필요한가?
  • 어떤 조건으로 필터링하는가?
  • 어떤 Key로 조인하는가?
  • 어떤 컬럼을 집계하는가?
  • 데이터의 예상 크기는 어느 정도인가?

이를 통해 Spark가 실행 방식을 최적화할 수 있음.

text
사용자가 표현한 흐름

모든 컬럼 읽기
→ status 필터
→ region, amount 선택
text
최적화된 흐름

status, region, amount만 읽기
→ 가능한 경우 파일 읽기 단계에서 필터

2.3 Spark Core and High-Level APIs

현재 Spark는 여러 고수준 API를 제공함.

text
Spark SQL
DataFrame
Dataset
Structured Streaming
MLlib
GraphX

이들은 서로 완전히 별개의 실행 엔진이 아니라, 대부분 Spark Core 위에서 동작함.

text
High-Level APIs
├─ Spark SQL
├─ DataFrame
├─ Structured Streaming
└─ MLlib
        ↓
Query Planning / Library Logic
        ↓
Spark Core
        ↓
Driver / Executor

Spark Core가 제공하는 핵심 기능

  • 분산 데이터 파티셔닝
  • Task 스케줄링
  • Executor 관리
  • Shuffle
  • Cache와 Persist
  • Broadcast
  • 장애 복구

DataFrame과 SQL이 RDD를 완전히 제거한 것은 아님.

text
DataFrame / SQL
       ↓
Logical Plan
       ↓
Physical Plan
       ↓
RDD[InternalRow]
또는
RDD[ColumnarBatch]
       ↓
Task Execution

정리하면:

text
RDD
→ 저수준 분산 데이터 API이자 Spark Core 실행 모델

DataFrame / SQL
→ Spark Core 위에서 실행 계획 최적화를 제공하는 고수준 API

2.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 / SparkContextApplication 코드의 진입점 객체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에서 재시도
  • 수량 관계
text
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로 구성
text
           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이 됨.

bash
spark-submit application.py
text
Spark Application
├─ Driver Process
├─ Executor Process 1
├─ Executor Process 2
└─ Executor Process 3

서로 다른 Spark Application은 일반적으로 독립적인 Driver와 Executor를 사용함.

text
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의 전체 실행을 조율하는 프로세스임.

text
Driver
├─ 사용자 main 코드 실행
├─ SparkSession / SparkContext 생성
├─ 실행 계획 생성
├─ Job 생성
├─ Stage 분리
├─ Task 생성
├─ Executor에 Task 할당
├─ 실행 상태 추적
└─ 실패한 작업 재시도

예시:

python
result = (
    orders
    .filter("status = 'PAID'")
    .groupBy("region")
    .count()
)

result.write.parquet(output_path)

Driver가 수행하는 과정:

  1. DataFrame 연산 관계 생성
  2. Logical Plan 분석
  3. Optimized Plan 생성
  4. Physical Plan 선택
  5. Action 호출 시 Job 생성
  6. Job을 Stage와 Task로 분리
  7. Executor에 Task 전달
  8. 실행 결과와 상태 관리

Driver는 전체 데이터를 직접 처리하지 않음.

text
Driver
→ 실행 계획과 메타데이터 관리

Executor
→ 실제 데이터 파티션 처리

Driver가 관리하는 정보

  • RDD Lineage
  • Logical Plan과 Physical Plan
  • Job, Stage, Task 상태
  • Executor 목록과 상태
  • Broadcast metadata
  • Task 재시도 정보

collect()를 호출하면 결과 데이터가 Driver로 이동함.

python
rows = dataframe.collect()
text
Executor 1 ─┐
Executor 2 ─┼→ Driver Memory
Executor 3 ─┘

결과 데이터가 너무 크면 Driver OOM이 발생할 수 있음.


2.7 The Role of Executors

  • Executor는 Worker Node에서 실행되는 Spark 프로세스
  • Driver가 전달한 Task를 실제로 수행
text
Worker Node
└─ Executor Process
   ├─ Core 1 → Task
   ├─ Core 2 → Task
   ├─ Core 3 → Task
   ├─ Core 4 → Task
   └─ Memory
      ├─ Execution Memory
      ├─ Cached Data
      └─ Shuffle Data

Executor의 역할

  • 데이터 Partition 읽기
  • map, filter, join, aggregate 실행
  • Shuffle Write
  • Shuffle Read
  • Cache와 Persist 데이터 저장
  • 결과와 실행 상태를 Driver에 보고
text
Driver
  ↓ Task
Executor
  ↓
Read Partition
  ↓
Transform
  ↓
Shuffle / Cache / Output

Executor의 CPU Core와 동시 Task 수

text
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를 동시에 실행 가능

용어 구분

text
Worker Node
= 물리 머신 또는 가상 머신

Executor
= Worker Node에서 실행되는 Spark 프로세스

Task
= Executor 안에서 실행되는 실제 작업 단위

하나의 Worker Node에 여러 Executor가 실행될 수도 있음.

text
Worker Node
├─ Executor A
└─ Executor B

Kubernetes에서는 보통 Executor 하나가 Pod 하나로 실행됨.

text
Kubernetes Node
├─ Executor Pod 1
├─ Executor Pod 2
└─ Other Pods

2.8 The Role of the Cluster Manager

  • Cluster Manager는 Spark 연산을 직접 수행하지 않음
  • 클러스터의 CPU와 Memory를 Spark Application에 할당하는 역할
text
Spark Driver
“Executor 5개가 필요합니다.”
          ↓
Cluster Manager
“사용 가능한 노드에 자원을 할당합니다.”
          ↓
Executor Processes

역할 구분

text
Cluster Manager
→ Executor를 어디에 얼마나 실행할지 결정

Spark Driver
→ 각 Executor가 어떤 Task를 실행할지 결정

예시:

text
Kubernetes
→ Executor Pod를 Node A에 배치

Spark Driver
→ 해당 Executor에 Partition 17 처리 Task 할당

현재 주요 Cluster Manager:

  1. Spark Standalone
  2. Hadoop YARN
  3. Kubernetes

2.9 Standalone, YARN, and Kubernetes

1. Spark Standalone

  • Spark가 자체적으로 제공하는 Cluster Manager
  • 별도의 YARN이나 Kubernetes 없이 Spark 클러스터 구성 가능
text
Spark Standalone Cluster

Master
├─ Worker 1
├─ Worker 2
└─ Worker 3

Master의 역할:

  • 클러스터 자원 관리
  • Worker 상태 추적
  • Application에 Executor 자원 할당

Worker의 역할:

  • Executor를 실행할 수 있는 노드

적합한 상황:

  • Spark 전용 클러스터
  • 비교적 단순한 온프레미스 환경
  • 학습과 테스트 환경

2. Hadoop YARN

  • Hadoop 생태계의 범용 자원 관리자
text
YARN Cluster
├─ Spark Application
├─ MapReduce Application
└─ Other Hadoop Applications
  • 기존에 HDFS와 Hadoop 클러스터를 운영 중인 환경에서 많이 사용
  • 여러 종류의 분산 Application이 동일한 자원을 공유

적합한 상황:

  • 기존 Hadoop 인프라 존재
  • HDFS 중심 데이터 처리
  • YARN Queue와 보안 정책 사용

3. Kubernetes

  • 컨테이너 기반의 범용 오케스트레이터
  • Spark Driver와 Executor를 Pod로 실행
text
Kubernetes Cluster
├─ Driver Pod
├─ Executor Pod 1
├─ Executor Pod 2
└─ Executor Pod 3

역할:

text
Kubernetes
→ Pod가 어느 Node에서 실행될지 결정

Spark Driver
→ Executor Pod 생성 요청
→ Executor에 Task 할당

적합한 상황:

  • 기존 Kubernetes 인프라 존재
  • Spark Application을 컨테이너 이미지로 관리
  • Spark와 다른 서비스가 동일한 클러스터를 공유
  • 클라우드 네이티브 운영 환경

2.10 Client Mode and Cluster Mode

Cluster Manager 선택과 Driver 실행 위치 선택은 서로 다른 개념임.

text
--master
→ 어떤 Cluster Manager를 사용할 것인가?

--deploy-mode
→ Driver를 어디에서 실행할 것인가?

Client Mode

Driver가 spark-submit을 실행한 머신에서 동작함.

text
Client Machine
├─ spark-submit
└─ Driver
       │
       │ Network
       ↓
Cluster
├─ Executor 1
├─ Executor 2
└─ Executor 3
bash
spark-submit \
  --master yarn \
  --deploy-mode client \
  application.py

장점:

  • Driver 로그와 출력에 바로 접근 가능
  • 디버깅과 대화형 실행에 편리
  • Notebook, Spark Shell 등에 적합

단점:

  • Client Machine이 종료되면 Application에 영향
  • Driver와 Executor 사이의 네트워크 거리가 멀 수 있음

Cluster Mode

Driver도 클러스터 내부에서 실행됨.

text
Client Machine
└─ spark-submit
       ↓ submit

Cluster
├─ Driver
├─ Executor 1
├─ Executor 2
└─ Executor 3
bash
spark-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 환경이 강하게 결합돼 있었음.

text
Classic PySpark

Python Application
       ↓ Py4J
Driver JVM
       ↓
Executors

Spark 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로 전달함.

python
spark = (
    SparkSession.builder
    .remote("sc://spark-server:15002")
    .getOrCreate()
)
text
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에서 지원하지 않음

의미:

text
Classic Spark
→ Spark를 Application 내부 라이브러리처럼 사용

Spark Connect
→ 원격 Spark 실행 서버를 여러 Client가 사용

2.12 The Complete Modern Spark Architecture

text
                   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

구성요소 정리

text
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에 전달

핵심 정리

text
Modern Spark
= DataFrame · SQL · Spark Connect
+ Query Planning and Optimization
+ Spark Core
+ Driver–Executor Architecture
Discussion