전체 그래프
Spark

RDD

data-engineeringsparkrddresilient-distributed-dataset

상위: Spark

요약

RDD(Resilient Distributed Dataset)는 Spark의 핵심 데이터 구조로, 여러 노드에 분산 저장되고 장애 시 자동 복구되는 불변 데이터셋입니다. 파티션 단위로 병렬 처리되며, Lazy Evaluation을 통해 성능을 최적화합니다.

RDD란?

  • R: Resilient (탄력적): 장애 발생 시 자동 복구
  • D: Distributed (분산): 여러 노드에 분산 저장
  • D: Dataset (데이터셋): 객체들의 모음 (item, record 단위)
  • → 데이터를 파티셔닝(partitioning)하여 병렬 처리 가능

RDD의 특징

1. 데이터의 추상화 (Data Abstraction)

  • 데이터는 여러 Partition으로 나누어 저장됨
  • 사용자는 전체 데이터처럼 사용하지만, 내부적으로는 분산 저장됨

RDD

연산 단계:

  • 연산 전: 논리적 데이터
  • 실행 시: 물리적으로 분산 처리
  • 저장 시: 하나의 파일로 통합

2. 탄력성 (Resilient) & 불변성 (Immutable)

  • RDD는 한 번 생성되면 변경 불가
  • 장애 발생 시 **의존성 정보(Lineage)**를 통해 자동 복구 가능

RDD

장애 복구 방식:

  1. RDD의 변환 과정을 Lineage로 기록
  2. 노드 장애 발생
  3. Lineage를 재실행하여 RDD 재생성

3. Type-safe 기능

  • RDD는 정해진 타입(예: Integer RDD, String RDD 등)의 객체를 가짐

  • 컴파일 시점에 타입 검사 → 성능 최적화유지보수성 향상

    ⇒ PySpark는 동적 타입 언어이므로 Type-safe 미지원

RDD

4. 정형 & 비정형 데이터 처리

  • Unstructured 데이터: 로그, 자연어 → sc.textFile() 사용
  • Structured 데이터: 테이블 → RDD.map() 또는 DataFrame 활용

RDD

5. 지연 평가 (Lazy Evaluation)

  • 중간 연산을 저장만 하고, 최종 Action이 일어나기 전까지 실행 안 함
  • 연산을 모아서 최적화된 실행 계획을 수립
  • 불필요한 연산 방지 → 성능 최적화 & 리소스 절약

RDD

실행 흐름:

  1. Transformation 연산 호출 → Lineage에 기록만
  2. 또 다른 Transformation 호출 → 계속 기록
  3. Action 호출 → 이때 모든 Transformation 실행

RDD 구성 요소

1. 의존성 정보 (Lineage)

  • 어떤 입력을 기반으로 어떻게 RDD가 생성되었는지 추적
  • 장애 발생 시 해당 정보를 기반으로 RDD 재생성 가능

Lineage 예시:

# textFile  map  filter  count
#  단계의 의존성이 기록됨
rdd1 = sc.textFile("data.txt")        # Lineage: textFile
rdd2 = rdd1.map(lambda x: x.upper())  # Lineage: textFile  map
rdd3 = rdd2.filter(lambda x: "SPARK" in x)  # Lineage: textFile  map  filter

2. 파티션 정보 (Partition)

  • 실행기(Executor)에 파티션 단위로 분산 → 병렬 연산 가능
  • 각 파티션은 독립적으로 처리됨
  • 파티션 수 = 병렬 처리 수준

3. 연산 함수

  • Partition Iterator[T] 형태
  • RDD 내 데이터는 반복자 형태로 전달되어 처리
  • 각 파티션마다 함수가 독립적으로 적용됨

RDD 생성 방법

1. 기존 메모리 데이터를 RDD로 변환 (parallelize)

Python:

wordsRDD = sc.parallelize(["fish", "cats", "dogs"])

Scala:

val wordsRDD = sc.parallelize(List("fish", "cats", "dogs"))

Java:

JavaRDD<String> wordsRDD = sc.parallelize(Arrays.asList("fish", "cats", "dogs"))
  • 클러스터로 데이터를 보낼 때 사용
  • 비효율적일 수 있어 소규모 분석에 적합

2. 외부 파일에서 RDD 생성 (textFile)

Python:

linesRDD = sc.textFile("/path/to/README.md")

Scala:

val linesRDD = sc.textFile("/path/to/README.md")

Java:

JavaRDD<String> linesRDD = sc.textFile("/path/to/README.md")
  • 실무에서는 보통 파일이나 DB에서 불러옴
  • 일반 텍스트(CSV, 로그 등), S3, HBase, Cassandra 등 다양한 소스 지원

RDD vs DataFrame

특징RDDDataFrame
스키마없음있음 (컬럼명, 타입)
최적화수동 최적화 필요Catalyst Optimizer 자동 최적화
타입 안정성컴파일 타임 (Scala/Java)런타임 (분석), 컴파일 타임 (문법)
언어 지원모든 언어모든 언어
사용 편의성저수준 API고수준 API, SQL 지원
성능직접 제어 가능하지만 최적화 어려움자동 최적화로 빠른 성능

RDD 사용 케이스

RDD를 사용하는 경우:

  1. Transformation / Action을 직접 제어해야 할 때
  2. 구조화되지 않은 스트림 데이터 처리 시
  3. 함수형 프로그래밍이 필요한 도메인일 때
  4. 스키마 변화가 필요 없는 경우
  5. DataFrame/Dataset이 성능 한계일 때

DataFrame을 사용하는 경우:

  1. 도메인 기반 고수준 API가 필요할 때
  2. SQL 쿼리를 사용하고 싶을 때
  3. 타입 안정성과 최적화된 실행 계획 필요할 때
  4. 일관성 있는 Spark API 사용을 원할 때