분산 데이터 처리 프레임워크: Spark, MapReduce, Flink 핵심 비교 분석

핵심 차이점 요약

비교 항목Hadoop MapReduceApache SparkApache Flink
처리 방식배치 전용배치 중심 + 마이크로배치이벤트 기반 스트리밍
실행 모델Map → Reduce (디스크 I/O)DAG + 인메모리 (RDD)상태 유지 스트림 처리
지연 시간분~시간 단위초 단위밀리초 단위
상태 관리상태 비저장제한적 상태 관리내장형 상태 관리
장애 복구작업 재시도RDD 계보 복구분산 스냅샷

실시간 사기 탐지 사례 비교

시나리오: 고객이 5분 내 3개 도시에서 거래 시 실시간 경고

MapReduce: 적용 불가

  • 정적 데이터 처리만 가능
  • 최소 10분 지연 발생
  • 실시간 스트림 처리 미지원

Spark: 제한적 적용

# Spark Structured Streaming 예제
from pyspark.sql import SparkSession
from pyspark.sql.types import *

spark = SparkSession.builder.appName("TransactionMonitor").getOrCreate()

schema = StructType([
    StructField("client_id", StringType()),
    StructField("location", StringType()),
    StructField("timestamp", TimestampType())
])

stream_df = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "kafka-server:9092") \
    .option("subscribe", "transactions") \
    .load() \
    .select(from_json(col("value").cast("string"), schema).alias("parsed")) \
    .select("parsed.*")

result = stream_df \
    .withWatermark("timestamp", "5 minutes") \
    .groupBy(window("timestamp", "5 minutes", "1 minute"), "client_id") \
    .agg(countDistinct("location").alias("loc_count")) \
    .filter(col("loc_count") >= 3)

output = result.writeStream \
    .format("console") \
    .outputMode("update") \
    .start()
  • 장점: 배치 코드와 유사한 개발 경험
  • 단점: 마이크로배치로 인한 지연, 상태 관리 복잡

Flink: 최적의 솔루션

// Flink DataStream API 예제
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

DataStream<TransactionEvent> input = env
    .addSource(new FlinkKafkaConsumer<>("transactions", new JSONDeserializer(), properties))
    .assignTimestampsAndWatermarks(WatermarkStrategy
        .<TransactionEvent>forBoundedOutOfOrderness(Duration.ofSeconds(10))
        .withTimestampAssigner((event, ts) -> event.getTimestamp()));

input.keyBy(event -> event.clientId)
    .window(SlidingEventTimeWindows.of(Time.minutes(5), Time.minutes(1)))
    .process(new LocationAnalyzer())
    .addSink(new AlertNotifier());

env.execute("FraudDetectionJob");

// 사용자 정의 처리 함수
public class LocationAnalyzer extends ProcessWindowFunction<TransactionEvent, Alert, String, TimeWindow> {
    @Override
    public void process(String key, Context ctx, Iterable<TransactionEvent> events, Collector<Alert> out) {
        Set<String> cities = new HashSet<>();
        for (TransactionEvent e : events) {
            cities.add(e.getCity());
        }
        if (cities.size() >= 3) {
            out.collect(new Alert(key, cities));
        }
    }
}
  • 장점: 이벤트 단위 실시간 처리, 내장 상태 관리
  • 장애 복구: 체크포인트 기반 Exactly-once 보장

추가 기술적 차이점

배치 처리 성능

  • MapReduce: EB급 대규모 데이터 처리에 적합
  • Spark: 고성능 SQL 최적화 엔진 보유
  • Flink: 배치를 스트림의 특수 사례로 처리

장애 허용 메커니즘

프레임워크복구 방식데이터 일관성
MapReduce작업 재실행At-least-once
SparkRDD 계보 재구성Exactly-once (조건부)
Flink분산 스냅샷기본 Exactly-once

프레임워크 선택 가이드

  • 대규모 배치 처리: MapReduce
  • 대화형 분석/ML: Spark
  • 밀리초 단위 실시간 처리: Flink
  • 통합 스트림-배치 아키텍처: Flink

태그: Spark MapReduce Flink 분산처리 스트리밍처리

7월 30일 11:55에 게시됨