MapReduce로 대용량 배치 처리를 설계하는 방법

MapReduce의 맵·셔플·리듀스 처리 구조와 HDFS 데이터 지역성, 장애 복구, Hadoop Streaming 운영 조건을 정리한다.

2026-08-14 · 최초 발행 2025-10-14

MapReduce는 불안정한 범용 하드웨어(commodity hardware)로 구성된 클러스터에서 대용량 데이터를 병렬 처리하기 위한 프로그래밍 모델이자 실행 프레임워크다. 작업을 맵 단계와 리듀스 단계로 나누고, 데이터가 있는 위치에서 계산을 수행하며, 장애가 발생해도 재시도·재실행으로 처리를 이어간다.

맵과 리듀스로 분리한 분산 처리 모델

MapReduce의 목적은 대용량 데이터 분산 병렬 처리(MPP)를 위한 단순하고 확장 가능한 모델을 제공하는 데 있다. 개발자는 변환과 집계를 위한 함수를 작성하고, 입력 분할·태스크 배치·중간 결과 전달·실패 복구는 프레임워크가 맡는다.

map(k1, v1) -> list(k2, v2)
reduce(k2, list(v2)) -> list(k3, v3)

실행은 입력 분할에서 시작해 Map, 로컬 Spill/Sort, Shuffle/Sort, Reduce, HDFS 출력 순으로 진행된다. 맵은 입력 레코드를 중간 키-값으로 바꾸고, 리듀스는 같은 키에 모인 값을 처리해 결과를 만든다.

사용자 정의 함수(UDF)로 맵과 리듀스를 분리하는 구조이므로 부수효과를 줄이고 키-값 불변성을 지키는 편이 좋다. 조합법칙(associative/commutative)이 성립하는 연산에 특히 적합하다.

셔플과 데이터 지역성이 만드는 처리 경로

맵 태스크의 출력은 로컬 디스크에 spill된 뒤 키 기준으로 정렬되고 파티셔닝된다. 이후 셔플은 각 키에 해당하는 데이터를 리듀서로 전달하고 병합한다. 이 구간은 노드 간 전송이 집중되는 병목이므로 압축·파티셔너·콤바이너가 성능 설계의 중심이 된다.

HDFS는 블록(통상 128MB/256MB) 단위로 입력을 나누며, NameNode는 메타데이터를, DataNode는 데이터 블록을 보관한다. 스케줄러는 가능한 한 데이터가 있는 노드에서 맵 태스크를 실행해 네트워크 I/O를 줄인다.

Node #2Node #1Submit JobHeartbeat/StatusYesRetry/SpeculativeNoClient LibraryJobTracker /ApplicationMasterNameNode (HDFSMetadata)TaskTracker / NodeManager#1TaskTracker / NodeManager#2Map Task(s)Spill & Local SortMap Task(s)Spill & Local SortShuffle & Merge by KeyReduce Task(s)HDFS OutputTask Failure?Re-run TaskJob Complete

실행 단계마다 달라지는 병목과 복구 방식

단계 역할 병렬성(확장성) 네트워크 I/O 상태/정렬 장애 처리
Map 입력 스플릿 처리, 키-값 생성 스플릿 수만큼 병렬 낮음(로컬 spill 위주) spill 시 키 정렬·파티션 태스크 재시도, 스펙큘러티브 실행
Shuffle/Sort 키 기준 재분배·머지 리듀서별 다중 페치 높음(노드 간 전송) 전역 키 정렬·그룹화 페치 실패 재시도, 맵 재실행
Reduce 그룹별 집계·출력 리듀서 수만큼 병렬 입력 수신·출력 쓰기 그룹 상태 유지 가능 커밋 프로토콜로 일관성 보장

신뢰하지 않는 노드를 전제로 태스크 재시도와 스펙큘러티브 실행(speculative execution)을 수행한다. 맵 출력이 사라지면 해당 맵만 다시 실행할 수 있으며, OutputCommitter는 리듀스 결과의 커밋 단계를 보장해 중복 실행 중에도 최종 산출물의 일관성을 유지한다.

Hadoop MRv1에서는 JobTracker가 작업 계획·스케줄·모니터링을 맡고 TaskTracker가 태스크를 실행한다. HDFS의 NameNode/DataNode와는 분리된 역할이다. Hadoop YARN(MRv2)에서는 ResourceManager, ApplicationMaster, NodeManager로 역할이 분리되며 MapReduce는 하나의 애플리케이션 프레임워크로 동작한다. 최신 정보 확인이 필요하다.

