상위: Flink
요약
Flink의 체크포인팅 메커니즘을 다룹니다. Barrier를 통한 분산 스냅샷 생성, 장애 복구 과정, 그리고 체크포인팅 최적화 기법들을 학습하여 안정적인 스트림 처리 시스템을 구축하는 방법을 알아봅니다.
Checkpointing
- 자동 저장 기능과 유사
- 현재까지의 처리 상태 & 데이터를 안전한 장소에 저장
- 장애 발생 시, 마지막 체크포인트부터 복구
- 데이터 유실 방지

- On checkpoint:
- 상태를 checkpoint storage에 저장
- On failure:
- 영향을 받은 연산자(operator)만 재시작
- After restart:
- checkpoint storage로부터 상태 복구
스냅샷(snapshot)
- 실행 중인 애플리케이션 상태를 주기적으로 백업
- 저장 위치: HDFS, S3 등
- 문제 발생 시 checkpoint를 불러와 복구

데이터 유실 방지 & 시스템 장애 대비
- Exactly-Once 보장 (중복 없이 정확히 한 번 처리)
- 정상 종료 시 체크포인트는 자동 삭제됨
상태 스냅샷 저장 위치
- 분산 스토리지 + 로컬 스토리지
- 로컬 저장소는 노드 장애 시 내구성 보장 ❌
- 다른 노드가 상태를 재분배할 권한도 없음
- 따라서 기본 저장소는 분산 스토리지(DFS, S3 등)
- 복구 시 Flink 동작
- 항상 로컬 저장소에서 먼저 복원 시도
- 단, 로컬 복구는 기본적으로 비활성화됨 (Flink 설정 필요)

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

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

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

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

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

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

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

- 체크포인트 이후 진행된 데이터는 다시 처리
- 소스에서 해당 시점 이후 데이터 재전송
- 하지만 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

관련 개념
- State Management - 상태 관리 기본 개념
- Savepoint - 수동 상태 저장
- Architecture - Flink 아키텍처 구조