체크포인트
체크포인트 (Checkpointing)
Flink의 모든 함수와 연산자는 상태 기반(stateful)일 수 있어요. 상태 기반 함수는 개별 요소/이벤트를 처리하는 동안 데이터를 저장해 두기 때문에, 상태는 좀 더 정교한 연산의 핵심 재료가 돼요. 그 상태를 내결함성 있게 만들려면 Flink가 **체크포인트(checkpoint)**를 찍어야 해요. 체크포인트를 통해 Flink는 상태와 스트림의 위치를 복구해서, 실패가 없었던 것과 같은 시맨틱스를 애플리케이션에 제공할 수 있어요.
체크포인트의 전제 조건
Flink의 체크포인트 메커니즘은 스트림과 상태를 위한 내구성 있는 저장소와 상호작용해요. 일반적으로 다음이 필요해요.
- 일정 시간 동안 레코드를 **재생(replay)**할 수 있는 영구적인(내구성 있는) 데이터 소스. 영구 메시지 큐(Apache Kafka, RabbitMQ, Amazon Kinesis, Google PubSub)나 파일 시스템(HDFS, S3, GFS, NFS, Ceph 등)이 그런 예시예요.
- 상태를 위한 영구 저장소. 보통 분산 파일 시스템(HDFS, S3, GFS, NFS, Ceph 등)이에요.
체크포인트 활성화와 설정
기본적으로 체크포인트는 비활성화되어 있어요. 활성화하려면 StreamExecutionEnvironment에서 enableCheckpointing(n)을 호출하면 되는데, *n**은 체크포인트 주기(밀리초)**예요.
체크포인트의 다른 파라미터로는 이런 것들이 있어요.
- checkpoint storage: 체크포인트 스냅샷을 어디에 내구성 있게 저장할지 정해요. 기본은 JobManager의 힙이고, 운영 배포에서는 내구성 있는 파일 시스템을 권장해요.
- exactly-once vs. at-least-once:
enableCheckpointing(n)메서드에 모드를 전달해 두 보장 수준을 고를 수 있어요. 대부분의 애플리케이션에서 exactly-once가 선호되고, 일관되게 초저지연(수 밀리초)을 요구하는 애플리케이션에 at-least-once가 유용할 수 있어요. - checkpoint timeout: 진행 중인 체크포인트가 이 시간 안에 완료되지 않으면 중단돼요.
- minimum time between checkpoints: 체크포인트 사이에 최소한의 시간을 정의해, 스트리밍 애플리케이션이 체크포인트 사이에 어느 정도 진행을 보장하게 해요. 예를 들어 5000으로 설정하면 이전 체크포인트가 완료된 뒤 최소 5초가 지나야 다음 체크포인트가 시작돼요. 이 값은 체크포인트가 때로 평균보다 오래 걸린다고 해도 영향을 받지 않아서, 체크포인트 주기보다 "체크포인트 사이 시간"으로 설정하는 편이 더 쉽기는 한데, 값이 설정되면 동시 체크포인트 수는 1로 고정돼요.
- tolerable checkpoint failure number: 연속 체크포인트 실패가 몇 번까지 허용된 뒤 전체 잡이 페일오버되는지를 정해요. 기본값은
0이라 체크포인트 실패가 처음 보고되면 잡이 실패해요. - number of concurrent checkpoints: 기본적으로 진행 중인 체크포인트가 있으면 다른 체크포인트를 트리거하지 않아요. 여러 겹치는 체크포인트를 허용할 수도 있는데, 처리 지연이 있는 파이프라인에서 유용해요.
- externalized checkpoints: 주기적 체크포인트를 외부에 영속시키도록 설정할 수 있어요. 외부화된 체크포인트는 메타데이터를 영구 저장소에 쓰고 잡이 실패해도 자동으로 정리되지 않아서, 실패 시 재개할 체크포인트가 남아요.
- unaligned checkpoints: 백프레셔 상황에서 체크포인팅 시간을 크게 줄여주는 정렬 없는 체크포인트를 활성화할 수 있어요. exactly-once 체크포인트에서만, 그리고 동시 체크포인트 1개일 때만 동작해요.
- checkpoints with finished tasks: 기본적으로 Flink는 DAG의 일부가 모든 레코드를 처리하고 끝났더라도 체크포인트를 계속 수행해요.
다음 코드가 대표적인 체크포인트 설정 예시예요.
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 1000ms마다 체크포인트 시작
env.enableCheckpointing(1000);
// 고급 옵션들:
// 모드를 exactly-once로 (기본값)
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
// 체크포인트 사이에 500ms의 진행 보장
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(500);
// 체크포인트는 1분 안에 완료되거나 폐기
env.getCheckpointConfig().setCheckpointTimeout(60000);
// 연속 체크포인트 실패 2번까지 허용
env.getCheckpointConfig().setTolerableCheckpointFailureNumber(2);
// 동시에 진행 중인 체크포인트는 1개만 허용
env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);
// 잡 취소 시에도 유지되는 외부화 체크포인트 활성화
env.getCheckpointConfig().setExternalizedCheckpointRetention(
ExternalizedCheckpointRetention.RETAIN_ON_CANCELLATION);
// 정렬 없는 체크포인트 활성화
env.getCheckpointConfig().enableUnalignedCheckpoints();
env = StreamExecutionEnvironment.get_execution_environment()
# 1000ms마다 체크포인트 시작
env.enable_checkpointing(1000)
# 모드를 exactly-once로 (기본값)
env.get_checkpoint_config().set_checkpointing_mode(CheckpointingMode.EXACTLY_ONCE)
# 체크포인트 사이에 500ms의 진행 보장
env.get_checkpoint_config().set_min_pause_between_checkpoints(500)
# 체크포인트는 1분 안에 완료되거나 폐기
env.get_checkpoint_config().set_checkpoint_timeout(60000)
# 연속 체크포인트 실패 2번까지 허용
env.get_checkpoint_config().set_tolerable_checkpoint_failure_number(2)
# 동시에 진행 중인 체크포인트는 1개만 허용
env.get_checkpoint_config().set_max_concurrent_checkpoints(1)
# 잡 취소 시에도 유지되는 외부화 체크포인트 활성화
env.get_checkpoint_config().set_externalized_checkpoint_retention(ExternalizedCheckpointRetention.RETAIN_ON_CANCELLATION)
# 정렬 없는 체크포인트 활성화
env.get_checkpoint_config().enable_unaligned_checkpoints()
체크포인트 저장소 선택
Flink의 체크포인트 메커니즘은 타이머와 상태 연산자(커넥터, 윈도우, 사용자 정의 상태 포함)에 있는 모든 상태의 일관된 스냅샷을 저장해요. 체크포인트가 어디에 저장되는지(JobManager 메모리, 파일 시스템, 데이터베이스 등)는 설정된 **체크포인트 저장소(Checkpoint Storage)**에 달려 있어요.
기본적으로 체크포인트는 JobManager의 메모리에 저장돼요. 큰 상태를 제대로 영속화하려면 Flink는 다른 위치에 상태를 체크포인트하는 여러 방법을 지원해요. 운영 배포에서는 체크포인트를 고가용성 파일 시스템에 저장하는 것을 강력히 권장해요.
Configuration config = new Configuration();
config.set(CheckpointingOptions.CHECKPOINT_STORAGE, "filesystem");
config.set(CheckpointingOptions.CHECKPOINTS_DIRECTORY, "...");
env.configure(config);
반복 잡에서의 상태 체크포인트
Flink는 현재 반복(iteration)이 없는 잡에 대해서만 처리 보장을 제공해요. 반복 잡에서 체크포인팅을 활성화하면 예외가 나요. 반복 프로그램에서 체크포인팅을 강제하려면 활성화할 때 특수 플래그를 설정해야 해요: env.enableCheckpointing(interval, CheckpointingMode.EXACTLY_ONCE, force = true).
루프 엣지의 인플라이트 레코드(그에 딸린 상태 변경 포함)는 실패 시 유실된다는 점을 알아 두세요.
그래프 일부가 끝난 상태의 체크포인트
Flink 1.14부터 잡 그래프 일부가 모든 데이터 처리를 마쳤어도 체크포인트를 계속 수행할 수 있어요. 유계 소스가 포함된 경우에 일어날 수 있죠. 이 기능은 1.15부터 기본 활성화되어 있고, 기능 플래그로 비활성화할 수 있어요. 태스크/서브태스크가 끝나면 더 이상 체크포인트에 기여하지 않아요. 커스텀 연산자나 UDF를 구현할 때 이 점을 유의해야 해요.
끝난 태스크와 함께 체크포인트하는 것을 지원하기 위해 Flink는 태스크 생명주기를 조정하고 StreamOperator#finish 메서드를 도입했어요. 이 메서드는 남은 버퍼드 상태를 플러시하는 명확한 절단 지점으로 기대돼요. finish 메서드가 호출된 뒤에 찍은 체크포인트는 (대부분의 경우) 비어 있어야 하고 버퍼드 데이터를 포함하지 않아야 해요. 다만 주목할 예외는 연산자가 외부 시스템의 트랜잭션을 가리키는 포인터를 갖는 경우(exactly-once 시맨틱 구현)예요. 그런 경우 finish() 메서드 호출 후 찍은 체크포인트는 연산자가 닫히기 전 최종 체크포인트에서 커밋될 마지막 트랜잭션(들)을 가리켜야 해요. 이에 대한 내장 사례가 exactly-once 싱크와 TwoPhaseCommitSinkFunction이에요.
연산자 상태에는 어떤 영향이 있나?
UnionListState는 외부 시스템의 오프셋 전역 뷰(Kafka 파티션의 현재 오프셋 저장 같은)를 구현하는 데 자주 쓰였는데 특별히 다뤄져요. close 메서드가 호출된 단일 서브태스크의 상태를 버리면 그 서브태스크에 배정된 파티션의 오프셋이 유실되므로, UnionListState를 쓰는 서브태스크가 모두 끝났거나 하나도 안 끝났을 때만 체크포인트가 성공하도록 해요. ListState는 비슷하게 쓰이는 걸 보지 못했지만, close 메서드 이후에 체크포인트된 상태는 버려져 복구 시에 없게 된다는 점을 알아 두세요.
태스크 종료 전 최종 체크포인트 대기
2단계 커밋(two-phase commit)을 쓰는 연산자의 모든 레코드가 커밋되도록, 태스크는 모든 연산자가 끝난 뒤 최종 체크포인트가 성공적으로 완료될 때까지 기다려요. 최종 체크포인트는 모든 연산자가 데이터 끝에 도달하면 주기 트리거를 기다리지 않고 즉시 트리거되지만, 잡은 이 최종 체크포인트가 완료될 때까지 기다려야 해요.
체크포인트 파일 병합 메커니즘 (Experimental)
Flink 1.20에 MVP(최소 기능 제품) 기능으로 도입된 unified file merging 메커니즘은 흩어진 작은 체크포인트 파일들을 큰 파일로 묶어서, 파일 생성·삭제 횟수를 줄여 체크포인트 중 파일 홍수 문제가 만드는 파일 시스템 메타데이터 관리 부담을 덜어줘요. execution.checkpointing.file-merging.enabled를 true로 설정하면 활성화돼요. 단, 이 메커니즘을 켜면 공간 증폭(space amplification)이 생길 수 있어요. 즉 실제 파일 시스템 점유가 실제 상태 크기보다 커질 수 있죠. execution.checkpointing.file-merging.max-space-amplification으로 공간 증폭의 상한을 제한할 수 있어요.
이 메커니즘은 Flink의 키 상태, 연산자 상태, 채널 상태에 적용돼요. 공유 범위 상태에는 서브태스크 수준 병합을, 전용 범위 상태에는 TaskManager 수준 병합을 제공해요. 단일 파일에 쓸 수 있는 최대 서브태스크 수는 execution.checkpointing.file-merging.max-subtasks-per-file로 설정할 수 있어요. 또한 execution.checkpointing.file-merging.across-checkpoint-boundary를 true로 설정하면 체크포인트 경계를 넘어 파일을 병합하는 것도 지원돼요.
더 알아보기 (Learn more)
- 상태 기반 스트림 처리 — 체크포인트가 상태를 내결함성 있게 만드는 원리
- 사용자 정의 함수 (Working with State) — 체크포인트로 저장되는 상태 API
- 윈도우 (Windows) — 상태를 쓰는 윈도우 연산
- 시간 기반 스트림 처리 — 이벤트 시간·워터마크