3. RDD: The Core Abstraction Behind Spark

불변성, 파티션, lineage로 이해하는 Spark의 핵심 추상화

  • Haram Lee
  • 2026-08-16
  • studies / Topics / Spark
  • RDD는 초기 Spark의 핵심 API이자, Spark의 분산 실행 모델을 이해하기 위한 가장 기본적인 추상화
  • 단순히 데이터를 여러 서버에 나눈 자료구조가 아니라 다음을 함께 표현함
flowchart TD
  rdd["RDD"] --> partitions["Partitioned Data"]
  rdd --> computation["Computation Logic"]
  rdd --> dependencies["Dependency Information"]
  rdd --> recovery["Recovery Information"]
  • 중요한 점
    • RDD 객체가 모든 실제 데이터를 직접 담고 있는 것은 아님
    • RDD는 분산된 데이터를 어떻게 계산할 것인지 설명하는 논리적 객체에 가까움
    • 실제 데이터는 외부 저장소에 있거나, 실행 시 계산되거나, Executor의 메모리·디스크에 캐시될 수 있음

3.1 What Is an RDD?

RDD는 Resilient Distributed Dataset의 약자

  • 클러스터의 여러 노드에 파티셔닝되어 병렬로 처리할 수 있는 요소들의 컬렉션
  • 파일 시스템의 데이터로부터 생성하거나, Driver의 기존 컬렉션을 분산시키거나, 다른 RDD를 변환해 만들 수 있음.
text
Resilient
→ 장애가 발생해도 손실된 데이터를 복구할 수 있음

Distributed
→ 데이터가 여러 Partition으로 나뉘어 클러스터에서 처리됨

Dataset
→ 처리 대상이 되는 데이터 요소들의 집합

예시:

python
numbers = spark.sparkContext.parallelize(
    [1, 2, 3, 4, 5, 6],
    3
)
text
numbers RDD

Partition 0: [1, 2]
Partition 1: [3, 4]
Partition 2: [5, 6]

각 Partition은 서로 다른 Executor에서 병렬로 계산될 수 있음.

text
Partition 0 → Executor A
Partition 1 → Executor B
Partition 2 → Executor C

RDD는 다음과 같은 메타데이터도 함께 가지고 있음.

text
RDD
├─ Partitions
├─ Dependencies
├─ Compute Function
├─ Partitioner
└─ Preferred Locations

3.2 RDD as a Computation Blueprint

다음 코드를 보자.

python
source = spark.sparkContext.textFile("s3://logs/")

errors = (
    source
    .filter(lambda line: "ERROR" in line)
    .map(parse_log)
)

이 코드가 실행됐다고 해서 errors의 모든 데이터가 즉시 계산되어 Driver 메모리에 들어오는 것은 아님.

Spark는 우선 다음과 같은 RDD 객체와 관계를 만듦.

text
Text File
    ↓
HadoopRDD
    ↓ filter
MapPartitionsRDD
    ↓ map
MapPartitionsRDD

각 RDD는 다음 질문에 대한 답을 가지고 있음.

text
이 데이터는 어디서 오는가?
부모 RDD는 무엇인가?
몇 개의 Partition으로 나뉘는가?
Partition 하나를 어떻게 계산하는가?

예를 들어 filter()map()을 호출하면 기존 RDD를 수정하는 것이 아니라 부모 RDD를 참조하는 새로운 MapPartitionsRDD가 만들어짐.

text
source RDD
    │
    │ filter
    ▼
filtered RDD
    │
    │ map
    ▼
parsed RDD

따라서 RDD는 실제 데이터 자체이면서 동시에:

필요한 데이터가 아직 계산되지 않았다면 어떻게 만들어야 하는지를 설명하는 계산 설계도

라고 볼 수 있음.


3.3 How an RDD Is Implemented

Spark Core에서 RDD는 Scala의 추상 클래스로 구현돼 있음.

실제 코드를 단순화하면 다음 구조에 가까움.

scala
abstract class RDD[T] {

  protected def getPartitions: Array[Partition]

  protected def getDependencies: Seq[Dependency[_]]

  def compute(
      split: Partition,
      context: TaskContext
  ): Iterator[T]

