전체 그래프
Flink

Checkpointing

data-engineeringflinkcheckpointingbarrierrecovery

상위: Flink

요약

Flink의 체크포인팅 메커니즘을 다룹니다. Barrier를 통한 분산 스냅샷 생성, 장애 복구 과정, 그리고 체크포인팅 최적화 기법들을 학습하여 안정적인 스트림 처리 시스템을 구축하는 방법을 알아봅니다.

Checkpointing

  • 자동 저장 기능과 유사
    • 현재까지의 처리 상태 & 데이터를 안전한 장소에 저장
    • 장애 발생 시, 마지막 체크포인트부터 복구
    • 데이터 유실 방지

flink-18.png

  1. On checkpoint:
    • 상태를 checkpoint storage에 저장
  2. On failure:
    • 영향을 받은 연산자(operator)만 재시작
  3. After restart:
    • checkpoint storage로부터 상태 복구

스냅샷(snapshot)

  • 실행 중인 애플리케이션 상태를 주기적으로 백업
    • 저장 위치: HDFS, S3 등
    • 문제 발생 시 checkpoint를 불러와 복구

flink-19.png

데이터 유실 방지 & 시스템 장애 대비

  • Exactly-Once 보장 (중복 없이 정확히 한 번 처리)
  • 정상 종료 시 체크포인트는 자동 삭제됨

상태 스냅샷 저장 위치

  • 분산 스토리지 + 로컬 스토리지
    • 로컬 저장소는 노드 장애 시 내구성 보장 ❌
    • 다른 노드가 상태를 재분배할 권한도 없음
    • 따라서 기본 저장소는 분산 스토리지(DFS, S3 등)
  • 복구 시 Flink 동작
    • 항상 로컬 저장소에서 먼저 복원 시도
    • 단, 로컬 복구는 기본적으로 비활성화됨 (Flink 설정 필요)

flink-20.png

Barrier

  • 데이터 스트림 사이에 '장벽'을 세움
    • 소스 노드가 일정 주기마다 **Checkpoint Barrier**를 전송
    • 연산자(operator)가 해당 barrier 도착 시, 자신의 상태를 저장
  • Barrier 동기화
    • 여러 스트림 중 하나만 barrier 도착 시 → 기다렸다가 모든 스트림에 barrier 도착 시 상태 저장
    • => 동일 시점의 스냅샷 보장 (Exactly-once semantics)

적용 순서

flink-21.png

  1. Checkpoint Coordinator가 모든 소스 노드에 barrier 전송

flink-22.png

  1. 각 Source는 barrier를 포함해 데이터를 downstream으로 전송

flink-23.png

  1. Operator는 barrier 도착 시까지 입력을 기다림

flink-24.png

  1. 모든 barrier 수신 시 → 상태 저장 & JobManager에 전송

flink-25.png

  1. 최종 Sink까지 모두 완료되면 메타데이터 기록 후 완료

flink-26.png

장애 발생 시 복구 과정

  • 어플리케이션을 다시 시작(restart)
  • 가장 최근에 성공한 체크포인트 데이터를 불러옴
    • 연산자(operator)의 상태 복원
    • 체크포인트에 기록된 offset부터 재시작

flink-27.png

  • 체크포인트 이후 진행된 데이터는 다시 처리
    • 소스에서 해당 시점 이후 데이터 재전송
  • 하지만 Barrier 정렬 덕분에 중복 처리되지 않음
    • Flink가 Exactly-once 보장

Checkpointing 최적화 기술

1. 비동기 체크포인팅

  • 데이터 처리를 멈추지 않고 별도 쓰레드로 스냅샷 저장
  • 지연(latency) 최소화 → 게임에서 '백그라운드 저장'처럼 작동

2. 증분 체크포인팅 ( Incremental Checkpointing )

  • GIT처럼 변경된 부분만 저장 → 저장 속도 개선 및 저장 공간 절약

3. Barrier Alignment 최적화

  • 느린 스트림 기다리며 정렬 = 지연 발생 가능
  • **Unaligned Checkpoint**는 느린 스트림 기다리지 않음

4. 체크포인트 주기 조정

  • 너무 자주: 안전하지만 성능 저하
  • 너무 드물게: 데이터 손실 위험
  • 최적 주기가 중요

5. 체크포인트 병렬성 & 최소 간격 (min pause) & 타임아웃

  • 동시에 여러 체크포인트 처리 가능 ( **MaxConcurrentCheckpoints** )
  • 최소 간격 및 타임아웃 설정으로 안정성 확보

Checkpointing 구현 및 관리

Checkpoint 활성화

  • StreamExecutionEnvironment
    • enable_checkpointing(10000) : 10,000 밀리초(10초)마다 체크포인트를 생성
  • 체크포인트 저장소 설정
    • set_checkpoint_storage(FileSystemCheckpointStorage(CHECKPOINT_PATH))
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.datastream.checkpoint_config import CheckpointConfig
from pyflink.datastream.checkpoint_storage import FileSystemCheckpointStorage

env = StreamExecutionEnvironment.get_execution_environment()
env.set_parallelism(1)

# Checkpoint 활성화 (10초 간격)
env.enable_checkpointing(10000)

# Checkpoint 저장소 설정 (로컬)
env.get_checkpoint_config().set_checkpoint_storage(
    FileSystemCheckpointStorage("file:///tmp/flink-checkpoints")
)

Checkpoint 설정

  • 동시 체크포인트 수 설정
    • 한 번에 하나의 체크포인트만 진행
  • 체크포인트 간 최소 간격 설정
    • 체크포인트 사이에 최소 5초의 간격을 둠
  • 체크포인트 타임아웃 설정
    • 체크포인트 타임아웃을 1분으로 지정
  • 증분 체크포인트 활성화
    • RocksDB 활용 필수
# 동시 체크포인트  설정 (최대 2개)
env.get_checkpoint_config().set_max_concurrent_checkpoints(2)

# 체크포인트 타임아웃 설정 (60초)
env.get_checkpoint_config().set_checkpoint_timeout(60000)

# 체크포인트  최소 간격 설정 (5초)
env.get_checkpoint_config().set_min_pause_between_checkpoints(5000)

# 증분 체크포인트 활성화 (RocksDB 사용)
env.set_state_backend(EmbeddedRocksDBStateBackend())

전체 코드 예제 (로컬 파일 활용)

import pandas as pd
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.datastream.checkpoint_config import CheckpointConfig
from pyflink.datastream.checkpoint_storage import FileSystemCheckpointStorage

# 실행 환경 설정
env = StreamExecutionEnvironment.get_execution_environment()
env.set_parallelism(1)

# Checkpoint 설정
env.enable_checkpointing(10000)
env.get_checkpoint_config().set_max_concurrent_checkpoints(2)
env.get_checkpoint_config().set_checkpoint_timeout(60000)
env.get_checkpoint_config().set_min_pause_between_checkpoints(5000)
env.get_checkpoint_config().set_checkpoint_storage(
    FileSystemCheckpointStorage("file:///tmp/flink-checkpoints")
)

# CSV 데이터 로드
df = pd.read_csv("../data/data.csv")
transactions = df[['transaction_id', 'amount']].dropna().values.tolist()

# Flink 데이터 스트림 생성
transaction_stream = env.from_collection(transactions)

# 데이터 출력
transaction_stream.print()

# 실행
env.execute("Enable Checkpointing Example")

체크포인트 생성 결과 (로컬 디렉토리 예시)

ls -lh /tmp/flink-checkpoints
drwxr-xr-x  chk-1
drwxr-xr-x  shared
drwxr-xr-x  taskowned

flink-28.png

관련 개념