Checkpoints
Checkpoints
개요
Checkpoint는 상태와 해당 스트림 위치를 복구할 수 있게 함으로써 Flink의 상태를 장애 허용(fault tolerant)하게 만듭니다. 이를 통해 애플리케이션은 실패 없는 실행과 동일한 의미론을 갖게 됩니다.
프로그램에서 checkpoint를 활성화하고 구성하는 방법은 Checkpointing을 참조하세요.
checkpoint와 savepoint의 차이를 이해하려면 checkpoints vs. savepoints를 보세요.
출처: 문서
본문
Checkpoint 저장소
checkpointing이 활성화되면 관리 상태(managed state)는 실패 시 일관된 복구를 보장하기 위해 유지됩니다. checkpointing 중 상태가 어디에 유지되는지는 선택한 Checkpoint Storage에 따라 달라집니다.
사용 가능한 Checkpoint Storage 옵션
기본적으로 Flink는 다음 checkpoint 저장소 유형을 번들합니다.
- JobManagerCheckpointStorage
- FileSystemCheckpointStorage
checkpoint 디렉터리가 구성되면
FileSystemCheckpointStorage가 사용되고, 그렇지 않으면 시스템이JobManagerCheckpointStorage를 사용합니다.
JobManagerCheckpointStorage
JobManagerCheckpointStorage는 checkpoint 스냅샷을 JobManager의 힙에 저장합니다.
특정 크기를 초과하면 checkpoint를 실패시키도록 구성하여 JobManager에서 OutOfMemoryError를 방지할 수 있습니다. 이 기능을 설정하려면 사용자는 해당 최대 크기로 JobManagerCheckpointStorage 인스턴스를 만들 수 있습니다.
new JobManagerCheckpointStorage(MAX_MEM_STATE_SIZE);
JobManagerCheckpointStorage의 제한 사항:
- 각 개별 상태의 크기는 기본적으로 5 MB로 제한됩니다. 이 값은
JobManagerCheckpointStorage생성자에서 늘릴 수 있습니다. - 구성된 최대 상태 크기와 관계없이 상태는 Pekko 프레임 크기보다 클 수 없습니다(Configuration 참조).
- 전체 상태는 JobManager 메모리에 맞아야 합니다.
JobManagerCheckpointStorage는 다음에 권장됩니다.
- 로컬 개발 및 디버깅
- 상태를 거의 사용하지 않는 작업. 예: 레코드 단위 함수(Map, FlatMap, Filter 등)만으로 구성된 작업. Kafka Consumer는 상태를 거의 사용하지 않습니다.
FileSystemCheckpointStorage
FileSystemCheckpointStorage는 "hdfs://namenode:40010/flink/checkpoints" 또는 "file:///data/flink/checkpoints" 같은 파일 시스템 URL(타입, 주소, 경로)로 구성됩니다.
checkpointing 시 구성된 파일 시스템과 디렉터리의 파일에 상태 스냅샷을 씁니다. 최소한의 메타데이터만 JobManager 메모리(또는 고가용성 모드에서는 메타데이터 checkpoint)에 저장됩니다.
checkpoint 디렉터리가 지정되면 FileSystemCheckpointStorage가 checkpoint 스냅샷을 유지하는 데 사용됩니다.
FileSystemCheckpointStorage는 다음에 권장됩니다.
- 모든 고가용성 설정.
유지된 Checkpoints
Checkpoint는 기본적으로 유지되지 않으며 작업을 실패에서 재개하는 데만 사용됩니다. 프로그램이 취소되면 삭제됩니다. 그러나 주기적 checkpoint를 유지하도록 구성할 수 있습니다. 구성에 따라 이러한 유지된(retained) checkpoint는 작업이 실패하거나 취소될 때 자동으로 정리되지 않습니다. 이렇게 하면 작업이 실패할 때 재개할 checkpoint를 확보할 수 있습니다.
CheckpointConfig config = env.getCheckpointConfig();
config.setExternalizedCheckpointRetention(ExternalizedCheckpointRetention.RETAIN_ON_CANCELLATION);
ExternalizedCheckpointRetention 모드는 작업을 취소할 때 checkpoint에 무슨 일이 일어나는지 구성합니다.
ExternalizedCheckpointRetention.RETAIN_ON_CANCELLATION: 작업이 취소될 때 checkpoint를 유지합니다. 이 경우 취소 후 checkpoint 상태를 수동으로 정리해야 합니다.ExternalizedCheckpointRetention.DELETE_ON_CANCELLATION: 작업이 취소될 때 checkpoint를 삭제합니다. checkpoint 상태는 작업이 실패할 때만 사용할 수 있습니다.
디렉터리 구조
savepoint와 유사하게 checkpoint는 메타데이터 파일과, 상태 백엔드에 따라 일부 추가 데이터 파일로 구성됩니다. 메타데이터 파일과 데이터 파일은 구성 파일에서 execution.checkpointing.dir로 구성된 디렉터리에 저장되며, 코드에서 작업별로 지정할 수도 있습니다.
현재 checkpoint 디렉터리 레이아웃(FLINK-8531에서 도입)은 다음과 같습니다.
/user-defined-checkpoint-dir
/{job-id}
|
+ --shared/
+ --taskowned/
+ --chk-1/
+ --chk-2/
+ --chk-3/
...
SHARED 디렉터리는 여러 checkpoint의 일부일 수 있는 상태용이고, TASKOWNED는 JobManager가 절대 삭제해서는 안 되는 상태용이며, EXCLUSIVE는 하나의 checkpoint에만 속하는 상태용입니다.
checkpoint 디렉터리는 공개 API가 아니며 향후 릴리스에서 변경될 수 있습니다.
구성 파일을 통한 전역 구성
execution.checkpointing.dir: hdfs:///checkpoints/
checkpoint 구성에서 작업별 구성
Configuration config = new Configuration();
config.set(CheckpointingOptions.CHECKPOINT_STORAGE, "filesystem");
config.set(CheckpointingOptions.CHECKPOINTS_DIRECTORY, "hdfs:///checkpoints-data/");
env.configure(config);
checkpoint 저장소 인스턴스로 구성
또는 원하는 checkpoint 저장소 인스턴스를 지정하여 checkpoint 저장소를 설정할 수 있습니다. 이를 통해 쓰기 버퍼 크기 같은 저수준 구성을 설정할 수 있습니다.
Configuration config = new Configuration();
config.set(CheckpointingOptions.CHECKPOINT_STORAGE, "filesystem");
config.set(CheckpointingOptions.CHECKPOINTS_DIRECTORY, "hdfs:///checkpoints-data/");
config.set(CheckpointingOptions.FS_WRITE_BUFFER_SIZE, FILE_SIZE_THESHOLD);
env.configure(config);
유지된 checkpoint에서 재개
작업은 savepoint에서와 마찬가지로 checkpoint의 메타데이터 파일을 사용해 checkpoint에서 재개될 수 있습니다(savepoint 복원 가이드 참조). 메타데이터 파일이 자체 포함(self-contained)이 아니라면 jobmanager는 그것이 참조하는 데이터 파일에 접근할 수 있어야 합니다(위의 Directory Structure 참조).
$ bin/flink run -s :checkpointMetaDataPath [:runArgs]