  val partitioner: Option[Partitioner]

  protected def getPreferredLocations(
      split: Partition
  ): Seq[String]
}

getPartitions

RDD가 어떤 Partition들로 구성되는지 반환함.

text
RDD
├─ Partition 0
├─ Partition 1
└─ Partition 2

Partition은 실제 데이터 전체를 담는 객체라기보다, RDD 안의 특정 데이터 조각을 식별하기 위한 논리적 정보에 가까움.

예를 들어 파일 기반 RDD의 Partition은 다음 정보를 가질 수 있음.

text
File Partition
├─ File Path
├─ Start Offset
└─ Length

compute()

특정 Partition을 실제로 계산하는 방법을 정의함.

scala
def compute(
    partition: Partition,
    context: TaskContext
): Iterator[T]

개념적으로 map()이 만든 RDD의 compute()는 다음과 유사함.

scala
override def compute(
    split: Partition,
    context: TaskContext
): Iterator[U] = {
  parent.iterator(split, context).map(f)
}

즉:

text
부모 RDD의 해당 Partition 읽기
→ map 함수 적용
→ 결과 Iterator 반환

반환값이 전체 ListArray가 아니라 Iterator라는 점이 중요함.

text
Partition 전체를 중간 컬렉션으로 생성
→ 다음 연산에 전달

하는 대신,

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

처럼 데이터를 한 레코드씩 흘려보낼 수 있음.

getDependencies

이 RDD가 어떤 부모 RDD에 의존하는지 표현함.

text
Child RDD
    ↓ depends on
Parent RDD

Spark는 이 Dependency 정보를 바탕으로 RDD DAG를 만들고, Shuffle 경계를 찾아 Stage를 구분함.

partitioner

Key–Value RDD에서 각 Key가 어느 Partition으로 들어갈지 결정하는 규칙임.

대표적으로:

text
HashPartitioner
RangePartitioner

등이 있음.

getPreferredLocations

특정 Partition을 어느 노드에서 계산하는 것이 유리한지 알려줌.

text
Data stored on Node A
        ↓
Prefer Task on Node A

데이터가 위치한 곳에 가까운 Executor에서 Task를 실행하면 네트워크를 통한 데이터 이동을 줄일 수 있음.


3.4 Resilient: Recovering from Failures

RDD의 Resilient데이터를 무조건 여러 개 복제해서 보관한다는 뜻이 아님.

RDD는 데이터가 만들어진 연산 관계인 Lineage를 기억해 손실된 Partition을 다시 계산할 수 있음.

python
source = sc.textFile("s3://logs/")
parsed = source.map(parse_log)
errors = parsed.filter(lambda x: x.level == "ERROR")
text
source RDD
    ↓ map(parse_log)
parsed RDD
    ↓ filter(level == ERROR)
errors RDD

이 관계가 Lineage임.

예를 들어 errors RDD의 Partition 2가 Executor 장애로 손실됐다고 하자.

text
errors RDD

Partition 0  정상
Partition 1  정상
Partition 2  손실
Partition 3  정상

Spark는 Lineage를 거슬러 올라가 필요한 Partition만 다시 계산할 수 있음.

text
source Partition 2 다시 읽기
        ↓ parse_log
parsed Partition 2 재생성
        ↓ filter
errors Partition 2 재생성
text
전체 RDD 재계산 ❌

손실된 Partition의 계보만 재계산 ⭕

Cache가 손실된 경우

python
errors.cache()

Cache된 Partition도 Executor가 종료되면 사라질 수 있음.

하지만 원본 데이터와 Lineage가 남아 있다면 다시 계산할 수 있음.

text
Cache
→ 빠른 재사용을 위한 임시 저장
→ 유실될 수 있음

Lineage
→ 유실된 Partition을 다시 만드는 방법

Lineage가 너무 길어지는 경우

Transformation이 너무 많이 연결되면 재계산 비용과 스케줄링 정보가 커질 수 있음.

text
RDD 1
 ↓
RDD 2
 ↓
RDD 3
 ↓
...
 ↓
RDD 1000

이럴 때 checkpoint()를 사용해 데이터를 안정적인 저장소에 실제로 기록하고 이전 Lineage를 끊을 수 있음.

