상태 스냅샷을 통한 장애 허용
상태 스냅샷을 통한 장애 허용 (Fault Tolerance via State Snapshots)
이 페이지는 Flink가 상태 스냅샷(state snapshot)을 통해 장애 허용(fault tolerance)을 어떻게 제공하는지 설명해요. 상태 백엔드, 체크포인트 저장소, 상태 스냅샷 동작 방식, exactly-once 보장에 대해 다뤄요.
출처: 문서
본문
상태 백엔드 (State Backends)
Flink가 관리하는 keyed state는 일종의 샤딩된(sharded) key/value 저장소이며, keyed state의 각 항목의 작업 복사본은 해당 key를 담당하는 taskmanager 어딘가에 로컬로 유지돼요. Operator state도 그것이 필요한 머신(들)에 로컬로 존재해요.
Flink가 관리하는 이 상태는 상태 백엔드(state backend)에 저장돼요. 두 가지 상태 백엔드 구현이 사용 가능해요 — 작업 상태를 디스크에 유지하는 내장 key/value 저장소인 RocksDB 기반의 것과, 작업 상태를 Java 힙의 메모리에 유지하는 또 다른 힙 기반 상태 백엔드예요.
| Name | Working State | Snapshotting |
|---|---|---|
| EmbeddedRocksDBStateBackend | Local disk (tmp dir) | Full / Incremental |
| 사용 가능한 메모리보다 큰 상태 지원. 일반적으로 힙 기반 백엔드보다 10x 느림 | ||
| HashMapStateBackend | JVM Heap | Full |
| 빠르며, 큰 힙 필요. GC에 영향받음 |
힙 기반 상태 백엔드에 유지되는 상태로 작업할 때 접근과 업데이트는 힙에서 객체를 읽고 쓰는 것을 포함해요. 하지만 EmbeddedRocksDBStateBackend에 유지되는 객체의 경우 접근과 업데이트는 직렬화와 역직렬화를 포함하므로 훨씬 더 비싸요. 그러나 RocksDB로 가질 수 있는 상태의 양은 로컬 디스크의 크기에 의해서만 제한돼요. 또한 EmbeddedRocksDBStateBackend만 증분 스냅샷(incremental snapshotting)을 할 수 있다는 점도 유의하세요. 이는 천천히 변하는 대량의 상태를 가진 애플리케이션에 중요한 이점이에요.
이 두 상태 백엔드 모두 비동기 스냅샷을 할 수 있어요. 즉, 진행 중인 스트림 처리를 방해하지 않고 스냅샷을 찍을 수 있어요.
체크포인트 저장소 (Checkpoint Storage)
Flink는 모든 연산자의 모든 상태에 대해 주기적으로 영구 스냅샷을 찍고, 이 스냅샷을 분산 파일 시스템 같은 더 내구성 있는 곳에 복사해요. 실패가 발생하면 Flink는 애플리케이션의 완전한 상태를 복원하고 아무 일도 없었던 것처럼 처리를 재개할 수 있어요.
이 스냅샷이 저장되는 위치는 job의 체크포인트 저장소(checkpoint storage)로 정의돼요. 두 가지 체크포인트 저장소 구현이 사용 가능해요 — 하나는 상태 스냅샷을 분산 파일 시스템에 지속시키고, 다른 하나는 JobManager의 힙을 사용해요.
| Name | State Backup |
|---|---|
| FileSystemCheckpointStorage | Distributed file system |
| 매우 큰 상태 크기 지원. 높은 내구성. 프로덕션 배포에 권장 | |
| JobManagerCheckpointStorage | JobManager JVM Heap |
| 작은 상태로 테스트·실험에 적합 (로컬) |
상태 스냅샷 (State Snapshots)
정의 (Definitions)
- Snapshot – Flink job의 상태에 대한 전역적으로 일관된 이미지를 가리키는 일반 용어. 스냅샷은 각 데이터 소스에 대한 포인터(예: 파일 또는 Kafka 파티션의 오프셋)와, 소스의 그 위치까지 모든 이벤트를 처리한 결과인 job의 각 상태 저장 연산자의 상태 복사본을 포함해요.
- Checkpoint – Flink가 장애로부터 복구할 수 있도록 자동으로 찍는 스냅샷. 체크포인트는 증분적일 수 있고, 빠르게 복원되도록 최적화돼요.
- Externalized Checkpoint – 보통 체크포인트는 사용자가 조작하도록 의도되지 않아요. Flink는 job이 실행되는 동안 가장 최근의 n개 체크포인트(n은 구성 가능)만 유지하고, job이 취소되면 삭제해요. 하지만 이를 대신 유지하도록 구성할 수 있으며, 이 경우 수동으로 그로부터 재개할 수 있어요.
- Savepoint – 상태 있는 재배포/업그레이드/재스케일링 같은 운영 목적으로 사용자(또는 API 호출)가 수동으로 트리거하는 스냅샷. Savepoint는 항상 완전(complete)하며, 운영 유연성에 최적화돼요.
상태 스냅샷은 어떻게 동작하나? (How does State Snapshotting Work?)
Flink는 비동기 배리어 스냅샷(asynchronous barrier snapshotting)으로 알려진 Chandy-Lamport 알고리즘의 변형을 사용해요.
체크포인트 코디네이터(job manager의 일부)가 task manager에 체크포인트 시작을 지시하면, 모든 소스가 자신의 오프셋을 기록하고 번호가 매겨진 체크포인트 배리어를 스트림에 삽입해요. 이 배리어들은 job 그래프를 흘러가며 각 체크포인트 전후의 스트림 부분을 나타내요.
체크포인트 n은 체크포인트 배리어 n 이전의 모든 이벤트를 소비한 결과인 각 연산자의 상태를 포함하고, 그 이후의 이벤트는 포함하지 않아요.
job 그래프의 각 연산자가 이 배리어 중 하나를 받으면 자신의 상태를 기록해요. 두 입력 스트림을 가진 연산자(예: CoProcessFunction)는 배리어 정렬(barrier alignment)을 수행해, 스냅샷이 두 입력 스트림 모두에서 두 배리어까지(지나지 않고) 이벤트를 소비한 결과인 상태를 반영하도록 해요.
Flink의 상태 백엔드는 copy-on-write 메커니즘을 사용해, 상태의 이전 버전이 비동기적으로 스냅샷되는 동안 스트림 처리가 방해받지 않고 계속되게 해요. 스냅샷이 영구적으로 지속된 후에만 상태의 이전 버전이 가비지 컬렉션돼요.
Exactly Once 보장 (Exactly Once Guarantees)
스트림 처리 애플리케이션에서 문제가 발생하면 결과가 손실되거나 중복될 수 있어요. Flink에서는 애플리케이션과 클러스터에 대한 선택에 따라 다음 중 어떤 결과든 가능해요.
- Flink가 실패로부터 복구하려는 노력을 하지 않음 (at most once)
- 아무것도 손실되지 않지만 중복된 결과가 발생할 수 있음 (at least once)
- 손실도 중복도 없음 (exactly once)
Flink가 소스 데이터 스트림을 되감고 재생하여 장애에서 복구한다는 점을 감안할 때, 이상적인 상황이 exactly once로 설명될 때 이것이 모든 이벤트가 정확히 한 번 처리된다는 뜻은 아니에요. 대신, 모든 이벤트가 Flink가 관리하는 상태에 정확히 한 번 영향을 준다는 뜻이에요.
배리어 정렬은 exactly once 보장을 제공하는 데에만 필요해요. 이게 필요 없다면, CheckpointingMode.AT_LEAST_ONCE를 사용하도록 Flink를 구성해 성능을 얻을 수 있으며, 이는 배리어 정렬을 비활성화하는 효과가 있어요.
End-to-end Exactly Once
모든 소스의 이벤트가 sink에 정확히 한 번 영향을 주기 위해 end-to-end exactly once를 달성하려면 다음이 참이어야 해요.
- 소스가 재생 가능(replayable)해야 하고,
- sink가 트랜잭션적(transactional)이어야 (또는 멱등(idempotent)이어야) 해요.
Hands-on
Flink Operations Playground에는 Observing Failure & Recovery에 대한 절이 포함되어 있어요.
추가 읽기 (Further Reading)
- Stateful Stream Processing
- State Backends
- Fault Tolerance Guarantees of Data Sources and Sinks
- Enabling and Configuring Checkpointing
- Checkpoints
- Savepoints
- Tuning Checkpoints and Large State
- Monitoring Checkpointing
- Task Failure Recovery