상위: Flink
요약
Kafka와 Flink를 결합하면 고신뢰 실시간 데이터 파이프라인을 구축할 수 있습니다. Kafka가 메시지 브로커 역할을, Flink가 이벤트 기반 스트림 처리를 담당하여 Exactly-once 보장과 높은 처리 탄력성을 제공합니다.
왜 실시간 처리가 필요한가?
전통적 배치 처리 방식의 한계
- 시간 지연 발생 → 사용자 행동에 즉각 대응 불가
실시간 처리의 장점
- 즉각적인 반응 및 자동화 조치 가능
- 예시: 실시간 상품 클릭 감지 후 추천 영역 개선
- 빠른 의사결정과 고객 경험 향상
- 예시: 실시간 이탈 예측 후 쿠폰 발급
사용자 행동 분석의 활용 예시
개인화 추천 시스템
- 사용자 실시간 클릭/검색 데이터를 바탕으로 상품 추천
- 예: 네이버 쇼핑, 쿠팡 추천 리스트
마케팅 자동화
- 예: 장바구니 이탈 시 자동 알림 또는 쿠폰 발송
- 이메일 마케팅 툴과 연계
사용자 흐름 분석 및 UX 개선
- 어떤 버튼을 자주 누르고, 어디서 이탈하는지 실시간 분석
- 제품/서비스 개선에 활용
Kafka + Flink 조합의 장점
| 항목 | Kafka | Flink | 조합 시 장점 |
|---|---|---|---|
| 메시징 | 고신뢰 메시지 브로커 | 메시지 소비 | 실시간 데이터 전달 |
| 처리 모델 | 로그 기반 스트림 | 이벤트 기반 스트림/배치 통합 | 높은 처리 탄력성과 확장성 |
| 보장 | At least once, Exactly once | Exactly once 지원 | 정확성 보장 |
| 유연성 | 다양한 Producer/Consumer 연동 | SQL, API 기반 처리 가능 | 사용자 맞춤 데이터 처리 가능 |
실습: Kafka + Flink 실시간 처리
1. Kafka & Zookeeper 서버 실행
# Kafka 서버 실행
/usr/local/kafka$ ./bin/kafka-server-start.sh config/server.properties
# Zookeeper 서버 실행
/usr/local/kafka$ ./bin/zookeeper-server-start.sh config/zookeeper.properties
2. Kafka Topic 생성
# Raw 소스 데이터를 담을 Topic
bin/kafka-topics.sh --create --topic user_behaviors \
--bootstrap-server localhost:9092 --partitions 1 --replication-factor 1
# Flink 처리 결과 데이터를 담을 Topic
bin/kafka-topics.sh --create --topic behavior_stats \
--bootstrap-server localhost:9092 --partitions 1 --replication-factor 1
3. Flink-Kafka 커넥터 다운로드
wget https://repo1.maven.org/maven2/org/apache/flink/flink-sql-connector-kafka/3.3.0-1.20/flink-sql-connector-kafka-3.3.0-1.20.jar
→ 해당 .jar 파일은 Python 스크립트 파일들과 동일한 디렉터리에 위치시킬 것
4. Kafka 프로듀서 스크립트 (flink_producer.py)
from kafka import KafkaProducer
import json
import random
from datetime import datetime
producer = KafkaProducer(
bootstrap_servers=['localhost:9092'],
value_serializer=lambda v: json.dumps(v).encode('utf-8')
)
user_ids = [f"user_{i}" for i in range(1, 100)]
item_ids = [f"item_{i}" for i in range(1, 200)]
categories = ["electronics", "books", "clothing", "home", "sports"]
behaviors = ["click", "view", "add_to_cart", "purchase"]
for _ in range(1000):
event = {
"user_id": random.choice(user_ids),
"item_id": random.choice(item_ids),
"category": random.choice(categories),
"behavior": random.choice(behaviors),
"ts": datetime.now().strftime("%Y-%m-%d %H:%M:%S.%f")[:-3]
}
producer.send("user_behaviors", event)
producer.close()
5. Flink 스트리밍 작업 파일 (kafka_flink_example.py)
환경 설정
import os
import logging
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment, EnvironmentSettings
def main():
logging.info("Flink 작업 시작...")
# 스트림 실행 환경 설정
env = StreamExecutionEnvironment.get_execution_environment()
env_settings = EnvironmentSettings.new_instance().in_streaming_mode().build()
table_env = StreamTableEnvironment.create(env, environment_settings=env_settings)
# 로깅 레벨 설정
table_env.get_config().get_configuration().set_string("pipeline.global-job-parameters.logger.level", "INFO")
# JAR 파일 추가 (전체 경로로 수정하세요)
kafka_jar = os.path.join(os.path.abspath('.'), 'flink-sql-connector-kafka-3.3.0-1.20.jar')
logging.info(f"사용하는 JAR 파일 경로: {kafka_jar}")
if not os.path.exists(kafka_jar):
logging.error(f"JAR 파일이 존재하지 않습니다: {kafka_jar}")
return
table_env.get_config().get_configuration().set_string("pipeline.jars", f"file://{kafka_jar}")
Kafka 소스 테이블 정의
# 소스 테이블 정의
try:
logging.info("Kafka 소스 테이블 생성 시도...")
table_env.execute_sql("""
CREATE TABLE kafka_source (
user_id STRING,
item_id STRING,
category STRING,
behavior STRING,
ts TIMESTAMP(3),
proctime AS PROCTIME()
) WITH (
'connector' = 'kafka',
'topic' = 'user_behaviors',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'flink-consumer-group',
'scan.startup.mode' = 'earliest-offset',
'format' = 'json',
'json.fail-on-missing-field' = 'false',
'json.ignore-parse-errors' = 'true'
)
""")
logging.info("Kafka 소스 테이블 생성 성공")
except Exception as e:
logging.error(f"소스 테이블 생성 실패: {e}")
return
→ Kafka 프로듀서가 보내는 user_behaviors 토픽의 JSON 데이터를 Flink에서 읽어오는 소스 정의
Kafka 싱크 테이블 정의
# 싱크 테이블 정의
try:
logging.info("Kafka 싱크 테이블 생성 시도...")
table_env.execute_sql("""
CREATE TABLE kafka_sink (
category STRING,
behavior STRING,
behavior_count BIGINT,
update_time TIMESTAMP(3),
PRIMARY KEY (category, behavior) NOT ENFORCED
) WITH (
'connector' = 'upsert-kafka',
'topic' = 'behavior_stats',
'properties.bootstrap.servers' = 'localhost:9092',
'key.format' = 'json',
'value.format' = 'json',
'properties.group.id' = 'flink-sink-group'
)
""")
logging.info("Kafka 싱크 테이블 생성 성공")
except Exception as e:
logging.error(f"싱크 테이블 생성 실패: {e}")
return
→ Flink가 집계한 데이터를 Kafka 토픽 behavior_stats로 다시 전송하는 싱크 정의
데이터 처리 및 집계
# 데이터 집계 쿼리
try:
logging.info("집계 쿼리 실행...")
table_env.execute_sql("""
INSERT INTO kafka_sink
SELECT
category,
behavior,
COUNT(*) as behavior_count,
CURRENT_TIMESTAMP as update_time
FROM kafka_source
GROUP BY category, behavior
""")
except Exception as e:
logging.error(f"쿼리 실행 실패: {e}")
if __name__ == "__main__":
main()
실습 실행 순서
1. 터미널 1: Kafka 프로듀서 실행
python flink_producer.py
2. 터미널 2: Kafka Consumer로 결과 확인
$KAFKA_HOME/bin/kafka-console-consumer.sh \
--topic behavior_stats \
--bootstrap-server localhost:9092 \
--from-beginning
3. 터미널 3: Flink 작업 실행
python kafka_flink_example.py
→ 터미널 2에 데이터가 출력되면 Flink → Kafka 연계가 정상 작동 중
프로젝트 구조