text
Long Lineage
      ↓
Checkpoint
      ↓
New Lineage Start
text
Cache
→ 빠른 재사용이 목적
→ Lineage는 유지됨

Checkpoint
→ 안정적인 저장소에 실제 데이터 기록
→ 긴 Lineage를 단절

3.5 Distributed: Partitioning Data Across the Cluster

RDD는 하나의 거대한 데이터 묶음으로 처리되지 않고 여러 Partition으로 나뉨.

text
RDD

├─ Partition 0
├─ Partition 1
├─ Partition 2
└─ Partition 3

Partition은 Spark의 병렬 처리 수준을 결정하는 핵심 단위

text
Partition 0 → Task 0
Partition 1 → Task 1
Partition 2 → Task 2
Partition 3 → Task 3

Partition과 저장 위치

RDD Partition이 항상 특정 Executor에 고정되어 있는 것은 아님.

text
Partition
= 논리적 데이터 조각

Executor
= Partition을 실제로 계산하는 프로세스

Task가 실행될 때 Scheduler가 Partition 계산을 특정 Executor에 할당함.

Cache한 경우에는 계산된 Partition 데이터가 특정 Executor의 BlockManager에 저장될 수 있음.

text
Before Computation

RDD Partition 0
→ 계산 방법만 존재
text
After cache()

Executor A
└─ Cached Block: RDD 3, Partition 0

Partition은 어떻게 만들어지는가?

  • 외부 파일의 Split
  • parallelize()에서 지정한 Slice 수
  • 부모 RDD의 Partition 상속
  • Shuffle 결과
  • repartition()
  • coalesce()

예시:

python
rdd = sc.parallelize(data, 4)
text
4 Partitions
python
repartitioned = rdd.repartition(10)
text
10 Partitions

3.6 Dataset: Representing Distributed Data

RDD에서 Dataset은 Spark SQL의 Dataset 클래스와 반드시 같은 의미로 보면 안 됨.

text
RDD의 Dataset
→ 데이터 요소들의 집합이라는 일반적인 의미

Spark SQL Dataset API
→ 스키마와 Encoder를 가진 고수준 API

RDD에는 어떤 형태의 객체든 담을 수 있음.

python
numbers: RDD[int]
lines: RDD[str]
pairs: RDD[tuple[str, int]]
orders: RDD[Order]
text
RDD[Int]
RDD[String]
RDD[(String, Int)]
RDD[Order]

이 유연성 덕분에 비정형 데이터나 사용자 정의 객체를 자유롭게 처리할 수 있음.


3.7 Immutability and Transformations

RDD는 Immutable, 즉 생성된 뒤 그 자체가 수정되지 않음.

python
numbers = sc.parallelize([1, 2, 3])

doubled = numbers.map(lambda x: x * 2)

이 코드에서 numbers의 값이 바뀌는 것이 아님.

text
numbers RDD
[1, 2, 3]

doubled RDD
[2, 4, 6]

map()은 기존 RDD를 수정하는 대신, 기존 RDD에 의존하는 새로운 RDD를 반환함.

text
numbers RDD
     │
     │ map(x × 2)
     ▼
doubled RDD

다음 Transformation도 동일함.

python
filtered = doubled.filter(lambda x: x > 2)
text
numbers RDD
     ↓ map
doubled RDD
     ↓ filter
filtered RDD

불변성이 필요한 이유

분산 환경에서 여러 Executor가 동일한 데이터를 직접 수정하도록 허용하면 다음 문제가 생길 수 있음.

text
Executor A가 데이터 수정
Executor B가 같은 데이터 수정
Executor C는 수정 전 데이터를 읽음
  • 동시성 제어
  • Lock
  • 데이터 버전 관리
  • 장애 복구
  • 변경 순서 관리

가 복잡해짐.

RDD는 기존 데이터 수정 대신 새 데이터셋을 만드는 모델을 사용함.

text
Input RDD
→ Transformation
→ New RDD

따라서 데이터가 어떻게 변했는지 Lineage로 명확하게 표현할 수 있음.

text
RDD A
  ↓ Transformation 1
