핵심 차이점 요약
| 비교 항목 | Hadoop MapReduce | Apache Spark | Apache 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 |
| Spark | RDD 계보 재구성 | Exactly-once (조건부) |
| Flink | 분산 스냅샷 | 기본 Exactly-once |
프레임워크 선택 가이드
- 대규모 배치 처리: MapReduce
- 대화형 분석/ML: Spark
- 밀리초 단위 실시간 처리: Flink
- 통합 스트림-배치 아키텍처: Flink