project/
├── flink_producer.py # Kafka 프로듀서
├── kafka_flink_example.py # Flink 스트리밍 작업
└── flink-sql-connector-kafka-3.3.0-1.20.jar # Kafka 커넥터
처리 흐름
- 데이터 생성: Kafka Producer → user_behaviors 토픽
- 실시간 처리: Flink가 토픽 구독 → 집계 수행
- 결과 전송: Flink → behavior_stats 토픽
- 결과 확인: Kafka Consumer → 터미널 출력
핵심 개념
Exactly-Once 보장
Flink는 Kafka와 통합 시 Exactly-once semantics 제공:
- Checkpoint와 두 단계 커밋 (2PC)
- 중복 없는 정확한 처리 보장
Backpressure
Flink가 자동으로 처리:
- Consumer 속도 > Producer 속도: 정상
- Producer 속도 > Consumer 속도: 자동 조절
Watermark
늦게 도착하는 데이터 처리:
ts TIMESTAMP(3),
WATERMARK FOR ts AS ts - INTERVAL '5' SECOND
실전 패턴
윈도우 집계
CREATE TABLE windowed_stats AS
SELECT
TUMBLE_START(ts, INTERVAL '1' MINUTE) as window_start,
category,
COUNT(*) as event_count
FROM kafka_source
GROUP BY TUMBLE(ts, INTERVAL '1' MINUTE), category
조인
# 스트림-스트림 조인
SELECT
a.user_id,
a.behavior,
b.user_info
FROM kafka_source a
JOIN user_stream b
ON a.user_id = b.user_id
WHERE a.ts BETWEEN b.ts - INTERVAL '1' HOUR AND b.ts
관련 주제
- Streaming Architecture - 전체 스트리밍 아키텍처
- Flink vs Spark - Spark와의 비교
- Structured Streaming - Spark 스트리밍 대안