State Backends
State Backends (상태 백엔드)
[Data Stream API]({{< ref "docs/dev/datastream/overview" >}})로 작성된 프로그램은 다양한 형태로 상태를 보유해요. 상태가 내부적으로 어떻게 표현되고 체크포인트에 어떻게, 어디에 저장되는지는 선택한 State Backend에 따라 달라져요.
출처: State Backends
본문
[Data Stream API]({{< ref "docs/dev/datastream/overview" >}})로 작성된 프로그램은 다양한 형태로 상태를 보유해요:
- 윈도우는 트리거될 때까지 요소나 집계를 모아요
- 변환 함수는 key/value 상태 인터페이스를 사용해 값을 저장할 수 있어요
- 변환 함수는
CheckpointedFunction인터페이스를 구현해 로컬 변수를 내결함성 있게 만들 수 있어요
스트리밍 API 가이드의 [state 섹션]({{< ref "docs/dev/datastream/fault-tolerance/state" >}})도 참고하세요.
체크포인팅이 활성화되면 그러한 상태는 데이터 손실을 막고 일관되게 복구하기 위해 체크포인트에 저장돼요. 상태가 내부적으로 어떻게 표현되고 체크포인트에 어떻게, 어디에 저장되는지는 선택한 State Backend에 따라 달라져요.
사용 가능한 State Backend
기본적으로 Flink는 다음 상태 백엔드를 번들로 제공해요:
- HashMapStateBackend
- EmbeddedRocksDBStateBackend
다른 것이 구성되지 않으면 시스템은 HashMapStateBackend을 사용해요.
HashMapStateBackend
HashMapStateBackend는 데이터를 Java 힙의 객체로 내부적으로 보유해요. Key/value 상태와 윈도우 연산자는 값, 트리거 등을 저장하는 해시 테이블을 보유해요.
HashMapStateBackend은 다음에 권장돼요:
- 큰 상태, 긴 윈도우, 큰 key/value 상태를 가진 잡.
- 모든 고가용성(HA) 구성.
또한 [managed memory]({{< ref "docs/deployment/memory/mem_setup_tm" >}}#managed-memory)를 0으로 설정하는 것이 권장돼요. 이렇게 하면 JVM의 사용자 코드에 최대 메모리 양이 할당되도록 보장돼요.
EmbeddedRocksDBStateBackend과 달리 HashMapStateBackend은 데이터를 힙의 객체로 저장하므로 객체를 재사용하는 것은 안전하지 않아요.
EmbeddedRocksDBStateBackend
EmbeddedRocksDBStateBackend는 (기본적으로) TaskManager 로컬 데이터 디렉터리에 저장되는 RocksDB 데이터베이스에 진행 중인 데이터를 보유해요. HashMapStateBackend처럼 Java 객체를 저장하는 것과 달리 데이터는 직렬화된 바이트 배열로 저장되는데, 주로 타입 직렬 변환기에 의해 정의되며, 키 비교는 Java의 hashCode()와 equals() 메서드를 사용하는 대신 바이트 단위로 수행돼요.
EmbeddedRocksDBStateBackend은 항상 비동기 스냅샷을 수행해요.
EmbeddedRocksDBStateBackend의 제한 사항:
- RocksDB의 JNI 브리지 API가 byte[] 기반이므로 키당 그리고 값당 최대 지원 크기는 각각 2^31 바이트예요. RocksDB에서 병합 연산을 사용하는 상태(예: ListState)는 값 크기가 2^31바이트를 초과해 조용히 누적될 수 있으며 다음 검색에서 실패해요. 이것은 현재 RocksDB JNI의 제한이에요.
EmbeddedRocksDBStateBackend은 다음에 권장돼요:
- 매우 큰 상태, 긴 윈도우, 큰 key/value 상태를 가진 잡.
- 모든 고가용성(HA) 구성.
유지할 수 있는 상태의 양은 사용 가능한 디스크 공간에 의해서만 제한된다는 점에 유의하세요. 이는 상태를 메모리에 유지하는 HashMapStateBackend와 비교해 매우 큰 상태를 유지할 수 있게 해줘요. 그러나 이는 이 상태 백엔드로 달성할 수 있는 최대 처리량이 더 낮다는 것을 의미하기도 해요. 이 백엔드로/백엔드에서의 모든 읽기/쓰기는 상태 객체를 검색/저장하기 위해 역직렬화/직렬화를 거쳐야 하며, 이는 힙 기반 백엔드처럼 항상 on-heap 표현으로 작업하는 것보다 더 비싸요. 역직렬화/직렬화로 인해 EmbeddedRocksDBStateBackend에서 객체를 재사용하는 것은 안전해요.
EmbeddedRocksDBStateBackend에 대한 [task executor 메모리 구성]({{< ref "docs/deployment/memory/mem_tuning" >}}#rocksdb-state-backend) 권장 사항도 확인하세요.
EmbeddedRocksDBStateBackend은 현재 증분 체크포인트를 제공하는 유일한 백엔드예요 ([여기]({{< ref "docs/ops/state/large_state_tuning" >}}) 참조).
일부 RocksDB 네이티브 메트릭은 사용할 수 있지만 기본적으로 비활성화되어 있어요. 전체 문서는 [여기]({{< ref "docs/deployment/config" >}}#rocksdb-native-metrics)에서 찾을 수 있어요.
슬롯당 RocksDB 인스턴스의 총 메모리 양도 제한할 수 있어요. 자세한 내용은 [여기]({{< ref "docs/ops/state/large_state_tuning" >}}#bounding-rocksdb-memory-usage) 문서를 참조하세요.
ForStStateBackend
ForStStateBackend는 ForSt 프로젝트에 기반한 상태 백엔드로, 이것 역시 RocksDB 위에 구축된 LSM-tree 구조의 key-value 저장소예요. 분리된 상태 관리(disaggregated state management)를 위해 설계되었으며, 자세한 내용은 [여기]({{< ref "docs/ops/state/disaggregated_state" >}})를 참조하세요. 가장 중요하게는 Flink가 지원하는 HDFS, S3 등 원격 파일 시스템에 sst 파일을 보유할 수 있어요. 이는 Flink가 상태 크기를 TaskManager의 로컬 디스크 용량을 넘어 확장할 수 있게 해줘요. 게다가 sst 파일을 원격 파일 시스템에 두면 체크포인트와 복구를 수행하는 더 가벼운 방법도 제공해요.
ForStStateBackend은 아직 실험 단계이며 프로덕션에 완전히 사용할 수는 없어요. 항상 비동기 증분 스냅샷을 수행해요.
ForStStateBackend은 다음에 권장돼요:
- 매우 큰 상태, 긴 윈도우, 큰 key/value 상태를 가진 잡. 로컬 디스크가 상태를 저장하기에 충분하지 않을 수 있어요.
- 모든 고가용성(HA) 구성.
- 비동기 상태 접근이 선호되는 경우. ForStStateBackend이 비동기 상태 접근을 지원하는 유일한 백엔드이기 때문이에요.
- 클라우드 네이티브 애플리케이션처럼 가벼운 체크포인트와 복구가 필요한 잡.
ForStStateBackend의 (현재) 제한 사항:
- EmbeddedRocksDBStateBackend과 마찬가지로 키당 그리고 값당 최대 지원 크기는 각각 2^31바이트예요.
- canonical savepoint, 전체 스냅샷, changelog, file-merging 체크포인트를 지원하지 않아요. 항상 증분 스냅샷을 수행해요.
EmbeddedRocksDBStateBackend과 비교해 ForStStateBackend은 데이터를 원격 파일 시스템에 저장하므로 유지할 수 있는 상태의 양은 무제한이에요. TaskManager의 로컬 디스크는 더 나은 성능을 위해 파일 캐시를 저장하는 데만 사용돼요. 대부분의 활성 상태가 원격 파일 시스템에 있을 때 상태 접근 성능은 네트워크 지연의 영향을 받을 수 있다는 점에 유의하세요. Flink는 이 문제를 완화하기 위해 비동기 상태 접근을 도입해요. State API V2에서 비동기 상태 메서드를 사용한다면 비동기 상태 접근의 이점을 누릴 수 있어요. State API V2에 익숙해지려면 [State API V2 문서]({{< ref "docs/dev/datastream/fault-tolerance/state_v2" >}})를 참조하세요.
올바른 State Backend 선택
HashMapStateBackend와 RocksDB 사이에서 결정할 때는 성능과 확장성 사이의 선택이에요. HashMapStateBackend은 각 상태 접근과 갱신이 Java 힙의 객체에 대해 동작하므로 매우 빠르지만, 상태 크기는 클러스터 내 사용 가능한 메모리에 의해 제한돼요. 반면 RocksDB는 사용 가능한 디스크 공간에 따라 확장될 수 있어요. 그러나 각 상태 접근과 갱신은 (역)직렬화와 잠재적으로 디스크 읽기를 필요로 하여 메모리 상태 백엔드보다 한 자릿수 느린 평균 성능을 초래해요. 사용 가능한 디스크 공간을 초과하는 매우 큰 상태를 다루거나 클라우드 네이티브 구성에서 빠른 재확장을 선호한다면 ForStStateBackend을 고려해야 해요.
{{< hint info >}} Flink 1.13에서 savepoint의 바이너리 형식을 통합했어요. 즉 savepoint를 찍은 다음 다른 상태 백엔드로 복원할 수 있어요. 모든 상태 백엔드는 1.13 버전부터 공통 형식을 생성해요. 따라서 상태 백엔드를 전환하려면 먼저 Flink 버전을 업그레이드한 다음 새 버전으로 savepoint를 찍고, 그 후에만 다른 상태 백엔드로 복원할 수 있어요. {{< /hint >}}
State Backend 구성
아무것도 지정하지 않으면 기본 상태 백엔드는 HashMapStateBackend이에요. 클러스터의 모든 잡에 대해 다른 기본값을 설정하려면 [Flink 구성 파일]({{< ref "docs/deployment/config#flink-configuration-file" >}})에서 새 기본 상태 백엔드를 정의하면 돼요. 기본 상태 백엔드는 아래처럼 잡별로 재정의할 수 있어요.
잡별 State Backend 설정
잡별 상태 백엔드는 아래 예제처럼 잡의 StreamExecutionEnvironment에 설정돼요:
Configuration config = new Configuration();
config.set(StateBackendOptions.STATE_BACKEND, "hashmap");
env.configure(config);
config = Configuration()
config.set_string('state.backend.type', 'hashmap')
env = StreamExecutionEnvironment.get_execution_environment(config)
IDE에서 EmbeddedRocksDBStateBackend을 사용하거나 Flink 잡에서 프로그래밍 방식으로 구성하려면 Flink 프로젝트에 다음 종속성을 추가해야 해요.
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-statebackend-rocksdb</artifactId>
<version>{{< version >}}</version>
<scope>provided</scope>
</dependency>
ForStStateBackend도 동일해요:
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-statebackend-forst</artifactId>
<version>{{< version >}}</version>
<scope>provided</scope>
</dependency>
{{< hint info >}}
RocksDB와 ForSt는 기본 Flink 배포의 일부이므로, 잡에서 RocksDB 코드를 사용하지 않고 [Flink 구성 파일]({{< ref "docs/deployment/config#flink-configuration-file" >}})에서 state.backend.type과 추가 [체크포인팅]({{< ref "docs/deployment/config" >}}#checkpointing) 및 [RocksDB 특정]({{< ref "docs/deployment/config" >}}#rocksdb-state-backend) 또는 [ForSt 특정]({{< ref "docs/deployment/config" >}}#forst-state-backend) 매개변수로 상태 백엔드를 구성한다면 이 종속성이 필요 없어요.
{{< /hint >}}
기본 State Backend 설정
[Flink 구성 파일]({{< ref "docs/deployment/config#flink-configuration-file" >}})에서 state.backend.type 구성 키를 사용해 기본 상태 백엔드를 구성할 수 있어요.
구성 항목에 대한 가능한 값은 hashmap (HashMapStateBackend), rocksdb (EmbeddedRocksDBStateBackend), forst (ForStStateBackend) 또는 상태 백엔드 팩토리 {{< gh_link file="flink-runtime/src/main/java/org/apache/flink/runtime/state/StateBackendFactory.java" name="StateBackendFactory" >}}를 구현하는 클래스의 정규화된 클래스 이름이에요. 예를 들어 EmbeddedRocksDBStateBackend은 org.apache.flink.state.rocksdb.EmbeddedRocksDBStateBackendFactory, ForStStateBackend은 org.apache.flink.state.forst.ForStStateBackendFactory와 같이요.
execution.checkpointing.dir 옵션은 모든 백엔드가 체크포인트 데이터와 메타데이터 파일을 쓰는 디렉터리를 정의해요. 체크포인트 디렉터리 구조에 대한 자세한 내용은 [여기]({{< ref "docs/ops/state/checkpoints" >}}#directory-structure)에서 확인할 수 있어요.
구성 파일의 샘플 섹션은 다음과 같을 수 있어요:
# The backend that will be used to store operator state checkpoints
state.backend: hashmap
# Directory for storing checkpoints
execution.checkpointing.dir: hdfs://namenode:40010/flink/checkpoints
RocksDB State Backend 세부 사항
이 섹션은 RocksDB 상태 백엔드를 더 자세히 설명해요.
증분 체크포인트 (Incremental Checkpoints)
RocksDB는 전체 체크포인트에 비해 체크포인트 시간을 극적으로 줄일 수 있는 증분 체크포인트를 지원해요. 상태 백엔드의 완전하고 자족적인 백업을 생성하는 대신 증분 체크포인트는 마지막 완료된 체크포인트 이후 발생한 변경 사항만 기록해요.
증분 체크포인트는 (일반적으로 여러 개의) 이전 체크포인트를 기반으로 구축돼요. Flink는 시간이 지남에 따라 자체적으로 통합되는 방식으로 RocksDB의 내부 압축 메커니즘을 활용해요. 결과적으로 Flink의 증분 체크포인트 기록은 무한정 증가하지 않으며, 오래된 체크포인트는 결국 자동으로 흡수되고 정리돼요.
증분 체크포인트의 복구 시간은 전체 체크포인트와 비교해 더 길거나 짧을 수 있어요. 네트워크 대역폭이 병목이라면 더 많은 데이터(더 많은 델타)를 가져오는 것을 의미하므로 증분 체크포인트에서 복원하는 데 조금 더 시간이 걸릴 수 있어요. CPU나 IOPs가 병목이라면 증분 체크포인트에서 복원하는 것이 더 빠른데, 로컬 RocksDB 테이블을 Flink의 canonical key/value 스냅샷 형식(savepoint와 전체 체크포인트에서 사용됨)으로 다시 구축하지 않아도 되기 때문이에요.
큰 상태에 대해 증분 체크포인트 사용을 권장하지만, 이 기능을 수동으로 활성화해야 해요:
- [Flink 구성 파일]({{< ref "docs/deployment/config#flink-configuration-file" >}})에서 기본값 설정:
execution.checkpointing.incremental: true는 애플리케이션이 코드에서 이 설정을 재정의하지 않는 한 증분 체크포인트를 활성화해요. - 코드에서 직접 구성할 수도 있어요(구성 기본값 재정의):
EmbeddedRocksDBStateBackend backend = new EmbeddedRocksDBStateBackend(true);
증분 체크포인트가 활성화되면 웹 UI에 표시되는 Checkpointed Data Size는 전체 상태 크기가 아니라 해당 체크포인트의 델타 체크포인트 데이터 크기만 나타낸다는 점에 유의하세요.
메모리 관리 (Memory Management)
Flink는 총 프로세스 메모리 소비를 제어하여 Flink TaskManager가 잘 동작하는 메모리 사용량을 가지도록 하는 것을 목표로 해요. 이는 환경(Docker/Kubernetes, Yarn 등)이 강제한 한도 내에 머물러 너무 많은 메모리를 소비한다고 종료되지 않도록 하면서도, 메모리를 과소 활용하지도 않는 것(불필요한 디스크 스필링, 낭비되는 캐싱 기회, 성능 저하)을 의미해요.
이를 달성하기 위해 Flink는 기본적으로 RocksDB의 메모리 할당을 TaskManager(더 정확히는 task slot)의 managed memory 양으로 구성해요. 이는 대부분의 애플리케이션에 좋은 out-of-the-box 경험을 주어야 하며, 대부분의 애플리케이션은 세부 RocksDB 설정을 조정할 필요가 없어야 해요. 메모리 관련 성능 문제를 개선하는 주요 메커니즘은 Flink의 managed memory를 단순히 늘리는 것이에요.
사용자는 이 기능을 비활성화하고 RocksDB가 ColumnFamily별로(연산자당 상태당 하나) 메모리를 독립적으로 할당하도록 할 수 있어요. 이는 전문 사용자에게 RocksDB에 대한 궁극적으로 더 세밀한 제어를 제공하지만, 사용자가 전체 메모리 소비가 환경의 한도를 초과하지 않도록 스스로 관리해야 함을 의미해요. 큰 상태 성능 튜닝에 대한 지침은 [large state tuning]({{< ref "docs/ops/state/large_state_tuning" >}}#tuning-rocksdb-memory)을 참조하세요.
RocksDB용 Managed Memory
이 기능은 기본적으로 활성화되어 있으며 state.backend.rocksdb.memory.managed 구성 키로 (비)활성화할 수 있어요.
Flink는 RocksDB의 네이티브 메모리 할당을 직접 관리하지 않지만, RocksDB가 Flink가 managed memory 예산으로 가진 메모리만큼 정확히 사용하도록 특정 방식으로 구성해요. 이는 슬롯 수준에서 수행돼요 (managed memory는 슬롯별로 계산됨).
RocksDB 인스턴스의 총 메모리 사용량을 설정하기 위해 Flink는 단일 슬롯의 모든 인스턴스 간에 공유되는 cache와 write buffer manager를 활용해요. 공유 캐시는 RocksDB에서 대부분의 메모리를 사용하는 세 구성 요소인 블록 캐시, 인덱스와 블룸 필터, 그리고 MemTable에 상한을 둬요.
고급 튜닝을 위해 Flink는 쓰기 경로(MemTable)와 읽기 경로(인덱스 & 필터, 나머지 캐시) 사이의 메모리 분할을 제어하는 두 매개변수도 제공해요. write buffer 메모리 부족(빈번한 플러시)이나 캐시 미스로 RocksDB가 성능이 나쁘다고 보인다면 이 매개변수를 사용해 메모리를 재분배할 수 있어요.
state.backend.rocksdb.memory.write-buffer-ratio, 기본값0.5는 주어진 메모리의 50%가 write buffer manager에 사용된다는 뜻이에요.state.backend.rocksdb.memory.high-prio-pool-ratio, 기본값0.1은 주어진 메모리의 10%가 공유 블록 캐시에서 인덱스와 필터에 대한 우선순위로 설정된다는 뜻이에요. 인덱스와 필터가 데이터 블록과 캐시 유지를 경쟁하여 성능 문제를 일으키지 않도록 이 값을 0으로 설정하지 않는 것을 강력히 권장해요. 게다가 L0 레벨 필터와 인덱스는 성능 문제를 완화하기 위해 기본적으로 캐시에 고정돼요. 자세한 내용은 RocksDB 문헌을 참조하세요.
{{< hint info >}}
위에서 설명한 메커니즘 (cache와 write buffer manager)이 활성화되면 PredefinedOptions와 RocksDBOptionsFactory를 통해 블록 캐시와 쓰기 버퍼에 대해 이루어진 사용자 지정 설정을 재정의해요.
{{< /hint >}}
{{< details "Expert Mode" >}}
메모리를 수동으로 제어하려면 state.backend.rocksdb.memory.managed를 false로 설정하고 ColumnFamilyOptions를 통해 RocksDB를 구성할 수 있어요. 또는 위에서 언급한 cache/buffer-manager 메커니즘을 사용하되 메모리 크기를 Flink의 managed memory 크기와 무관한 고정 양으로 설정할 수 있어요(state.backend.rocksdb.memory.fixed-per-slot 또는 state.backend.rocksdb.memory.fixed-per-tm 옵션). 두 경우 모두 JVM 외부에 RocksDB용 충분한 메모리가 있는지 사용자가 직접 확인해야 해요.
{{< /details >}}
타이머 (Heap vs. RocksDB)
타이머는 나중에(이벤트 시간 또는 처리 시간) 작업을 예약하는 데 사용돼요. 예를 들어 윈도우를 실행하거나 ProcessFunction을 콜백하는 데요.
RocksDB 상태 백엔드를 선택하면 타이머는 기본적으로도 RocksDB에 저장돼요. 이는 애플리케이션이 많은 타이머로 확장할 수 있게 해주는 견고하고 확장 가능한 방식이에요. 그러나 RocksDB에서 타이머를 유지하는 것은 일정 비용이 들 수 있으며, 그래서 Flink는 다른 상태를 저장하는 데 RocksDB를 사용하더라도 타이머를 JVM 힙에 저장하는 옵션을 제공해요. 힙 기반 타이머는 타이머 수가 적을 때 더 나은 성능을 가질 수 있어요.
구성 옵션 state.backend.rocksdb.timer-service.factory를 (기본값인 rocksdb 대신) heap으로 설정하여 타이머를 힙에 저장하세요.
{{< hint warning >}} 힙 기반 타이머가 있는 RocksDB 상태 백엔드 조합은 현재 타이머 상태에 대해 비동기 스냅샷을 지원하지 않아요. keyed state 같은 다른 상태는 여전히 비동기로 스냅샷됩니다. {{< /hint >}}
{{< hint warning >}} 힙 기반 타이머로 RocksDB 상태 백엔드를 사용할 때, 애플리케이션에 raw keyed state에 쓰는 연산자가 있으면 체크포인팅과 savepoint 찍기가 실패할 것으로 예상돼요. 이는 커스텀 스트림 연산자를 작성하는 고급 사용자에게만 관련이 있어요. {{< /hint >}}
RocksDB 네이티브 메트릭 활성화
특정 메트릭을 선택적으로 활성화하여 Flink의 메트릭 시스템을 통해 RockDB의 네이티브 메트릭에 접근할 수 있어요. 자세한 내용은 [구성 문서]({{< ref "docs/deployment/config" >}}#rocksdb-native-metrics)를 참조하세요.
{{< hint warning >}} RocksDB 네이티브 메트릭을 활성화하면 애플리케이션에 부정적인 성능 영향이 있을 수 있어요. {{< /hint >}}
고급 RocksDB 메모리 튜닝
{{< hint info >}} Flink는 대부분의 사용 사례에 작동해야 하는 정교한 기본 RocksDB 메모리 관리를 제공해요. 아래 메커니즘은 주로 전문 튜닝이나 문제 해결에 사용되어야 해요. {{< /hint >}}
사전 정의된 Per-ColumnFamily 옵션
Predefined Options로 사용자는 각 RocksDB Column Family에 일부 사전 정의 구성 프로필을 적용할 수 있어요. 예를 들어 메모리 사용, 스레드, 압축 설정 등을 구성할 수 있어요. 현재 각 연산자의 각 상태에 대해 하나의 Column Family가 있어요.
적용할 사전 정의 옵션을 선택하는 두 가지 방법이 있어요:
- [Flink 구성 파일]({{< ref "docs/deployment/config#flink-configuration-file" >}})에서
state.backend.rocksdb.predefined-options를 통해 옵션 이름을 설정해요. - 프로그래밍 방식으로 사전 정의 옵션 설정:
EmbeddedRocksDBStateBackend.setPredefinedOptions(PredefinedOptions.SPINNING_DISK_OPTIMIZED_HIGH_MEM).
이 옵션의 기본값은 DEFAULT이며 PredefinedOptions.DEFAULT로 변환돼요.
프로그래밍 방식으로 설정된 사전 정의 옵션은 [Flink 구성 파일]({{< ref "docs/deployment/config#flink-configuration-file" >}})을 통해 구성된 것을 재정의해요.
Flink 구성 파일에서 Column Family 옵션 읽기
RocksDB 상태 백엔드는 [여기 정의된]({{< ref "docs/deployment/config" >}}#advanced-rocksdb-state-backends-options) 모든 구성 옵션을 읽어요. 따라서 RocksDB에서 managed memory를 끄고 구성에 관련 항목을 넣어 저수준 Column Family 옵션을 간단히 구성할 수 있어요.
Options Factory를 RocksDB에 전달
RocksDB의 옵션을 수동으로 제어하려면 RocksDBOptionsFactory를 구성해야 해요. 이 메커니즘은 Column Family 설정(예: 메모리 사용, 스레드, 압축 설정)을 세밀하게 제어할 수 있게 해줘요. 현재 각 연산자의 각 상태에 대해 하나의 Column Family가 있어요.
RocksDB 상태 백엔드에 RocksDBOptionsFactory를 전달하는 두 가지 방법이 있어요:
- [Flink 구성 파일]({{< ref "docs/deployment/config#flink-configuration-file" >}})에서
state.backend.rocksdb.options-factory를 통해 옵션 팩토리 클래스 이름을 구성해요. - 프로그래밍 방식으로 옵션 팩토리 설정, 예:
EmbeddedRocksDBStateBackend.setRocksDBOptions(new MyOptionsFactory());
프로그래밍 방식으로 설정된 옵션 팩토리는 [Flink 구성 파일]({{< ref "docs/deployment/config#flink-configuration-file" >}})을 통해 구성된 것을 재정의하고, 옵션 팩토리는 설정된 경우 사전 정의 옵션보다 우선순위가 높아요.
RocksDB는 프로세스에서 직접 메모리를 할당하는 네이티브 라이브러리이며, JVM에서 할당하지 않아요. RocksDB에 할당하는 모든 메모리는 일반적으로 TaskManager의 JVM 힙 크기를 같은 양만큼 줄여 회계 처리해야 해요. 그렇게 하지 않으면 YARN 등이 구성된 것보다 더 많은 메모리를 할당했다며 JVM 프로세스를 종료할 수 있어요.
아래는 커스텀 ConfigurableOptionsFactory를 정의하는 예제예요 (state.backend.rocksdb.options-factory 아래에 클래스 이름 설정).
public class MyOptionsFactory implements ConfigurableRocksDBOptionsFactory {
public static final ConfigOption<Integer> BLOCK_RESTART_INTERVAL = ConfigOptions
.key("my.custom.rocksdb.block.restart-interval")
.intType()
.defaultValue(16)
.withDescription(
" Block restart interval. RocksDB has default block restart interval as 16. ");
private int blockRestartInterval = BLOCK_RESTART_INTERVAL.defaultValue();
@Override
public DBOptions createDBOptions(DBOptions currentOptions,
Collection<AutoCloseable> handlesToClose) {
return currentOptions
.setIncreaseParallelism(4)
.setUseFsync(false);
}
@Override
public ColumnFamilyOptions createColumnOptions(ColumnFamilyOptions currentOptions,
Collection<AutoCloseable> handlesToClose) {
return currentOptions.setTableFormatConfig(
new BlockBasedTableConfig()
.setBlockRestartInterval(blockRestartInterval));
}
@Override
public RocksDBOptionsFactory configure(ReadableConfig configuration) {
this.blockRestartInterval = configuration.get(BLOCK_RESTART_INTERVAL);
return this;
}
}
Still not supported in Python API.
Changelog 활성화
소개
Changelog는 체크포인트 시간과 그에 따른 exactly-once 모드의 end-to-end 지연을 줄이는 것을 목표로 하는 기능이에요.
가장 일반적으로 체크포인트 기간은 다음에 의해 영향을 받아요:
- 배리어 이동 시간과 정렬 — [Unaligned checkpoints]({{< ref "docs/ops/state/checkpointing_under_backpressure#unaligned-checkpoints" >}})와 [Buffer debloating]({{< ref "docs/ops/state/checkpointing_under_backpressure#buffer-debloating" >}})으로 해결
- 스냅샷 생성 시간(소위 동기 단계) — 비동기 스냅샷으로 해결 (위에서 언급)
- 스냅샷 업로드 시간(비동기 단계)
업로드 시간은 [증분 체크포인트]({{< ref "#incremental-checkpoints" >}})로 줄일 수 있어요. 그러나 대부분의 증분 상태 백엔드는 주기적으로 일부 형태의 압축을 수행하여 새 변경 사항 외에도 이전 상태를 재업로드하게 돼요. 대규모 배포에서는 매 체크포인트마다 적어도 하나의 태스크가 많은 데이터를 업로드할 확률이 매우 높은 경향이 있어요.
Changelog를 활성화하면 Flink는 상태 변경을 지속적으로 업로드하여 changelog를 형성해요. 체크포인트에서 이 changelog의 관련 부분만 업로드하면 돼요. 구성된 상태 백엔드는 배경에서 주기적으로 스냅샷됩니다. 업로드가 성공하면 changelog가 잘립니다.
결과적으로 비동기 단계 기간이 줄고, 디스크에 플러시할 데이터가 없기 때문에 동기 단계도 줄어들어요. 특히 long-tail 지연이 개선돼요. 동시에 몇 가지 다른 이점도 얻을 수 있어요:
- 더 안정적이고 낮은 end-to-end 지연.
- 장애 복구 후 데이터 재생(replay) 감소.
- 리소스 사용률의 더 안정적 활용.
그러나 리소스 사용량은 더 높아져요:
- DFS에 더 많은 파일이 생성됨
- 상태 변경을 업로드하는 데 더 많은 IO 대역폭이 사용됨
- 상태 변경을 직렬화하는 데 더 많은 CPU가 사용됨
- 상태 변경을 버퍼링하기 위해 Task Manager가 더 많은 메모리를 사용함
changelog가 소량의 일일 CPU 및 네트워크 대역폭 리소스를 추가하지만, 피크 CPU 및 네트워크 대역폭 사용량을 줄인다는 점은 주목할 가치가 있어요.
복구 시간도 고려해야 할 또 다른 사항이에요. state.changelog.periodic-materialize.interval 설정에 따라 changelog가 길어질 수 있고 재생하는 데 더 많은 시간이 걸릴 수 있어요. 그러나 체크포인트 기간과 결합한 복구 시간은 changelog가 없는 구성보다 여전히 낮을 가능성이 높으며, 장애 조치 경우에도 더 낮은 end-to-end 지연을 제공해요. 그러나 앞서 언급한 시간의 실제 비율에 따라 유효 복구 시간이 증가할 수도 있어요.
자세한 내용은 FLIP-158을 참조하세요.
설치
Changelog JAR은 표준 Flink 배포에 포함되어 있어요.
필요한 파일시스템 플러그인을 [추가]({{< ref "docs/deployment/filesystems/overview" >}})했는지 확인하세요.
구성
다음은 YAML의 예시 구성이에요:
state.changelog.enabled: true
state.changelog.storage: filesystem # currently, only filesystem and memory (for tests) are supported
state.changelog.dstl.dfs.base-path: s3://<bucket-name> # similar to execution.checkpointing.dir
다음 기본값을 유지하세요 (제한 사항 참조):
execution.checkpointing.max-concurrent-checkpoints: 1
다른 옵션은 [구성 섹션]({{< ref "docs/deployment/config#state-changelog-options" >}})을 참조하세요.
Changelog는 잡별로 프로그래밍 방식으로 활성화/비활성화할 수도 있어요:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableChangelogStateBackend(true);
env = StreamExecutionEnvironment.get_execution_environment()
env.enable_changelog_statebackend(true)
모니터링
사용 가능한 메트릭은 [여기]({{< ref "docs/ops/metrics#state-changelog" >}})에 나열되어 있어요.
태스크가 상태 변경을 쓰는 것으로 인해 지연(backpressure)되면 UI에서 busy(빨간색)로 표시돼요.
기존 잡 업그레이드
Changelog 활성화
savepoint와 checkpoint 양쪽 모두에서 재개가 지원돼요:
- 기존의 non-changelog 잡이 있다고 가정
- [savepoint]({{< ref "docs/ops/state/savepoints#resuming-from-savepoints" >}})나 [checkpoint]({{< ref "docs/ops/state/checkpoints#resuming-from-a-retained-checkpoint" >}}) 중 하나를 찍음
- 구성 변경(Changelog 활성화)
- 찍은 스냅샷에서 재개
Changelog 비활성화
savepoint와 checkpoint 양쪽 모두에서 재개가 지원돼요:
- 기존의 changelog 잡이 있다고 가정
- [savepoint]({{< ref "docs/ops/state/savepoints#resuming-from-savepoints" >}})나 [checkpoint]({{< ref "docs/ops/state/checkpoints#resuming-from-a-retained-checkpoint" >}}) 중 하나를 찍음
- 구성 변경(Changelog 비활성화)
- 찍은 스냅샷에서 재개
제한 사항
- 최대 하나의 동시 체크포인트만 허용
- Flink 1.15 기준
filesystemchangelog 구현만 사용 가능 - [NO_CLAIM]({{< ref "docs/deployment/config#execution-savepoint-restore-mode" >}}) 모드 지원 안 됨
레거시 백엔드에서 마이그레이션
Flink 1.13부터 커뮤니티는 로컬 상태 저장과 체크포인트 저장의 분리를 사용자가 더 잘 이해할 수 있도록 공개 상태 백엔드 클래스를 재작업했어요. 이 변경은 Flink 상태 백엔드나 체크포인트 프로세스의 런타임 구현이나 특성에 영향을 주지 않아요. 단지 의도를 더 잘 전달하기 위한 것이에요. 사용자는 상태나 일관성을 잃지 않고 기존 애플리케이션을 새 API로 마이그레이션할 수 있어요.
MemoryStateBackend
레거시 MemoryStateBackend은 HashMapStateBackend와 [JobManagerCheckpointStorage]({{< ref "docs/ops/state/checkpoints#the-jobmanagercheckpointstorage" >}})를 사용하는 것과 동일해요.
Flink 구성 파일을 사용한 구성
state.backend: hashmap
# Optional, Flink will automatically default to JobManagerCheckpointStorage
# when no checkpoint directory is specified.
execution.checkpointing.storage: jobmanager
코드 구성
Configuration config = new Configuration();
config.set(StateBackendOptions.STATE_BACKEND, "hashmap");
config.set(CheckpointingOptions.CHECKPOINT_STORAGE, "jobmanager");
env.configure(config);
config = Configuration()
config.set_string('state.backend.type', 'hashmap')
config.set_string('execution.checkpointing.storage', 'jobmanager')
env = StreamExecutionEnvironment.get_execution_environment(config)
FsStateBackend
레거시 FsStateBackend은 HashMapStateBackend와 [FileSystemCheckpointStorage]({{< ref "docs/ops/state/checkpoints#the-filesystemcheckpointstorage" >}})를 사용하는 것과 동일해요.
Flink 구성 파일을 사용한 구성
state.backend: hashmap
execution.checkpointing.dir: file:///checkpoint-dir/
# Optional, Flink will automatically default to FileSystemCheckpointStorage
# when a checkpoint directory is specified.
execution.checkpointing.storage: filesystem
코드 구성
Configuration config = new Configuration();
config.set(StateBackendOptions.STATE_BACKEND, "hashmap");
config.set(CheckpointingOptions.CHECKPOINT_STORAGE, "filesystem");
config.set(CheckpointingOptions.CHECKPOINTS_DIRECTORY, "file:///checkpoint-dir");
env.configure(config);
// Advanced FsStateBackend configurations, such as write buffer size
// can be set manually by using CheckpointingOptions.
config.set(CheckpointingOptions.FS_WRITE_BUFFER_SIZE, 4 * 1024);
env.configure(config);
config = Configuration()
config.set_string('state.backend.type', 'hashmap')
config.set_string('execution.checkpointing.storage', 'filesystem')
config.set_string('execution.checkpointing.dir', 'file:///checkpoint-dir')
env = StreamExecutionEnvironment.get_execution_environment(config)
# Advanced FsStateBackend configurations, such as write buffer size
# can be set manually by using CheckpointingOptions.
config.set_string('state.storage.fs.write-buffer-size', '4096');
env.configure(config);
RocksDBStateBackend
레거시 RocksDBStateBackend은 EmbeddedRocksDBStateBackend와 [FileSystemCheckpointStorage]({{< ref "docs/ops/state/checkpoints#the-filesystemcheckpointstorage" >}})를 사용하는 것과 동일해요.
Flink 구성 파일을 사용한 구성
state.backend: rocksdb
execution.checkpointing.dir: file:///checkpoint-dir/
# Optional, Flink will automatically default to FileSystemCheckpointStorage
# when a checkpoint directory is specified.
execution.checkpointing.storage: filesystem
코드 구성
Configuration config = new Configuration();
config.set(StateBackendOptions.STATE_BACKEND, "rocksdb");
config.set(CheckpointingOptions.CHECKPOINT_STORAGE, "filesystem");
config.set(CheckpointingOptions.CHECKPOINTS_DIRECTORY, "file:///checkpoint-dir");
env.configure(config);
// If you manually passed FsStateBackend into the RocksDBStateBackend constructor
// to specify advanced checkpointing configurations such as write buffer size,
// you can achieve the same results by using CheckpointingOptions.
config.set(CheckpointingOptions.FS_WRITE_BUFFER_SIZE, 4 * 1024);
env.configure(config);
config = Configuration()
config.set_string('state.backend.type', 'rocksdb')
config.set_string('execution.checkpointing.storage', 'filesystem')
config.set_string('execution.checkpointing.dir', 'file:///checkpoint-dir')
env = StreamExecutionEnvironment.get_execution_environment(config)
# If you manually passed FsStateBackend into the RocksDBStateBackend constructor
# to specify advanced checkpointing configurations such as write buffer size,
# you can achieve the same results by using CheckpointingOptions.
config.set_string('state.storage.fs.write-buffer-size', '4096');
env.configure(config);