Hadoop / Spark
정의
Hadoop 은 2006 야후에서 시작된 분산 저장 + 분산 처리 오픈소스 플랫폼. Google 이 2003-2004 년에 발표한 GFS (분산 파일시스템) 와 MapReduce (분산 처리 모델) 논문을 오픈소스로 구현한 것이 출발점.
Spark 는 2009 UC Berkeley AMPLab 에서 시작된 인메모리 분산 처리 엔진. MapReduce 의 디스크 I/O 병목을 해결하기 위해 RDD (Resilient Distributed Dataset) 개념을 도입해 반복 계산에서 10-100 배 빠른 성능 을 제공.
두 시스템 모두 단일 서버로 처리 불가능한 규모의 데이터 를 여러 서버로 나눠 처리하는 것이 목적. 오늘날 데이터 레이크 / 데이터 웨어하우스 위의 대규모 처리 엔진은 대부분 이 두 계보를 잇는다.
Hadoop 의 두 축
Hadoop 은 원래 저장 (HDFS) + 처리 (MapReduce) 두 컴포넌트로 시작. 이후 YARN (자원 관리) 이 분리되어 3 요소가 됨.
flowchart TB
subgraph Hadoop
subgraph Storage["저장"]
HDFS["HDFS<br/>(분산 파일시스템)"]
end
subgraph Resource["자원 관리"]
YARN["YARN<br/>(스케줄러)"]
end
subgraph Compute["처리"]
MR["MapReduce"]
OtherEngines["Spark, Hive, Flink, Tez..."]
end
end
Compute --> YARN
Compute --> HDFS
YARN --> Storage
HDFS (Hadoop Distributed File System)
파일을 고정 크기 블록 (기본 128 MB) 으로 나눠 여러 노드에 분산 저장 + 복제.
flowchart TB
NN["NameNode<br/>(메타데이터: 파일→블록→노드 매핑)"]
DN1["DataNode 1"]
DN2["DataNode 2"]
DN3["DataNode 3"]
DN4["DataNode 4"]
NN -.->|heartbeat| DN1
NN -.->|heartbeat| DN2
NN -.->|heartbeat| DN3
NN -.->|heartbeat| DN4
File["file.parquet 1 GB"] --> B1["Block 1 (128MB)"]
File --> B2["Block 2"]
File --> B3["Block 3"]
File --> B8["... Block 8"]
B1 -.->|"replica 3"| DN1
B1 -.-> DN2
B1 -.-> DN3
B2 -.-> DN2
B2 -.-> DN3
B2 -.-> DN4
핵심 특성:
- 큰 블록 (128 MB 기본): 파일당 블록 수 감소 -> NameNode 부담 감소, sequential I/O 최적
- 복제 계수 3 (기본): 노드/랙/데이터센터 장애에 대비
- Write-Once-Read-Many: append 만 가능, 임의 수정 불가 (log-structured 사고)
- Rack awareness: 복제본을 다른 랙에 분산 (물리 장애 격리)
한계:
- NameNode 는 단일 지점 병목 (HA 로 완화)
- 작은 파일 다수 는 부적합 (파일당 블록 오버헤드)
- 저지연 요구 (< 100 ms) 부적합 (배치 지향)
클라우드 시대에는 S3 같은 객체 저장소가 사실상 HDFS 를 대체 (EMR 은 EMRFS 로 S3 를 HDFS 처럼 접근).
MapReduce
큰 데이터를 작은 조각으로 나눠 병렬 처리 -> 결과를 합치는 프로그래밍 모델. Google 이 2004 년에 발표.
flowchart LR
Input[("입력 데이터<br/>여러 블록")]
subgraph Map["Map 단계 (병렬)"]
M1["Mapper 1<br/>(k,v) 방출"]
M2["Mapper 2"]
M3["Mapper N"]
end
subgraph Shuffle["Shuffle 단계<br/>(네트워크 재분배)"]
S["같은 key 는 같은 reducer 로"]
end
subgraph Reduce["Reduce 단계 (병렬)"]
R1["Reducer 1<br/>집계"]
R2["Reducer 2"]
end
Output[("결과")]
Input --> M1
Input --> M2
Input --> M3
M1 --> S
M2 --> S
M3 --> S
S --> R1
S --> R2
R1 --> Output
R2 --> Output
예: WordCount
Map: "hello world" → [("hello", 1), ("world", 1)]
Shuffle: 같은 word 를 같은 reducer 로 보냄
Reduce: [("hello", 1), ("hello", 1), ...] → ("hello", 42)
MapReduce 의 한계 (Spark 등장의 이유):
- 단계마다 디스크 write (Map 결과 -> disk -> Shuffle -> disk -> Reduce)
- 반복 알고리즘 (ML 학습, PageRank) 이 매우 느림 (매 iteration 마다 디스크 왕복)
- 표현력 제한 (map + reduce 로 표현 어려운 알고리즘 다수)
YARN (Yet Another Resource Negotiator)
컴퓨트 자원 (CPU, 메모리) 을 여러 애플리케이션에 분배 하는 스케줄러. Hadoop 2.x 부터 MapReduce 로부터 분리.
구성:
- ResourceManager (마스터): 전체 자원 스케줄링
- NodeManager (워커): 각 노드의 컨테이너 관리
- ApplicationMaster (앱별): 애플리케이션의 태스크 조율
YARN 위에서 MapReduce, Spark, Flink, Hive 등이 함께 실행 가능 -> Hadoop 이 단일 처리 엔진이 아닌 플랫폼 이 됨.
Spark 의 등장
“디스크 I/O 를 줄이자” 가 핵심 통찰.
RDD (Resilient Distributed Dataset)
Spark 의 원류 추상. 메모리에 분산 저장된 불변 컬렉션.
“Resilient” 의 의미: 데이터 자체를 복제해서 저장하지 않고, 데이터를 만드는 계보 (lineage) 만 저장. 노드 장애 시 lineage 를 다시 실행해 복구.
raw = sc.textFile("s3://.../logs") # RDD 1
words = raw.flatMap(lambda l: l.split()) # RDD 2 (lineage: raw → flatMap)
counts = words.map(lambda w: (w, 1)) # RDD 3
result = counts.reduceByKey(lambda a, b: a+b) # RDD 4
result.saveAsTextFile("s3://.../out") # Action
- Transformation (map, filter, join): 지연 실행 (lazy), lineage 만 기록
- Action (count, collect, save): 실제 실행 트리거
DAG 스케줄러
Transformation 을 DAG (Directed Acyclic Graph) 로 표현 -> stage (shuffle 경계) 로 나눔 -> task 로 분해 -> executor 에서 병렬 실행.
flowchart LR
subgraph Stage1["Stage 1 (narrow deps)"]
T1[Read] --> T2[Filter] --> T3[Map]
end
subgraph Boundary["Shuffle 경계"]
S["ShuffleWrite/Read"]
end
subgraph Stage2["Stage 2 (narrow deps)"]
T4[Combine] --> T5[Output]
end
Stage1 --> S --> Stage2
- Narrow dependency (map, filter): 파티션 1 → 파티션 1 (shuffle X)
- Wide dependency (groupBy, join): 파티션 다수 → 파티션 재분배 (shuffle O)
Shuffle 이 성능 병목이므로 shuffle 을 최소화 하는 게 Spark 튜닝의 핵심.
DataFrame / Dataset (2015+)
RDD 는 저수준. Catalyst Optimizer 가 최적화할 수 있도록 스키마와 SQL 문법 을 얹은 것이 DataFrame / Dataset.
# RDD (저수준, 최적화 여지 없음)
raw.filter(lambda r: r[2] > 100).map(lambda r: (r[0], r[3]))
# DataFrame (Catalyst 가 파티션 프루닝, 컬럼 프루닝, predicate pushdown 자동)
df.select("id", "amount").filter("amount > 100")
# SQL (같은 최적화 자동)
spark.sql("SELECT id, amount FROM t WHERE amount > 100")
Catalyst 가 자동으로:
- Predicate pushdown (Parquet statistics 활용)
- Column pruning
- Constant folding
- Join reordering
- Whole-stage code generation (Java bytecode 생성)
현대에는 DataFrame API + SQL 이 사실상 표준. RDD 는 직접 사용 드묾.
아키텍처
flowchart TB
subgraph Driver["Driver 프로세스"]
Main["사용자 코드 main()"]
DAG["DAG Scheduler"]
Task["Task Scheduler"]
end
subgraph CM["Cluster Manager<br/>(YARN / K8s / Standalone)"]
RM["자원 할당"]
end
subgraph Workers["Executor JVM 들"]
E1["Executor 1<br/>(N task slot)"]
E2["Executor 2"]
E3["Executor N"]
end
Main --> DAG
DAG --> Task
Task -->|"task 제출"| CM
CM --> E1
CM --> E2
CM --> E3
E1 -.->|"결과 반환"| Driver
- Driver: 사용자 코드 실행, DAG 생성, 태스크 분배
- Executor: 실제 데이터 처리 (여러 task 동시 실행 = slot)
- Cluster Manager: 자원 할당 (YARN / EMR / K8s / Standalone)
Spark 컴포넌트
Spark 는 하나의 엔진 위에 여러 상위 라이브러리 제공.
| 컴포넌트 | 용도 |
|---|---|
| Spark Core | RDD, DAG, 실행 엔진 |
| Spark SQL | DataFrame + SQL 쿼리 (Catalyst 최적화) |
| Structured Streaming | DataFrame API 로 스트리밍 (micro-batch or continuous) |
| MLlib | 분산 머신러닝 (일부 알고리즘) |
| GraphX | 그래프 분석 (거의 사용 안 됨, GraphFrames 로 대체) |
MapReduce vs Spark
| 축 | MapReduce | Spark |
|---|---|---|
| 저장 | 단계마다 디스크 | 메모리 우선 (부족 시 디스크 spill) |
| 모델 | Map + Reduce 2 단계 | 임의 DAG (map, filter, join, groupBy, …) |
| 속도 | 기준 | 10-100 배 (반복 알고리즘) |
| API | Java 위주 | Scala, Python, R, SQL, Java |
| 상호작용 | 배치 only | REPL / 노트북 (Zeppelin, Jupyter) |
| 스트리밍 | X (Storm 별도) | Structured Streaming |
| ML | Mahout 별도 | MLlib 내장 |
| 표현력 | 낮음 | 높음 (SQL, DataFrame, MLlib) |
| 디버깅 | 어려움 | Spark UI |
결론: 새 프로젝트는 사실상 Spark. MapReduce 는 레거시 유지만.
관리형 Spark / Hadoop
| 서비스 | 특징 |
|---|---|
| Amazon EMR | Hadoop 생태계 완전 관리형, EC2/Serverless/EKS |
| AWS Glue | 서버리스 Spark ETL |
| Databricks | Spark 창시자들의 상용, Delta Lake |
| Google Dataproc | GCP 관리형 |
| Azure HDInsight / Synapse | Azure |
| 자체 호스팅 | Cloudera, Hortonworks (인수 통합) |
언제 Spark 를 쓰나
적합:
- 수 GB ~ PB 규모 배치 처리
- ETL 파이프라인 (raw -> 정제)
- ML 특징 엔지니어링 + 학습
- 대규모 join / aggregation
- 반복 알고리즘 (K-means, GraphX, PageRank)
- 로그/이벤트 스트림 분석 (Structured Streaming)
부적합 (더 나은 대안):
- 초저지연 (< 100 ms) 처리 -> Flink / Kafka Streams
- 트랜잭션 처리 -> Postgres / OLTP DB
- 작은 데이터 (< 10 GB) -> pandas / DuckDB / Polars
- 단순 SQL 애드혹 -> Athena / Redshift
Spark 성능 튜닝의 큰 축
- 파일 포맷: Parquet + ZSTD (CSV/JSON 대비 10-100 배)
- 파티셔닝: partition 컬럼으로 프루닝
- Shuffle 최소화:
repartition남용 X, broadcast join 활용 - Broadcast join: 작은 테이블 (~ 100 MB) 을 executor 에 복제
- AQE (Adaptive Query Execution): 런타임 통계로 자동 최적화 (Spark 3.x 기본 on)
- Executor 크기: 너무 크면 GC 문제, 너무 작으면 오버헤드 (경험적 4-8 core, 16-32 GB)
- Cache / Persist: 반복 참조 데이터만
- Skew 처리: Salting, AQE skew join
함정
WARNING
collect() 를 큰 데이터에 = Driver OOM. Driver 는 단일 JVM. 큰 결과는 write / foreachPartition 로.
CAUTION
Shuffle 남용. groupBy, join, distinct, repartition 는 shuffle. Wide dependency 를 반복하면 성능 급락.
WARNING
작은 파일 (Small files) 수천 개 = task 수 폭증 + 스케줄링 오버헤드. Compaction 후 read.
IMPORTANT
UDF 는 Catalyst 최적화 우회. 가능하면 built-in 함수 사용. Python UDF 는 특히 느림 (직렬화 오버헤드).
CAUTION
Skewed key = 특정 task 만 오래 걸림 (전체가 그 task 에 blocked). AQE skew join 또는 salting.
IMPORTANT
RDD 직접 사용 지양. DataFrame / SQL 이 Catalyst 최적화 받아 훨씬 빠름. 특수한 경우만 RDD.
관련 위키
- MPP (대량 병렬 처리) - 다른 계보의 분산 처리
- Data Warehouse - Spark 의 소비 대상
- Data Lake - Spark 의 저장소
- Apache Parquet - Spark 의 표준 포맷
- ETL / ELT - Spark 의 대표 워크로드
- Amazon EMR - AWS 관리형 Hadoop/Spark
- AWS Glue - 서버리스 Spark ETL
- Amazon Athena - Trino 기반 SQL 대안
- AWS S3 - HDFS 대체 저장소
이 글의 용어 (12개)
- [AWS] Amazon Redshiftcloud
- 정의 Amazon Redshift 는 AWS 가 관리하는 페타바이트 규모 컬럼형 데이터 웨어하우스 입니다. 2012년 PostgreSQL 8.0.2 를 기반으로 시작해 MPP (…
- [AWS] Athena: 서버리스 SQL on S3cloud
- 정의 Amazon Athena 는 S3 에 있는 데이터를 서버 없이 SQL 로 쿼리 하는 서비스. 클러스터 프로비저닝, 스키마 로딩, 인덱스 생성 없이 데이터 파일 (Parque…
- [AWS] EMR: 관리형 빅데이터 클러스터cloud
- 정의 Amazon EMR (Elastic MapReduce) 은 Hadoop / Spark 등 오픈소스 빅데이터 프레임워크를 관리형 클러스터로 제공 하는 AWS 서비스. AWS …
- [AWS] Glue: 서버리스 ETL + Data Catalogcloud
- 정의 AWS Glue 는 서버리스 ETL + 통합 메타데이터 카탈로그 플랫폼. Apache Spark 로 데이터 변환을 실행하고, Hive 호환 Data Catalog 로 스키마…
- [AWS] Managed Service for Apache Flink (구 Kinesis Data Analytics)cloud
- 정의 Amazon Managed Service for Apache Flink (구 Amazon Kinesis Data Analytics) 는 스트리밍 데이터를 실시간으로 처리 (…
- [AWS] S3: object storage, storage classes, lifecyclecloud
- 정의 S3 = AWS 의 object storage. bucket + key + object. 11 9's durability (99.999999999%), 무한 확장. 2026…
- [DB] PostgreSQL: 프로세스 모델, MVCC, WAL, 확장성database-internals
- 정의 PostgreSQL 은 오픈소스 ORDBMS. 1986 UC Berkeley POSTGRES 의 후예. MVCC, 확장 가능 타입, JSONB, full-text searc…
- 데이터 레이크data-engineering
- 정의 데이터 레이크 (Data Lake) 는 구조 여부와 무관하게 모든 종류의 데이터를 원본 형식 그대로 저장하는 중앙 저장소. 스키마를 미리 강제하지 않고 (schema-on-…
- 데이터 웨어하우스data-engineering
- 정의 데이터 웨어하우스 (Data Warehouse, DW) 는 여러 운영 시스템에서 흘러 들어온 이력 데이터를 통합해, 대규모 분석 쿼리를 빠르게 실행하도록 최적화된 중앙 저장…
- Apache Parquetdata-engineering
- 정의 Apache Parquet 은 분석 쿼리에 최적화된 오픈소스 컬럼형 이진 파일 포맷. 2013년 Twitter + Cloudera 가 Google Dremel 논문 (201…
- ETL / ELTdata-engineering
- 정의 ETL = Extract (추출) + Transform (변환) + Load (적재). 여러 소스에서 데이터를 뽑아 정제한 뒤 데이터 웨어하우스 에 저장하는 파이프라인. E…
- MPP (대량 병렬 처리)data-engineering
- 정의 MPP (Massively Parallel Processing, 대량 병렬 처리) 는 하나의 대규모 쿼리를 수십 - 수백 노드에 나눠 동시에 실행하는 아키텍처. 각 노드가 …
💬 댓글