배치 집계와 ETL에서의 설계 선택

로그 집계와 KPI 배치 리포팅에서는 웹·앱 로그를 시간·사용자 기준으로 모아 일/시간 단위 KPI를 산출할 수 있다. 중간 카운트를 콤바이너로 줄이고 Snappy/Zstd 압축을 적용하면 셔플 I/O를 줄일 수 있다.

대규모 조인 기반 ETL에서는 대용량 fact 테이블과 사전 테이블을 조인한 뒤 정제·파티션 출력한다. 소테이블은 분산 캐시를 이용한 맵사이드 조인을 우선하고, 키 편중은 커스텀 파티셔너로 완화한다.

역색인(Inverted Index) 생성에서는 문서를 토큰과 문서 리스트의 매핑으로 바꾼다. 동일 키 내부의 순서가 필요하면 Secondary sort를 적용하고, stop-word 필터링은 맵 단계에서 먼저 수행한다.

텍스트보다 Parquet/ORC와 블록 압축을 선택하면 셔플과 저장 비용을 낮출 수 있지만 CPU 압축 비용이 생긴다. 커스텀 파티셔너와 샘플링은 키 편중을 완화하는 대신 구현 복잡도를 높인다.

태스크 튜닝에서는 split 크기(예: 256MB 이상)와 메모리 정렬 버퍼(mapreduce.task.io.sort.mb)를 조정한다. 핫스팟 환경에서는 스펙큘러티브 실행이 중복 셔플 비용을 만들 수 있으므로 비활성화를 고려한다.

반복·인터랙티브 워크로드에는 MapReduce가 맞지 않을 수 있어 Spark/Flink 등 메모리 중심 엔진을 고려한다. 반대로 배치 처리의 안정성과 재현성이 우선인 작업에서는 MapReduce를 선택할 수 있다.

Hadoop Streaming으로 WordCount 실행하기

전제조건은 Hadoop 3.3+ 클러스터와 HDFS 접근 권한이다. Python 3.8+를 노드에 배포하고, 입력은 HDFS 경로 /data/input, 출력은 /data/output_wordcount로 둔다.

mapper.py

#!/usr/bin/env python3
import sys
for line in sys.stdin:
    for word in line.strip().split():
        if word:
            print(f"{word}\t1")

reducer.py

#!/usr/bin/env python3
import sys

current_word = None
current_count = 0

for line in sys.stdin:
    try:
        word, count_str = line.rstrip("\n").split("\t", 1)
        count = int(count_str)
    except ValueError:
        continue
    if current_word == word:
        current_count += count
    else:
        if current_word is not None:
            print(f"{current_word}\t{current_count}")
        current_word = word
        current_count = count

if current_word is not None:
    print(f"{current_word}\t{current_count}")

실행 명령은 다음과 같다.

hadoop jar $HADOOP_HOME/share/hadoop/tools/lib/hadoop-streaming-*.jar \
  -D mapreduce.job.reduces=2 \
  -D mapreduce.map.output.compress=true \
  -D mapreduce.output.fileoutputformat.compress=true \
  -files mapper.py,reducer.py \
  -mapper "python3 mapper.py" \
  -combiner "python3 reducer.py" \
  -reducer "python3 reducer.py" \
  -input /data/input \
  -output /data/output_wordcount

콤바이너는 합계처럼 교환/결합 법칙이 성립하는 연산에만 사용한다. 소규모 파일이 다수인 경우에는 HDFS 소파일 문제가 생길 수 있으므로 사전 병합(archive) 또는 HDFS concat을 고려한다.

수백 노드 규모까지 선형에 근접한 처리량 확장성을 달성할 수 있으며, 데이터 지역성은 네트워크 I/O를 50~90% 절감할 수 있다(워크로드·압축 설정에 의존). 장애 자동 복구는 재처리율을 낮추고 운영자 개입을 최소화한다. 단순한 프로그래밍 모델은 재현성과 감사 추적성을 확보하게 하며, 컴모디티 하드웨어 기반 운영은 TCO와 벤더 종속을 줄이고 배치 파이프라인 표준화는 데이터 거버넌스 체계화에 연결된다.

MapReduceHadoopHDFS분산 처리빅데이터