MapReduce와 Spark RDD로 설계하는 분산 병렬 데이터 처리
MapReduce와 Spark RDD의 실행 모델, 셔플과 파티셔닝, 장애 복구 및 분산 병렬 처리 튜닝 방식을 실무 관점에서 정리한다.
2026-08-14 · 최초 발행 2024-04-29
대용량 처리는 파티션과 셔플에서 성패가 갈린다
분산 처리 시스템은 데이터를 여러 노드에 나눠 배치하고 태스크를 병렬 실행해 확장성을 확보한다. 이때 처리량만큼 중요한 것이 데이터 로컬리티, 셔플 비용, 장애 복구와 출력 일관성이다.
MapReduce는 키-값 기반의 분산 처리 모델이자 런타임이다. Map 단계에서 레코드를 변환하고, Shuffle/Sort를 거쳐 Reduce 단계에서 집계한다. HDFS나 오브젝트 스토리지의 대규모 배치를 디스크 중심으로 다루며, Combiner·Partitioner·Speculative Execution 같은 안정성 중심의 설계를 활용한다.
Spark RDD(Resilient Distributed Dataset)는 불변이며 파티션된 컬렉션 추상화다. Transformations와 Actions를 기준으로 지연 실행하고 DAG 스케줄링을 수행한다. 메모리 중심 처리와 캐시, Lineage 재계산을 이용하므로 반복적이거나 대화형인 워크로드에서 성능을 낼 수 있다.
분산 병렬 처리에서는 데이터 파티셔닝과 태스크 병렬화로 수평 확장을 구현한다. 데이터 로컬리티와 셔플 최적화, 장애 복구, 일관성 보장을 어느 한쪽으로 치우치지 않게 설계해야 한다.
실행 모델이 만드는 처리 경로
입력 데이터는 스플릿 또는 파티션으로 나뉜다. 노드 로컬과 랙 로컬 배치를 우선하면 네트워크 비용을 줄일 수 있다. 해시 또는 레인지 파티셔너의 선택, 키 스큐 완화는 셔플 병목을 줄이는 핵심이다.
MapReduce는 Job에서 Map과 Reduce로 이어지는 단계형 파이프라인이며 잡 단위 배리어가 존재한다. Spark는 DAG를 스테이지와 태스크로 나누어 스케줄링하고, 파이프라이닝과 넓은 의존성(Shuffle) 경계를 활용해 병렬 실행을 최적화한다.
장애 대응 방식도 다르다. MapReduce는 HDFS 복제, 태스크 재시도·추정 실행(Speculative Execution), 출력의 원자적 커밋을 사용한다. Spark는 Lineage에 따라 데이터를 재계산하고 체크포인팅·캐시, 태스크 재시도와 셔플 재생성, 출력 커밋 프로토콜을 조합한다.
셔플 구간에서는 Map-side Combine, 압축, 정렬기(External Sort)로 디스크와 네트워크 부담을 낮춘다. Spark는 Tungsten/Project Hydra(최적화 엔진, 최신 정보 확인 필요), 셔플 파일 병합, AQE(Adaptive Query Execution, DataFrame)로 런타임 최적화를 수행한다. CPU·메모리·디스크·네트워크 자원은 YARN/Mesos/Kubernetes로 관리하며, 병렬도·파티션 수·GC·압축 코덱·I/O 블록 크기·스토리지 계층을 조정해 SLO를 맞춘다.
배치 안정성과 반복 연산의 선택지
| 항목 | MapReduce | Spark RDD |
|---|---|---|
| 성능 | 디스크 중심, 반복 연산 느림, 대규모 셔플 안정 | 메모리 중심, 반복·대화형 고성능 |
| 확장성 | 페타바이트급 검증, 단순 모델 | 스케일아웃 용이, 셔플 병목 관리 필요 |
| 일관성 | 쓰기-일회성, 원자적 커밋, 강한 잡 경계 | 라인리지 일관성, 체크포인트로 보강 |
| 안정성 | 높은 내고장성(HDFS 복제, 재시도) | 재계산 기반, 셔플 장애 영향도 존재 |
| 운영 편의 | 구성 단순, 개발 장황, 튜닝 포인트 적음 | 개발 생산성 높음, 풍부한 API·툴 |
로그 ETL과 집계 배치에서는 원천 로그 압축 해제, 정규화, 세션 집계, 파티션 저장의 흐름을 구성할 수 있다. 이때 대량 셔플을 줄이기 위해 맵사이드 프리-집계를 적용한다.
사용자 행동 카운팅과 세그먼트 생성에서는 Top-N 키 집중 문제가 발생할 수 있다. 키 솔팅과 커스텀 파티셔닝으로 키 스큐를 완화한다. 머신러닝 전처리나 반복 연산에서는 RDD 캐시와 영속화를 사용해 반복 피처 엔지니어링을 가속하고, 조인 순서와 브로드캐스트로 대소조인을 최적화한다.
PageRank 등의 그래프 처리에서는 Spark의 반복적 메시지 패싱에 캐시와 체크포인트를 혼용할 수 있다. MapReduce는 배치 반복 방식으로 안정적으로 처리한다. 사실·차원 테이블을 조인하는 대규모 집계형 DWH 오프로드에서는 브로드캐스트 조인, 파티션 프루닝, Parquet 파일 포맷으로 I/O를 최적화한다.
설계와 운영에서 확인할 지점
입력은 수십~수백 MB 단위로 스플릿 가능한 파일을 권장하며, 스몰 파일은 사전에 Compaction으로 합친다. 파티셔닝은 키 분포와 시간축 조회 패턴을 고려해 해시와 레인지 중 선택한다. 핫 파티션이 생기면 Salting이나 스플릿 키 재설계가 필요하다.
파티션 수는 코어 수×2~3을 권장하되, 워크로드와 I/O 비율에 따라 조정한다. 셔플에는 Map-side Combine/Pre-aggregate, Snappy/LZ4 압축 코덱, 정렬 버퍼와 셔플 파일 병합을 활용한다. Spark에서는 executor 메모리·오버헤드·스토리지/실행 메모리 비율을 조정하고, GC 로그를 분석해 Young/Old GC를 튜닝한다.
출력은 작업별 임시 디렉터리에 기록한 뒤 원자적 rename으로 커밋한다. 동일 입력을 재실행해도 동일한 출력이 나오도록 멱등성을 확보한다. 느린 태스크에는 스펙ulative 실행을 활성화하고, RDD 체크포인트로 Lineage가 과도하게 깊어지는 일을 막는다. 상위 N 키를 샘플링해 분포를 확인하고, 유효성 검증과 오염 레코드 격리(에러 레이크/사이드카 파일)도 함께 구성한다.
PySpark RDD 워드카운트
전제조건: Spark 3.3+, Python 3.9+, 클러스터/YARN/K8s 중 하나
실행:
- 로컬: spark-submit wordcount.py input_dir output_dir
- YARN: --master yarn --deploy-mode cluster 권장
# wordcount.py
from pyspark import SparkConf, SparkContext
import sys
if __name__ == "__main__":
if len(sys.argv) != 3:
print("Usage: wordcount.py <input> <output>", file=sys.stderr)
sys.exit(1)
conf = SparkConf().setAppName("RDDWordCount")
sc = SparkContext(conf=conf)
input_path, output_path = sys.argv[1], sys.argv[2]
rdd = (sc.textFile(input_path)
.repartition(200) # 병렬도 조정
.flatMap(lambda line: line.split())
.map(lambda w: (w, 1))
.reduceByKey(lambda a, b: a + b)) # 맵사이드 프리-집계 포함
# 스큐 키 솔팅(예: 상위 키에 접두어 추가). 필요 시 활성화.
# salted = (rdd.map(lambda kv: (str(hash(kv[0]) % 10) + "_" + kv[0], kv[1]))
# .reduceByKey(lambda a, b: a + b)
# .map(lambda kv: (kv[0].split("_", 1)[1], kv[1]))
# .reduceByKey(lambda a, b: a + b))
rdd.saveAsTextFile(output_path)
sc.stop()
repartition으로 병렬도를 조정하며, reduceByKey에는 맵사이드 집계가 포함되어 셔플 데이터를 줄인다. 스큐가 발생하면 키 솔팅이나 커스텀 파티셔너를 적용한다.
Spark RDD는 반복성·대화형 워크로드에서 MapReduce 대비 310배, 캐시 활용 시 최대 수십 배 성능 향상이 가능하다. 데이터와 클러스터 특성에 의존한다. 셔플 최적화와 브로드캐스트 조인으로 네트워크 I/O 3070% 절감 달성 사례도 다수다.
자원 튜닝과 캐시 전략을 전제로 동일 SLA 기준 노드 수를 20~50% 절감할 수 있다. 멱등 파이프라인과 자동 재시도를 구성하면 재처리 시간을 줄이고 야간 배치 안정성을 높일 수 있다.
대규모 배치의 단순성과 보수적 안정성이 우선이면 MapReduce를 유지할 수 있다. 반복 분석과 ETL, 대규모 집계에서는 Spark RDD를 우선 검토하되, 파티셔닝 설계, 셔플 최소화, 멱등 출력 커밋, 장애와 스큐 대응을 순서대로 갖춰야 한다.