본문으로 건너뛰기
김신건의 로그

Hadoop / Spark

· 수정 · 📖 약 5분 · 1,772자/단어 #data-engineering #hadoop #spark #big-data #distributed
Apache Hadoop, Hadoop, Apache Spark, Spark, Hadoop MapReduce, Spark RDD, PySpark

정의

Hadoop2006 야후에서 시작된 분산 저장 + 분산 처리 오픈소스 플랫폼. Google 이 2003-2004 년에 발표한 GFS (분산 파일시스템) 와 MapReduce (분산 처리 모델) 논문을 오픈소스로 구현한 것이 출발점.

Spark2009 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 CoreRDD, DAG, 실행 엔진
Spark SQLDataFrame + SQL 쿼리 (Catalyst 최적화)
Structured StreamingDataFrame API 로 스트리밍 (micro-batch or continuous)
MLlib분산 머신러닝 (일부 알고리즘)
GraphX그래프 분석 (거의 사용 안 됨, GraphFrames 로 대체)

MapReduce vs Spark

MapReduceSpark
저장단계마다 디스크메모리 우선 (부족 시 디스크 spill)
모델Map + Reduce 2 단계임의 DAG (map, filter, join, groupBy, …)
속도기준10-100 배 (반복 알고리즘)
APIJava 위주Scala, Python, R, SQL, Java
상호작용배치 onlyREPL / 노트북 (Zeppelin, Jupyter)
스트리밍X (Storm 별도)Structured Streaming
MLMahout 별도MLlib 내장
표현력낮음높음 (SQL, DataFrame, MLlib)
디버깅어려움Spark UI

결론: 새 프로젝트는 사실상 Spark. MapReduce 는 레거시 유지만.

관리형 Spark / Hadoop

서비스특징
Amazon EMRHadoop 생태계 완전 관리형, EC2/Serverless/EKS
AWS Glue서버리스 Spark ETL
DatabricksSpark 창시자들의 상용, Delta Lake
Google DataprocGCP 관리형
Azure HDInsight / SynapseAzure
자체 호스팅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 성능 튜닝의 큰 축

  1. 파일 포맷: Parquet + ZSTD (CSV/JSON 대비 10-100 배)
  2. 파티셔닝: partition 컬럼으로 프루닝
  3. Shuffle 최소화: repartition 남용 X, broadcast join 활용
  4. Broadcast join: 작은 테이블 (~ 100 MB) 을 executor 에 복제
  5. AQE (Adaptive Query Execution): 런타임 통계로 자동 최적화 (Spark 3.x 기본 on)
  6. Executor 크기: 너무 크면 GC 문제, 너무 작으면 오버헤드 (경험적 4-8 core, 16-32 GB)
  7. Cache / Persist: 반복 참조 데이터만
  8. 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.

WARNING

HDFS 는 클라우드에서 임시. EMR 등에서 클러스터 종료 시 데이터 손실. 결과는 S3 에.

IMPORTANT

RDD 직접 사용 지양. DataFrame / SQL 이 Catalyst 최적화 받아 훨씬 빠름. 특수한 경우만 RDD.

관련 위키

이 글의 용어 (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, 대량 병렬 처리) 는 하나의 대규모 쿼리를 수십 - 수백 노드에 나눠 동시에 실행하는 아키텍처. 각 노드가 …

💬 댓글

사이트 검색 / 명령어

검색

스크롤 = 확대/축소 · 드래그 = 이동 · 0 = 원래 크기 · ESC = 닫기