상위: Data Engineering
요약
스트리밍 아키텍처는 실시간 데이터 처리를 위한 시스템 설계 패턴입니다. Kafka가 메시지 브로커로 데이터를 전달하고, Flink 또는 Spark가 스트림 처리를 수행하여 즉각적인 인사이트를 제공합니다.
왜 실시간 처리가 필요한가?
전통적 배치 처리 방식의 한계
- 시간 지연 발생 → 사용자 행동에 즉각 대응 불가
- 일정 시간마다 데이터를 모아서 처리
- 실시간 의사결정 불가능
실시간 처리의 장점
1. 즉각적인 반응
- 예시: 실시간 상품 클릭 감지 후 추천 영역 개선
- 사용자 행동 즉시 분석
- 개인화된 경험 제공
2. 빠른 의사결정
- 예시: 실시간 이탈 예측 후 쿠폰 발급
- 비즈니스 기회 즉시 포착
- 고객 이탈 방지
3. 자동화 조치
- 이상 탐지 즉시 알림
- 자동 스케일링
- 실시간 모니터링
스트리밍 아키텍처 구성 요소
1. 메시지 브로커 (Kafka)
역할:
- 데이터 생산자(Producer)와 소비자(Consumer) 분리
- 데이터 버퍼링 및 내구성 보장
- 여러 소비자에게 동시 전달
특징:
- 고처리량, 저지연
- 수평 확장 가능
- 데이터 복제 및 내구성
2. 스트림 처리 엔진
Flink
적합한 경우:
- 초저지연 요구 (밀리초)
- 복잡한 이벤트 처리
- 강력한 상태 관리
특징:
- True Streaming
- Exactly-once 보장
- 이벤트 타임 처리
Spark Structured Streaming
적합한 경우:
- 배치 + 스트리밍 통합
- ML 파이프라인 연계
- 기존 Spark 인프라 활용
특징:
- 마이크로 배치
- DataFrame API
- MLlib 통합
3. 데이터 저장소
- 실시간 DB: Redis, Cassandra
- 분석 DB: ClickHouse, Druid
- 데이터 레이크: S3, HDFS
일반적인 스트리밍 패턴
1. Event-Driven 아키텍처
Producer → Kafka → Stream Processor → Sink
│ │ │
│ │ └→ Dashboard
│ └→ Alerts
└→ Logs
사용 케이스:
- 사용자 행동 추적
- 실시간 로그 분석
- IoT 센서 데이터
2. Lambda / Kappa 아키텍처
두 패턴의 구조·장단점·선택 기준은 Lambda vs Kappa Architecture 참고.
적합:
- 재처리 용이한 경우
- 실시간 우선
사용자 행동 분석 활용
1. 개인화 추천 시스템
흐름:
사용자 클릭 → Kafka → Flink/Spark → 추천 모델 → 실시간 추천
예시:
- 네이버 쇼핑 실시간 추천
- 쿠팡 개인화 리스트
- YouTube 동영상 추천
구현:
# 실시간 추천 업데이트
user_clicks.join(item_features, "item_id") \
.groupBy("user_id") \
.apply(recommendation_model) \
.sink_to_redis()
2. 마케팅 자동화
흐름:
장바구니 이탈 감지 → 규칙 엔진 → 쿠폰 발급 → 이메일/푸시
예시:
- 장바구니 이탈 시 10분 내 쿠폰 발송
- 검색만 하고 구매 안 한 사용자에게 할인 알림
구현:
# 이탈 감지 및 자동 쿠폰 발급
abandoned_carts = user_events \
.filter(lambda x: x.event_type == "add_to_cart") \
.windowBy(session_window("5 minutes")) \
.apply(detect_abandonment) \
.sink_to_marketing_system()
3. 사용자 흐름 분석
흐름:
페이지 뷰 → Kafka → 실시간 집계 → Dashboard
예시:
- 어떤 버튼을 자주 클릭하는지
- 어디서 이탈이 많이 발생하는지
- 실시간 전환율 모니터링
구현:
# 실시간 Funnel 분석
page_views.groupBy("page", window("1 minute")) \
.count() \
.join(conversions, "session_id") \
.calculate_conversion_rate()
Kafka + Flink/Spark 통합 패턴
Exactly-Once 보장
Flink:
# Checkpoint 설정
env.enable_checkpointing(60000) # 1분마다
env.get_checkpoint_config().set_checkpointing_mode(CheckpointingMode.EXACTLY_ONCE)
Spark:
# Kafka Offset 관리
query = df.writeStream \
.option("checkpointLocation", "/path/to/checkpoint") \
.start()
Backpressure 처리
Flink: 자동 조절
- Consumer 처리 속도에 맞춰 Producer 조절
Spark: Trigger 설정
.trigger(processingTime='10 seconds')
상태 관리
Flink: Keyed State
# 사용자별 상태 유지
.keyBy("user_id") \
.process(StatefulFunction())
Spark: Window + Watermark
df.withWatermark("timestamp", "10 minutes") \
.groupBy("user_id", window("1 hour"))
프레임워크 선택 기준
요구사항 기반 선택
| 요구사항 | Flink | Spark |
|---|---|---|
| 지연 시간 < 100ms | ✅ | ❌ |
| 배치 + 스트리밍 통합 | △ | ✅ |
| 복잡한 상태 관리 | ✅ | △ |
| ML 파이프라인 연계 | △ | ✅ |
| 학습 곡선 | 높음 | 낮음 |
| 운영 복잡도 | 높음 | 낮음 |
실전 조합
많은 기업이 하이브리드 접근:
Flink:
- Fraud detection
- 실시간 알림
- 복잡한 CEP
Spark:
- 배치 ETL
- ML 학습
- 대시보드 집계
예시 아키텍처:
Kafka → Flink (실시간 알림)
→ Spark (배치 집계) → Data Lake
모니터링 및 운영
핵심 메트릭
-
처리량 (Throughput)
- 초당 처리 메시지 수
- 목표: Producer rate와 동일
-
지연 시간 (Latency)
- End-to-end latency
- 목표: 요구사항에 따라 (ms ~ s)
-
백로그 (Backlog)
- 처리되지 않은 메시지 수
- 목표: 0에 가깝게 유지
-
에러율
- 실패한 메시지 비율
- 목표: < 0.1%
장애 복구
Flink:
- Checkpoint에서 재시작
- Savepoint로 버전 관리
Spark:
- Checkpoint에서 재시작
- Write-ahead log
실전 예시
실시간 Fraud Detection
# Flink 패턴 매칭
pattern = Pattern.begin("login") \
.where(lambda x: x.event == "login") \
.next("purchase") \
.where(lambda x: x.amount > 1000) \
.within(Time.minutes(5))
alerts = CEP.pattern(transactions, pattern) \
.select(fraud_alert_function) \
.addSink(AlertSink())
실시간 추천
# Spark 스트리밍
user_actions = spark.readStream.format("kafka").load()
recommendations = user_actions \
.join(item_features, "item_id") \
.groupBy("user_id") \
.agg(collect_list("item_id").alias("recent_items")) \
.join(recommendation_model, "user_id")
recommendations.writeStream \
.format("redis") \
.start()
관련 주제
- Flink - 실시간 스트리밍 프레임워크
- Spark - 통합 데이터 처리 프레임워크
- Kafka Integration - Kafka + Flink 통합
- Structured Streaming - Kafka + Spark 통합
- Batch vs Real-time - 배치 vs 실시간 처리