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를 변환해 만들 수 있음.
Resilient
→ 장애가 발생해도 손실된 데이터를 복구할 수 있음
Distributed
→ 데이터가 여러 Partition으로 나뉘어 클러스터에서 처리됨
Dataset
→ 처리 대상이 되는 데이터 요소들의 집합예시:
numbers = spark.sparkContext.parallelize(
[1, 2, 3, 4, 5, 6],
3
)numbers RDD
Partition 0: [1, 2]
Partition 1: [3, 4]
Partition 2: [5, 6]각 Partition은 서로 다른 Executor에서 병렬로 계산될 수 있음.
Partition 0 → Executor A
Partition 1 → Executor B
Partition 2 → Executor CRDD는 다음과 같은 메타데이터도 함께 가지고 있음.
RDD
├─ Partitions
├─ Dependencies
├─ Compute Function
├─ Partitioner
└─ Preferred Locations3.2 RDD as a Computation Blueprint
다음 코드를 보자.
source = spark.sparkContext.textFile("s3://logs/")
errors = (
source
.filter(lambda line: "ERROR" in line)
.map(parse_log)
)이 코드가 실행됐다고 해서 errors의 모든 데이터가 즉시 계산되어 Driver 메모리에 들어오는 것은 아님.
Spark는 우선 다음과 같은 RDD 객체와 관계를 만듦.
Text File
↓
HadoopRDD
↓ filter
MapPartitionsRDD
↓ map
MapPartitionsRDD각 RDD는 다음 질문에 대한 답을 가지고 있음.
이 데이터는 어디서 오는가?
부모 RDD는 무엇인가?
몇 개의 Partition으로 나뉘는가?
Partition 하나를 어떻게 계산하는가?예를 들어 filter()나 map()을 호출하면 기존 RDD를 수정하는 것이 아니라 부모 RDD를 참조하는 새로운 MapPartitionsRDD가 만들어짐.
source RDD
│
│ filter
▼
filtered RDD
│
│ map
▼
parsed RDD따라서 RDD는 실제 데이터 자체이면서 동시에:
필요한 데이터가 아직 계산되지 않았다면 어떻게 만들어야 하는지를 설명하는 계산 설계도
라고 볼 수 있음.
3.3 How an RDD Is Implemented
Spark Core에서 RDD는 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들로 구성되는지 반환함.
RDD
├─ Partition 0
├─ Partition 1
└─ Partition 2Partition은 실제 데이터 전체를 담는 객체라기보다, RDD 안의 특정 데이터 조각을 식별하기 위한 논리적 정보에 가까움.
예를 들어 파일 기반 RDD의 Partition은 다음 정보를 가질 수 있음.
File Partition
├─ File Path
├─ Start Offset
└─ Lengthcompute()
특정 Partition을 실제로 계산하는 방법을 정의함.
def compute(
partition: Partition,
context: TaskContext
): Iterator[T]개념적으로 map()이 만든 RDD의 compute()는 다음과 유사함.
override def compute(
split: Partition,
context: TaskContext
): Iterator[U] = {
parent.iterator(split, context).map(f)
}즉:
부모 RDD의 해당 Partition 읽기
→ map 함수 적용
→ 결과 Iterator 반환반환값이 전체 List나 Array가 아니라 Iterator라는 점이 중요함.
Partition 전체를 중간 컬렉션으로 생성
→ 다음 연산에 전달하는 대신,
Record 1 읽기 → map → filter
Record 2 읽기 → map → filter
Record 3 읽기 → map → filter처럼 데이터를 한 레코드씩 흘려보낼 수 있음.
getDependencies
이 RDD가 어떤 부모 RDD에 의존하는지 표현함.
Child RDD
↓ depends on
Parent RDDSpark는 이 Dependency 정보를 바탕으로 RDD DAG를 만들고, Shuffle 경계를 찾아 Stage를 구분함.
partitioner
Key–Value RDD에서 각 Key가 어느 Partition으로 들어갈지 결정하는 규칙임.
대표적으로:
HashPartitioner
RangePartitioner등이 있음.
getPreferredLocations
특정 Partition을 어느 노드에서 계산하는 것이 유리한지 알려줌.
Data stored on Node A
↓
Prefer Task on Node A데이터가 위치한 곳에 가까운 Executor에서 Task를 실행하면 네트워크를 통한 데이터 이동을 줄일 수 있음.
3.4 Resilient: Recovering from Failures
RDD의 Resilient는 데이터를 무조건 여러 개 복제해서 보관한다는 뜻이 아님.
RDD는 데이터가 만들어진 연산 관계인 Lineage를 기억해 손실된 Partition을 다시 계산할 수 있음.
source = sc.textFile("s3://logs/")
parsed = source.map(parse_log)
errors = parsed.filter(lambda x: x.level == "ERROR")source RDD
↓ map(parse_log)
parsed RDD
↓ filter(level == ERROR)
errors RDD이 관계가 Lineage임.
예를 들어 errors RDD의 Partition 2가 Executor 장애로 손실됐다고 하자.
errors RDD
Partition 0 정상
Partition 1 정상
Partition 2 손실
Partition 3 정상Spark는 Lineage를 거슬러 올라가 필요한 Partition만 다시 계산할 수 있음.
source Partition 2 다시 읽기
↓ parse_log
parsed Partition 2 재생성
↓ filter
errors Partition 2 재생성전체 RDD 재계산 ❌
손실된 Partition의 계보만 재계산 ⭕Cache가 손실된 경우
errors.cache()Cache된 Partition도 Executor가 종료되면 사라질 수 있음.
하지만 원본 데이터와 Lineage가 남아 있다면 다시 계산할 수 있음.
Cache
→ 빠른 재사용을 위한 임시 저장
→ 유실될 수 있음
Lineage
→ 유실된 Partition을 다시 만드는 방법Lineage가 너무 길어지는 경우
Transformation이 너무 많이 연결되면 재계산 비용과 스케줄링 정보가 커질 수 있음.
RDD 1
↓
RDD 2
↓
RDD 3
↓
...
↓
RDD 1000이럴 때 checkpoint()를 사용해 데이터를 안정적인 저장소에 실제로 기록하고 이전 Lineage를 끊을 수 있음.
Long Lineage
↓
Checkpoint
↓
New Lineage StartCache
→ 빠른 재사용이 목적
→ Lineage는 유지됨
Checkpoint
→ 안정적인 저장소에 실제 데이터 기록
→ 긴 Lineage를 단절3.5 Distributed: Partitioning Data Across the Cluster
RDD는 하나의 거대한 데이터 묶음으로 처리되지 않고 여러 Partition으로 나뉨.
RDD
├─ Partition 0
├─ Partition 1
├─ Partition 2
└─ Partition 3Partition은 Spark의 병렬 처리 수준을 결정하는 핵심 단위
Partition 0 → Task 0
Partition 1 → Task 1
Partition 2 → Task 2
Partition 3 → Task 3Partition과 저장 위치
RDD Partition이 항상 특정 Executor에 고정되어 있는 것은 아님.
Partition
= 논리적 데이터 조각
Executor
= Partition을 실제로 계산하는 프로세스Task가 실행될 때 Scheduler가 Partition 계산을 특정 Executor에 할당함.
Cache한 경우에는 계산된 Partition 데이터가 특정 Executor의 BlockManager에 저장될 수 있음.
Before Computation
RDD Partition 0
→ 계산 방법만 존재After cache()
Executor A
└─ Cached Block: RDD 3, Partition 0Partition은 어떻게 만들어지는가?
- 외부 파일의 Split
parallelize()에서 지정한 Slice 수- 부모 RDD의 Partition 상속
- Shuffle 결과
repartition()coalesce()
예시:
rdd = sc.parallelize(data, 4)4 Partitionsrepartitioned = rdd.repartition(10)10 Partitions3.6 Dataset: Representing Distributed Data
RDD에서 Dataset은 Spark SQL의 Dataset 클래스와 반드시 같은 의미로 보면 안 됨.
RDD의 Dataset
→ 데이터 요소들의 집합이라는 일반적인 의미
Spark SQL Dataset API
→ 스키마와 Encoder를 가진 고수준 APIRDD에는 어떤 형태의 객체든 담을 수 있음.
numbers: RDD[int]
lines: RDD[str]
pairs: RDD[tuple[str, int]]
orders: RDD[Order]RDD[Int]
RDD[String]
RDD[(String, Int)]
RDD[Order]이 유연성 덕분에 비정형 데이터나 사용자 정의 객체를 자유롭게 처리할 수 있음.
3.7 Immutability and Transformations
RDD는 Immutable, 즉 생성된 뒤 그 자체가 수정되지 않음.
numbers = sc.parallelize([1, 2, 3])
doubled = numbers.map(lambda x: x * 2)이 코드에서 numbers의 값이 바뀌는 것이 아님.
numbers RDD
[1, 2, 3]
doubled RDD
[2, 4, 6]map()은 기존 RDD를 수정하는 대신, 기존 RDD에 의존하는 새로운 RDD를 반환함.
numbers RDD
│
│ map(x × 2)
▼
doubled RDD다음 Transformation도 동일함.
filtered = doubled.filter(lambda x: x > 2)numbers RDD
↓ map
doubled RDD
↓ filter
filtered RDD불변성이 필요한 이유
분산 환경에서 여러 Executor가 동일한 데이터를 직접 수정하도록 허용하면 다음 문제가 생길 수 있음.
Executor A가 데이터 수정
Executor B가 같은 데이터 수정
Executor C는 수정 전 데이터를 읽음- 동시성 제어
- Lock
- 데이터 버전 관리
- 장애 복구
- 변경 순서 관리
가 복잡해짐.
RDD는 기존 데이터 수정 대신 새 데이터셋을 만드는 모델을 사용함.
Input RDD
→ Transformation
→ New RDD따라서 데이터가 어떻게 변했는지 Lineage로 명확하게 표현할 수 있음.
RDD A
↓ Transformation 1
RDD B
↓ Transformation 2
RDD C불변성과 Lineage는 서로 연결됨.
Immutability
→ 이전 RDD가 변하지 않음
→ Transformation 관계를 안정적으로 추적 가능
→ 손실된 Partition을 재계산 가능초기 RDD 논문도 RDD를 immutable하고 partitioned된 레코드 컬렉션으로 정의함. (MIT CSAIL)
3.8 Transformations, Actions, and Lazy Evaluation
RDD 연산은 크게 Transformation과 Action으로 나뉨.
Transformation
기존 RDD로부터 새로운 RDD를 만드는 연산.
mapped = rdd.map(...)
filtered = mapped.filter(...)
pairs = filtered.flatMap(...)대표적인 Transformation:
mapfilterflatMapuniondistinctreduceByKeygroupByKeyjoinrepartition
RDD A
↓ Transformation
RDD BTransformation을 호출해도 바로 전체 데이터가 계산되지 않음.
Action
RDD의 실제 결과를 요구하는 연산임.
rdd.count()
rdd.collect()
rdd.take(10)
rdd.reduce(...)
rdd.saveAsTextFile(...)대표적인 Action:
countcollecttakefirstreducesaveAsTextFileforeach
RDD Lineage
↓ Action
Spark Job
↓
Actual Computation예시:
result = (
sc.textFile("s3://logs/")
.filter(lambda x: "ERROR" in x)
.map(parse_log)
)여기까지는 Transformation만 존재함.
No Job Yet
Text File
↓ filter
↓ map
Result RDD이후:
result.count()가 호출되면 실제 Job이 시작됨.
Action: count()
↓
Job Created
↓
Tasks Sent to Executors왜 지연 실행하는가?
1. 여러 연산을 파이프라인으로 실행
rdd.map(f).filter(g).map(h)Record 1 → f → g → h
Record 2 → f → g → h
Record 3 → f → g → h각 연산 결과 전체를 별도로 저장할 필요가 없음.
2. 필요하지 않은 데이터는 계산하지 않을 수 있음
rdd.take(10)전체 데이터가 아니라 결과 10개를 얻는 데 필요한 Partition만 먼저 계산할 수 있음.
3. 실행 전에 전체 의존 관계를 확인
Spark는 최종 RDD에서 부모 RDD 방향으로 Dependency를 따라가며 어떤 작업이 필요한지 파악함.
Final RDD
↓ parent
RDD C
↓ parent
RDD B
↓ parent
RDD A다만 RDD의 Lazy Evaluation을 DataFrame의 Catalyst Optimization과 완전히 동일하게 보면 안 됨.
RDD Lazy Evaluation
→ 연산 연결, 필요한 계산 결정, Stage 구성의 기반
DataFrame Lazy Evaluation
→ 위 기능에 더해 SQL 의미를 이용한 계획 최적화 가능RDD 내부의 임의 Python·Scala 함수는 Spark가 의미를 분석하기 어렵기 때문에, 컬럼 가지치기나 조인 재배치 같은 고수준 최적화에는 제한이 있음.
3.9 RDD Lineage
RDD Lineage는 RDD가 어떤 Source와 Transformation을 거쳐 만들어졌는지를 나타냄.
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는 다음과 같음.
TextFile RDD
↓ filter
MapPartitionsRDD
↓ map
MapPartitionsRDD
↓ reduceByKey
ShuffledRDD각 RDD는 직접 모든 이전 연산 문자열을 저장하는 단순한 기록이 아니라, 부모 RDD에 대한 Dependency를 가지고 있음.
RDD C
└─ Dependency on RDD B
RDD B
└─ Dependency on RDD A이 연결을 따라가면 전체 Lineage가 만들어짐.
Lineage의 역할
- 실행할 연산 관계 표현
- Stage를 나누기 위한 Dependency 정보 제공
- 장애 발생 시 손실된 Partition 재계산
- Cache가 없는 RDD를 필요할 때 다시 계산
Lineage
├─ Execution Planning
├─ Stage Construction
└─ Fault RecoveryLineage와 DAG
Lineage는 일반적으로 DAG 형태로 구성됨.
RDD A
/ \
/ \
RDD B RDD C
\ /
\ /
RDD DDAG는 Directed Acyclic Graph, 즉 방향은 있지만 순환하지 않는 그래프임.
Directed
→ 부모 RDD에서 자식 RDD 방향으로 연산 관계가 존재
Acyclic
→ RDD D가 다시 자기 부모로 돌아가는 순환 관계가 없음
Graph
→ 여러 RDD와 Dependency가 노드와 간선으로 표현됨RDD는 불변이고 Transformation은 새 RDD를 생성하기 때문에 DAG 관계를 안정적으로 구성할 수 있음.
3.10 Building a DAG from RDD Operations
다음 코드를 보자.
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 관계:
orders
↓ map(parse_order)
parsed_orders
↓ filter(status == PAID)
paid_orders
↓ map(region, amount)
region_amounts
↓ reduceByKey
totalsDAG로 표현하면:
TextFileRDD
↓
MapPartitionsRDD
↓
MapPartitionsRDD
↓
MapPartitionsRDD
↓
ShuffledRDDSpark는 Action인 collect()가 호출되면 최종 RDD인 totals부터 부모 RDD 방향으로 Dependency를 추적함.
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 종류에 따라 서로 다른 구현체가 사용됨.
RDD
├─ ParallelCollectionRDD
├─ HadoopRDD
├─ MapPartitionsRDD
├─ ShuffledRDD
├─ UnionRDD
└─ CoalescedRDDParallelCollectionRDD
Driver의 로컬 컬렉션을 여러 Partition으로 나눌 때 사용됨.
rdd = sc.parallelize([1, 2, 3, 4], 2)ParallelCollectionRDD
├─ Partition 0: [1, 2]
└─ Partition 1: [3, 4]HadoopRDD
HDFS, S3 등 Hadoop InputFormat과 호환되는 데이터 소스를 읽을 때 사용됨.
rdd = sc.textFile("s3://bucket/logs/")HadoopRDD
├─ File Split 0
├─ File Split 1
└─ File Split 2MapPartitionsRDD
map, filter, flatMap 같은 다수의 Transformation에서 사용됨.
mapped = rdd.map(lambda x: x * 2)Parent RDD
↓ map
MapPartitionsRDDShuffledRDD
Key별 재분배가 필요한 Shuffle 결과를 표현함.
result = pairs.reduceByKey(lambda a, b: a + b)Parent Partitions
↓ Shuffle
ShuffledRDDUnionRDD
여러 RDD를 하나로 합침.
combined = rdd1.union(rdd2)RDD 1 ─┐
├→ UnionRDD
RDD 2 ─┘이처럼 RDD는 하나의 고정된 데이터 컨테이너가 아니라:
서로 다른 데이터 소스와 연산을 동일한 Partition·Dependency·
compute()모델로 실행할 수 있게 만드는 공통 추상 클래스
라고 볼 수 있음.
3.12 The Relationship Between RDDs and DataFrames
현재 Spark 사용자는 대부분 RDD보다 DataFrame과 SQL을 사용함.
orders = spark.read.parquet("s3://orders/")
result = (
orders
.filter("status = 'PAID'")
.groupBy("region")
.sum("amount")
)DataFrame은 Spark에 다음 정보를 제공함.
Schema
├─ Column Names
├─ Data Types
└─ Nullable Information
Operations
├─ Filter Expressions
├─ Join Conditions
├─ Aggregations
└─ ProjectionsSpark는 이를 Logical Plan으로 표현하고 Catalyst Optimizer로 최적화함.
DataFrame / SQL
↓
Logical Plan
↓
Optimized Logical Plan
↓
Physical Plan
↓
Distributed ExecutionRDD에서는 사용자 함수가 블랙박스에 가까움.
rdd.filter(lambda x: custom_condition(x))Spark
→ 함수가 호출된다는 사실은 알 수 있음
→ 함수 내부의 의미는 분석하기 어려움DataFrame에서는 필터가 명시적인 Expression임.
df.filter(df.status == "PAID")Spark
→ status 컬럼을 사용함
→ PAID 조건으로 필터링함
→ Predicate Pushdown 가능 여부 판단RDD와 DataFrame 비교
| 구분 | RDD | DataFrame |
|---|---|---|
| 데이터 표현 | 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)
DataFrame / SQL
↓
Query Planning and Optimization
↓
Spark Physical Operators
↓
Distributed Partition Processing
↓
Tasks on Executors따라서 다음처럼 정리할 수 있음.
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 재계산"]