전체 그래프
Concepts

Streaming Architecture

data-engineeringstreamingarchitecturekafkareal-time

상위: 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"))

프레임워크 선택 기준

요구사항 기반 선택

요구사항FlinkSpark
지연 시간 < 100ms
배치 + 스트리밍 통합
복잡한 상태 관리
ML 파이프라인 연계
학습 곡선높음낮음
운영 복잡도높음낮음

실전 조합

많은 기업이 하이브리드 접근:

Flink:

  • Fraud detection
  • 실시간 알림
  • 복잡한 CEP

Spark:

  • 배치 ETL
  • ML 학습
  • 대시보드 집계

예시 아키텍처:

Kafka  Flink (실시간 알림)
      Spark (배치 집계)  Data Lake

모니터링 및 운영

핵심 메트릭

  1. 처리량 (Throughput)

    • 초당 처리 메시지 수
    • 목표: Producer rate와 동일
  2. 지연 시간 (Latency)

    • End-to-end latency
    • 목표: 요구사항에 따라 (ms ~ s)
  3. 백로그 (Backlog)

    • 처리되지 않은 메시지 수
    • 목표: 0에 가깝게 유지
  4. 에러율

    • 실패한 메시지 비율
    • 목표: < 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()

관련 주제