RDD B
  ↓ Transformation 2
RDD C

불변성과 Lineage는 서로 연결됨.

text
Immutability
→ 이전 RDD가 변하지 않음
→ Transformation 관계를 안정적으로 추적 가능
→ 손실된 Partition을 재계산 가능

초기 RDD 논문도 RDD를 immutable하고 partitioned된 레코드 컬렉션으로 정의함. (MIT CSAIL)


3.8 Transformations, Actions, and Lazy Evaluation

RDD 연산은 크게 TransformationAction으로 나뉨.

Transformation

기존 RDD로부터 새로운 RDD를 만드는 연산.

python
mapped = rdd.map(...)
filtered = mapped.filter(...)
pairs = filtered.flatMap(...)

대표적인 Transformation:

  • map
  • filter
  • flatMap
  • union
  • distinct
  • reduceByKey
  • groupByKey
  • join
  • repartition
text
RDD A
  ↓ Transformation
RDD B

Transformation을 호출해도 바로 전체 데이터가 계산되지 않음.

Action

RDD의 실제 결과를 요구하는 연산임.

python
rdd.count()
rdd.collect()
rdd.take(10)
rdd.reduce(...)
rdd.saveAsTextFile(...)

대표적인 Action:

  • count
  • collect
  • take
  • first
  • reduce
  • saveAsTextFile
  • foreach
text
RDD Lineage
    ↓ Action
Spark Job
    ↓
Actual Computation

예시:

python
result = (
    sc.textFile("s3://logs/")
      .filter(lambda x: "ERROR" in x)
      .map(parse_log)
)

여기까지는 Transformation만 존재함.

text
No Job Yet

Text File
   ↓ filter
   ↓ map
Result RDD

이후:

python
result.count()

가 호출되면 실제 Job이 시작됨.

text
Action: count()
       ↓
Job Created
       ↓
Tasks Sent to Executors

왜 지연 실행하는가?

1. 여러 연산을 파이프라인으로 실행

python
rdd.map(f).filter(g).map(h)
text
Record 1 → f → g → h
Record 2 → f → g → h
Record 3 → f → g → h

각 연산 결과 전체를 별도로 저장할 필요가 없음.

2. 필요하지 않은 데이터는 계산하지 않을 수 있음

python
rdd.take(10)

전체 데이터가 아니라 결과 10개를 얻는 데 필요한 Partition만 먼저 계산할 수 있음.

3. 실행 전에 전체 의존 관계를 확인

Spark는 최종 RDD에서 부모 RDD 방향으로 Dependency를 따라가며 어떤 작업이 필요한지 파악함.

text
Final RDD
   ↓ parent
RDD C
   ↓ parent
RDD B
   ↓ parent
RDD A

다만 RDD의 Lazy Evaluation을 DataFrame의 Catalyst Optimization과 완전히 동일하게 보면 안 됨.

text
RDD Lazy Evaluation
→ 연산 연결, 필요한 계산 결정, Stage 구성의 기반

DataFrame Lazy Evaluation
→ 위 기능에 더해 SQL 의미를 이용한 계획 최적화 가능

RDD 내부의 임의 Python·Scala 함수는 Spark가 의미를 분석하기 어렵기 때문에, 컬럼 가지치기나 조인 재배치 같은 고수준 최적화에는 제한이 있음.


3.9 RDD Lineage

RDD Lineage는 RDD가 어떤 Source와 Transformation을 거쳐 만들어졌는지를 나타냄.

python
lines = sc.textFile("logs.txt")

errors = (
    lines
    .filter(lambda x: "ERROR" in x)
    .map(lambda x: (extract_service(x), 1))
    .reduceByKey(lambda a, b: a + b)
)

Lineage는 다음과 같음.

text
TextFile RDD
     ↓ filter
MapPartitionsRDD
     ↓ map
MapPartitionsRDD
     ↓ reduceByKey
ShuffledRDD

각 RDD는 직접 모든 이전 연산 문자열을 저장하는 단순한 기록이 아니라, 부모 RDD에 대한 Dependency를 가지고 있음.

text
RDD C
└─ Dependency on RDD B

RDD B
└─ Dependency on RDD A

