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 변환이 장애율이 훨씬 낮았다.
- 1편 — 금융권은 왜 하둡을 쓰는가: 도입 배경과 진화
- 2편 — 하둡 구성도: HDFS·YARN·Hive·Spark 운영 흐름
- 3편 — 운영 자동화 스크립트 모음 (금융권 표준)
- 4편 — 하둡 보안 아키텍처: 계정·권한·감사·데이터 보호
- 5편 — Job Template·배포 체계 표준화
- 6편 — Hadoop + Kafka 실시간 분석 아키텍처 (현재 글)
- 7편 — 클러스터 운영 체계: 폐쇄망·보안·권한 관리
- 8편 — Sqoop·Oozie·Spark Batch 적재 파이프라인
- 9편 — 실사용 ①: 신용평가 / 여신 리스크 모델링
- 10편 — 실사용 ②: FDS 이상거래탐지 로그 분석
'빅데이터 플랫폼 & 아키텍처' 카테고리의 다른 글
| 8편. 데이터 적재 파이프라인 구축 – Sqoop/Oozie/Spark Batch 실전 (0) | 2025.12.01 |
|---|---|
| 7편. 금융권 하둡 클러스터 운영 체계 – 내부망·보안·권한 관리 완전 가이드 (0) | 2025.11.28 |
| 5편. Hadoop 운영 자동화 – 스크립트·Job Template·배포 체계 표준화 (0) | 2025.11.28 |
| 4편. 금융권 Hadoop 보안 아키텍처 – 계정·권한·감사·데이터 보호 실전 가이드 (0) | 2025.11.28 |
| 3편. Hadoop 운영 자동화 스크립트 모음 – 금융권 표준 실전 스크립트 공개 (0) | 2025.11.28 |