전체 그래프
Flink

Kafka Integration

data-engineeringflinkkafkastreamingreal-time

상위: Flink

요약

Kafka와 Flink를 결합하면 고신뢰 실시간 데이터 파이프라인을 구축할 수 있습니다. Kafka가 메시지 브로커 역할을, Flink가 이벤트 기반 스트림 처리를 담당하여 Exactly-once 보장과 높은 처리 탄력성을 제공합니다.

왜 실시간 처리가 필요한가?

전통적 배치 처리 방식의 한계

  • 시간 지연 발생 → 사용자 행동에 즉각 대응 불가

실시간 처리의 장점

  • 즉각적인 반응 및 자동화 조치 가능
    • 예시: 실시간 상품 클릭 감지 후 추천 영역 개선
  • 빠른 의사결정과 고객 경험 향상
    • 예시: 실시간 이탈 예측 후 쿠폰 발급

사용자 행동 분석의 활용 예시

개인화 추천 시스템

  • 사용자 실시간 클릭/검색 데이터를 바탕으로 상품 추천
  • 예: 네이버 쇼핑, 쿠팡 추천 리스트

마케팅 자동화

  • 예: 장바구니 이탈 시 자동 알림 또는 쿠폰 발송
  • 이메일 마케팅 툴과 연계

사용자 흐름 분석 및 UX 개선

  • 어떤 버튼을 자주 누르고, 어디서 이탈하는지 실시간 분석
  • 제품/서비스 개선에 활용

Kafka + Flink 조합의 장점

항목KafkaFlink조합 시 장점
메시징고신뢰 메시지 브로커메시지 소비실시간 데이터 전달
처리 모델로그 기반 스트림이벤트 기반 스트림/배치 통합높은 처리 탄력성과 확장성
보장At least once, Exactly onceExactly 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 연계가 정상 작동 중

프로젝트 구조

Kafka Integration

project/
├── flink_producer.py              # Kafka 프로듀서
├── kafka_flink_example.py         # Flink 스트리밍 작업
└── flink-sql-connector-kafka-3.3.0-1.20.jar  # Kafka 커넥터

처리 흐름

  1. 데이터 생성: Kafka Producer → user_behaviors 토픽
  2. 실시간 처리: Flink가 토픽 구독 → 집계 수행
  3. 결과 전송: Flink → behavior_stats 토픽
  4. 결과 확인: 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

관련 주제