이 연결을 따라가면 전체 Lineage가 만들어짐.

Lineage의 역할

  1. 실행할 연산 관계 표현
  2. Stage를 나누기 위한 Dependency 정보 제공
  3. 장애 발생 시 손실된 Partition 재계산
  4. Cache가 없는 RDD를 필요할 때 다시 계산
text
Lineage
├─ Execution Planning
├─ Stage Construction
└─ Fault Recovery

Lineage와 DAG

Lineage는 일반적으로 DAG 형태로 구성됨.

text
    RDD A
    /   \
   /     \
RDD B   RDD C
   \     /
    \   /
     RDD D

DAG는 Directed Acyclic Graph, 즉 방향은 있지만 순환하지 않는 그래프임.

text
Directed
→ 부모 RDD에서 자식 RDD 방향으로 연산 관계가 존재

Acyclic
→ RDD D가 다시 자기 부모로 돌아가는 순환 관계가 없음

Graph
→ 여러 RDD와 Dependency가 노드와 간선으로 표현됨

RDD는 불변이고 Transformation은 새 RDD를 생성하기 때문에 DAG 관계를 안정적으로 구성할 수 있음.


3.10 Building a DAG from RDD Operations

다음 코드를 보자.

python
orders = sc.textFile("orders.csv")

paid_orders = (
    orders
    .map(parse_order)
    .filter(lambda x: x.status == "PAID")
)

region_amounts = paid_orders.map(
    lambda x: (x.region, x.amount)
)

totals = region_amounts.reduceByKey(
    lambda a, b: a + b
)

totals.collect()

RDD 관계:

text
orders
  ↓ map(parse_order)
parsed_orders
  ↓ filter(status == PAID)
paid_orders
  ↓ map(region, amount)
region_amounts
  ↓ reduceByKey
totals

DAG로 표현하면:

text
TextFileRDD
     ↓
MapPartitionsRDD
     ↓
MapPartitionsRDD
     ↓
MapPartitionsRDD
     ↓
ShuffledRDD

Spark는 Action인 collect()가 호출되면 최종 RDD인 totals부터 부모 RDD 방향으로 Dependency를 추적함.

text
collect(totals)
       ↓
totals는 region_amounts에 의존
       ↓
region_amounts는 paid_orders에 의존
       ↓
paid_orders는 orders에 의존

그 후 Dependency 종류에 따라 실제 실행 Stage를 나눔. RDD DAG와 Stage DAG의 구분은 4.3에서 실행 관점으로 정리함.


3.11 Common RDD Implementations

RDD는 추상 클래스이고, 데이터 소스나 Transformation 종류에 따라 서로 다른 구현체가 사용됨.

text
RDD
├─ ParallelCollectionRDD
├─ HadoopRDD
├─ MapPartitionsRDD
├─ ShuffledRDD
├─ UnionRDD
└─ CoalescedRDD

ParallelCollectionRDD

Driver의 로컬 컬렉션을 여러 Partition으로 나눌 때 사용됨.

python
rdd = sc.parallelize([1, 2, 3, 4], 2)
text
ParallelCollectionRDD
├─ Partition 0: [1, 2]
└─ Partition 1: [3, 4]

HadoopRDD

HDFS, S3 등 Hadoop InputFormat과 호환되는 데이터 소스를 읽을 때 사용됨.

python
rdd = sc.textFile("s3://bucket/logs/")
text
HadoopRDD
├─ File Split 0
├─ File Split 1
└─ File Split 2

MapPartitionsRDD

map, filter, flatMap 같은 다수의 Transformation에서 사용됨.

python
mapped = rdd.map(lambda x: x * 2)
text
Parent RDD
    ↓ map
MapPartitionsRDD

ShuffledRDD

Key별 재분배가 필요한 Shuffle 결과를 표현함.

python
result = pairs.reduceByKey(lambda a, b: a + b)
text
Parent Partitions
       ↓ Shuffle
ShuffledRDD

UnionRDD

여러 RDD를 하나로 합침.

python
combined = rdd1.union(rdd2)
text
RDD 1 ─┐
       ├→ UnionRDD
RDD 2 ─┘

