본문 바로가기

빅데이터 플랫폼 & 아키텍처

6편. Hadoop + Kafka 기반 실시간 분석 아키텍처 – End-to-End 구현 가이드

반응형

 

6편. Hadoop + Kafka 기반 실시간 분석 아키텍처 – End-to-End 구현 가이드

실시간 데이터 처리는 이제 금융권에서도 필수 영역이다.
과거 Hadoop 배치 중심의 데이터 플랫폼에서 벗어나

Kafka + Spark Streaming + Hadoop(HDFS/Hive) + Elastic/Grafana 조합

을 활용해 거의 실시간(near-real-time) 분석 파이프라인이 구축되는 시대다.

이 글에서는 금융권 환경에서 실제 적용 가능한 엔터프라이즈 실시간 분석 아키텍처 전체 흐름을 정리했다.


🧩 1. 전체 구조(End-to-End) 아키텍처

아래는 은행권에서 실제 운영하는 구조와 거의 동일한 형태다.

[Source System]
     ↓ (로그, 거래 데이터, 이벤트)
[Kafka Producer]
     ↓
┌───────────────────────────┐
│         Kafka Cluster      │
│  - Topic 분리(raw/clean)   │
│  - Partition / Replication │
└───────────────────────────┘
     ↓
[Spark Streaming / Flink]
 - 실시간 파싱
 - 스코어링(ML 모델)
 - 필터링/정제
 - 실시간 집계
     ↓
[Hadoop HDFS/Hive]
 - Raw layer 저장
 - Clean layer 저장
 - Near real-time 집계 테이블 생성
     ↓
[ElasticSearch / Druid / ClickHouse]
 - 시각화
 - 즉시 조회
     ↓
[Grafana / Kibana]
 - 실시간 Dashboard
     ↓
[Alerting: Slack / Email / SMS]

금융권 특성상 **Hadoop(HDFS/Hive)**가 중심 저장소로,
빠른 검색용으로 ElasticSearch/Kibana를 추가하는 형태가 가장 많다.


🛠️ 2. Kafka 구성 – 금융권에서 중요한 설정

금융권 Kafka 운영 시 아래 항목은 필수다.

✔ Topic 분리 전략

  • raw_topic
  • clean_topic
  • agg_topic (집계 데이터)
  • alert_topic (이상 거래 탐지)

✔ Partition 전략

  • 거래 데이터: 키 기반 partitioning (customer_id)
  • 로그 데이터: 날짜 기반 partitioning

✔ Replication Factor

  • 금융권 최소 3
  • 장애 대응을 고려해 Rack A/B 분리 필수

✔ Retention 정책

  • raw: 3일
  • clean: 7일
  • agg: 1~3개월

⚙️ 3. Spark Streaming 실전 코드 예제

아래는 Kafka → Spark → HDFS로 실시간 적재하는 금융권 표준 코드 구성이다.

from pyspark.sql import SparkSession
from pyspark.sql.functions import *

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

df_raw = (spark.readStream
    .format("kafka")
    .option("kafka.bootstrap.servers", "broker1:9092,broker2:9092")
    .option("subscribe", "raw_topic")
    .load())

df_parsed = df_raw.selectExpr("CAST(value AS STRING)") \
    .select(from_json(col("value"), schema).alias("data")) \
    .select("data.*")

df_enriched = df_parsed.withColumn("ingest_ts", current_timestamp())

query = (df_enriched.writeStream
    .format("parquet")
    .option("path", "/data/stream/raw")
    .option("checkpointLocation", "/chk/raw")
    .trigger(processingTime="30 seconds")
    .start())

query.awaitTermination()

특징

  • checkpoint 기반 장애 복구 가능
  • schema evolution 대응 용이
  • HDFS에 실시간 적재 가능

📦 4. 실시간 집계(aggregation) 구조

Spark Streaming에서 가장 자주 쓰는 방식은 Tumbling/Sliding Window 기반 집계이다.

예: 5분 단위 거래 합계

agg = df_parsed \
    .groupBy(window(col("ts"), "5 minutes"), col("account_id")) \
    .agg(sum("amount").alias("sum_amt"))

용도

  • 실시간 거래 모니터링
  • 고객 행동 분석
  • 이상거래(Fraud) 탐지

📊 5. 실시간 분석 결과 저장 구조

금융권에서 가장 많이 쓰는 방식:

✔ HDFS/Hive (근본 저장소)

  • Raw 데이터
  • Clean 데이터
  • 집계 결과

✔ ElasticSearch (실시간 조회)

  • 고객 행동 조회
  • 트랜잭션 모니터링
  • 대시보드 연동
  • Fraud 탐지 시각화

✔ Druid/ClickHouse (대량 이벤트 빠른 조회)

고객 행동/로그 분석에 매우 효율적.


🔧 6. Alerting(이상 거래 탐지) 구성

Spark Streaming에서 특정 조건 생성 후 alert_topic으로 보내는 방식:

alert = df_parsed.filter("amount > 10000000")  # 1천만원 초과

Kafka alert_topic →
Python Alert Consumer →
Slack/Email/SMS로 전송.

금융권에서는 “신속성”을 위해 Slack + SMS 조합을 가장 선호한다.


📡 7. 배치 + 실시간 통합 아키텍처

실시간 파이프라인은 결국 배치(Hive/Spark)와 통합되어야 한다.

대표적인 통합 구조:

[실시간]
Kafka → Spark Streaming → HDFS Clean → Elastic

[배치]
HDFS Clean → Hive → Spark Batch → Mart Layer

금융권은 "원장 데이터"가 Batch 기반이므로
실시간 분석 결과도 Overnight Batch에서 반드시 재검증한다.


🧨 실전 Tip 3개

Tip 1. Kafka 스키마 변화를 역추적 가능한 구조를 유지하라.

Schema Registry(Confluent) 또는 JSON 버전 필드 필수.

Tip 2. Checkpoint를 반드시 HDFS에 두고 삭제하지 마라.

Spark Streaming 실패 시 복구 속도 차이가 매우 크다.

Tip 3. Elastic은 ingest pipeline보다 Spark에서 정제하는 것이 안정적이다.

금융권에서는 ES ingest 파이프라인보다 Spark 변환이 장애율이 훨씬 낮았다.

 

반응형