이처럼 RDD는 하나의 고정된 데이터 컨테이너가 아니라:

서로 다른 데이터 소스와 연산을 동일한 Partition·Dependency·compute() 모델로 실행할 수 있게 만드는 공통 추상 클래스

라고 볼 수 있음.


3.12 The Relationship Between RDDs and DataFrames

현재 Spark 사용자는 대부분 RDD보다 DataFrame과 SQL을 사용함.

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

result = (
    orders
    .filter("status = 'PAID'")
    .groupBy("region")
    .sum("amount")
)

DataFrame은 Spark에 다음 정보를 제공함.

text
Schema
├─ Column Names
├─ Data Types
└─ Nullable Information

Operations
├─ Filter Expressions
├─ Join Conditions
├─ Aggregations
└─ Projections

Spark는 이를 Logical Plan으로 표현하고 Catalyst Optimizer로 최적화함.

text
DataFrame / SQL
       ↓
Logical Plan
       ↓
Optimized Logical Plan
       ↓
Physical Plan
       ↓
Distributed Execution

RDD에서는 사용자 함수가 블랙박스에 가까움.

python
rdd.filter(lambda x: custom_condition(x))
text
Spark
→ 함수가 호출된다는 사실은 알 수 있음
→ 함수 내부의 의미는 분석하기 어려움

DataFrame에서는 필터가 명시적인 Expression임.

python
df.filter(df.status == "PAID")
text
Spark
→ status 컬럼을 사용함
→ PAID 조건으로 필터링함
→ Predicate Pushdown 가능 여부 판단

RDD와 DataFrame 비교

구분RDDDataFrame
데이터 표현JVM·Python 객체스키마가 있는 Row·Column
API 성격함수형·저수준선언적·고수준
Spark의 데이터 이해제한적컬럼과 타입을 이해
최적화제한적Catalyst Optimizer
직렬화 비용객체 중심내부 Binary·Columnar 표현 가능
일반적인 사용특수한 저수준 처리ETL·SQL·분석의 기본 선택

PySpark 공식 Quickstart는 DataFrame이 Lazy하게 평가되고 RDD 위에 구현되어 있다고 설명함. 현대 Spark SQL에서는 Logical·Physical Plan이 최적화된 뒤 내부적으로 InternalRow 또는 컬럼형 배치를 처리하는 분산 실행으로 연결됨. (Apache Spark)

text
DataFrame / SQL
       ↓
Query Planning and Optimization
       ↓
Spark Physical Operators
       ↓
Distributed Partition Processing
       ↓
Tasks on Executors

따라서 다음처럼 정리할 수 있음.

text
RDD
→ Spark의 저수준 분산 데이터 API
→ Partition, Dependency, Lineage, Compute 모델 제공

DataFrame
→ 스키마와 선언적 연산을 제공
→ Spark가 실행 계획을 최적화할 수 있게 함

Spark Core
→ 최종 작업을 Partition과 Task 단위로 분산 실행

RDD가 일반 사용자의 주력 API는 아니더라도 여전히 중요한 이유:

  • Spark의 Partition 개념을 설명함
  • Lineage와 장애 복구 모델을 설명함
  • Narrow·Wide Dependency를 설명함
  • Stage와 Task가 만들어지는 기반을 설명함
  • DataFrame 아래의 분산 실행 원리를 이해하게 함

3.13 The Complete RDD Model

flowchart TD
  rdd["RDD"] --> partitions["Partitions<br/>데이터 조각"]
  rdd --> dependencies["Dependencies<br/>부모 RDD 관계"]
  rdd --> compute["Compute Function<br/>Partition 계산법"]
  partitions --> lineage["RDD Lineage<br/>데이터가 만들어지는 관계"]
  dependencies --> lineage
  compute --> lineage
  lineage --> lazy["Lazy Evaluation<br/>Action 전까지 실행 보류"]
  lazy -- "Action" --> scheduler["DAGScheduler 분석<br/>Dependency와 Shuffle 경계 확인"]
  scheduler --> tasks["Stage / Task 생성"]
  tasks --> executors["Executors가 Partition 계산"]
  executors -- "failure" --> recompute["Lineage로 손실 Partition 재계산"]
Discussion