체크포인트와 대규모 상태 튜닝
체크포인트와 대규모 상태 튜닝 (Tuning Checkpoints and Large State)
이 페이지는 대규모 상태를 사용하는 애플리케이션을 구성하고 튜닝하는 방법에 대한 가이드를 제공해요.
출처: 문서
본문
개요 (Overview)
Flink 애플리케이션이 대규모로 안정적으로 실행되려면 두 조건이 충족되어야 해요:
- 애플리케이션이 체크포인트를 안정적으로 찍을 수 있어야 함
- 리소스가 실패 후 입력 데이터 스트림을 따라잡기에 충분해야 함
첫 번째 섹션들은 대규모에서 성능 좋은 체크포인트를 얻는 방법을 논의해요. 마지막 섹션은 얼마나 많은 리소스를 사용할지 계획하는 것에 관한 몇 가지 모범 사례를 설명해요.
상태와 체크포인트 모니터링 (Monitoring State and Checkpoints)
체크포인트 동작을 모니터링하는 가장 쉬운 방법은 UI의 체크포인트 섹션이에요. 체크포인트 모니터링 문서는 사용 가능한 체크포인트 메트릭에 접근하는 방법을 보여줘요.
체크포인트를 확장할 때 특히 관심 있는 두 숫자(Task 수준 메트릭과 웹 인터페이스 모두로 노출)는:
- 연산자가 첫 번째 체크포인트 배리어를 받을 때까지의 시간. 체크포인트를 트리거할 시간이 지속적으로 매우 높으면 체크포인트 배리어가 소스에서 연산자로 이동하는 데 오랜 시간이 걸린다는 뜻이에요. 이는 일반적으로 시스템이 지속적인 백프레셔 하에서 작동하고 있음을 나타내요.
- 정렬(alignment) 기간. 첫 번째와 마지막 체크포인트 배리어를 받는 사이의 시간으로 정의돼요. 정렬되지 않은 exactly-once 체크포인트와 at-least-once 체크포인트 동안에는 하위 작업들이 중단 없이 업스트림 하위 작업의 모든 데이터를 처리해요. 그러나 정렬된(aligned) exactly-once 체크포인트에서는 이미 체크포인트 배리어를 받은 채널이 나머지 모든 채널이 따라잡아 자신의 체크포인트 배리어를 받을 때까지 추가 데이터 전송이 차단돼요(정렬 시간).
이 두 값 모두 이상적으로 낮아야 해요. 값이 높다는 것은 체크포인트 배리어가 백프레셔(들어오는 레코드를 처리할 리소스 부족) 때문에 작업 그래프를 천천히 이동한다는 뜻이에요. 이는 처리된 레코드의 엔드투엔드 지연 증가로도 관찰할 수 있어요. 이 숫자들은 일시적 백프레셔, 데이터 치우침, 네트워크 문제가 있을 때 가끔 높을 수 있다는 점에 주의해요.
정렬되지 않은 체크포인트를 사용해 체크포인트 배리어의 전파 시간을 빠르게 할 수 있어요. 하지만 이는 애초에 백프레셔를 유발하는 근본 문제를 해결하지 못한다는 점에 유의하세요(엔드투엔드 레코드 지연은 높게 유지됨).
체크포인팅 튜닝 (Tuning Checkpointing)
체크포인트는 애플리케이션이 구성할 수 있는 규칙적인 간격으로 트리거돼요. 체크포인트 완료가 체크포인트 간격보다 오래 걸리면 진행 중인 체크포인트가 완료되기 전에는 다음 체크포인트가 트리거되지 않아요. 기본적으로 진행 중인 체크포인트가 완료되면 다음 체크포인트가 즉시 트리거돼요.
체크포인트가 기본 간격보다 길어지는 일이 자주 발생하면(예: 상태가 계획보다 커지거나 체크포인트 저장소가 일시적으로 느려져서) 시스템이 지속적으로 체크포인트를 찍어요(진행 중 것이 끝나면 즉시 새 것이 시작). 이는 너무 많은 리소스가 체크포인팅에 묶여 연산자가 너무 적게 진행할 수 있음을 의미해요. 이 동작은 비동기적으로 체크포인트되는 상태를 사용하는 스트리밍 애플리케이션에 미치는 영향은 적지만 전체 애플리케이션 성능에는 여전히 영향을 줄 수 있어요.
이런 상황을 막기 위해 애플리케이션은 체크포인트 사이의 최소 기간을 정의할 수 있어요:
StreamExecutionEnvironment.getCheckpointConfig().setMinPauseBetweenCheckpoints(milliseconds)
이 기간은 최신 체크포인트 끝과 다음 시작 사이에 반드시 경과해야 하는 최소 시간 간격이에요. 아래 그림은 이것이 체크포인팅에 미치는 영향을 보여줘요.
참고: 애플리케이션은 CheckpointConfig를 통해 여러 체크포인트가 동시에 진행 중일 수 있도록 구성할 수 있어요. Flink에서 대규모 상태를 가진 애플리케이션의 경우 이는 종종 너무 많은 리소스를 체크포인팅에 묶어요. savepoint가 수동으로 트리거되면 진행 중인 체크포인트와 동시에 진행될 수 있어요.
RocksDB 또는 ForSt 튜닝 (Tuning RocksDB or ForSt)
많은 대규모 Flink 스트리밍 애플리케이션의 상태 저장 '일꾼'은 RocksDB State Backend예요. 이 백엔드는 주 메모리를 훨씬 넘어 확장되며 대규모 keyed state를 안정적으로 저장해요.
TaskManager의 로컬 디스크 공간을 초과할 정도로 매우 큰 상태를 처리한다면 분리형(disaggregated) 상태 저장소 ForStStateBackend 사용을 고려할 수 있어요. 이 백엔드는 상태를 HDFS나 S3 같은 별도 저장 시스템에 저장하고 TaskManager에는 상태 메타데이터와 캐시만 유지해요. 또한 대규모 상태 애플리케이션에서는 State API V2를 ForStStateBackend와 함께 쓰는 것이 권장돼요.
RocksDB의 성능은 구성에 따라 달라질 수 있으며, 이 섹션은 RocksDB State Backend를 사용하는 작업을 튜닝하기 위한 몇 가지 모범 사례를 설명해요.
ForSt의 설계는 RocksDB와 매우 유사하고 구성 가능한 옵션도 거의 같으므로, 다음 섹션들을 참고해 ForSt를 구성할 수 있어요. 다음 글은 RocksDB 관점에서 소개된 것이에요. ForSt를 비슷한 방식으로 구성하려면 ForSt 아래의 해당 구성을 사용해야 해요.
증분 체크포인트 (Incremental Checkpoints)
체크포인트가 걸리는 시간을 줄이는 것과 관련해서는 증분 체크포인트를 활성화하는 것이 첫 번째 고려사항 중 하나여야 해요. 증분 체크포인트는 상태 백엔드의 완전하고 자체 포함된 백업을 만드는 대신 이전 완료 체크포인트와 비교한 변경 사항만 기록하므로 전체 체크포인트와 비교해 체크포인팅 시간을 극적으로 줄일 수 있어요. 배경 정보는 RocksDB의 증분 체크포인트를 참고해요.
RocksDB 또는 JVM 힙의 타이머 (Timers in RocksDB or on JVM Heap)
타이머는 기본적으로 RocksDB에 저장되며, 이는 더 견고하고 확장 가능한 선택이에요. 타이머가 거의 없는 작업만 성능 튜닝할 때(윈도우 없음, ProcessFunction에서 타이머 미사용) 그 타이머를 힙에 두면 성능이 향상될 수 있어요. 힙 기반 타이머는 체크포인팅 시간을 늘릴 수 있고 본질적으로 메모리를 넘어 확장할 수 없으므로 신중하게 사용해요. 힙 기반 타이머 구성 방법은 이 섹션을 참고해요.
RocksDB 메모리 튜닝 (Tuning RocksDB Memory)
RocksDB State Backend의 성능은 사용 가능한 메모리 양에 크게 의존해요. 성능을 높이려면 메모리를 추가하는 것이 큰 도움이 되거나, 메모리가 어떤 기능에 가는지 조정할 수 있어요.
기본적으로 RocksDB State Backend는 Flink의 managed memory 예산을 RocksDB 버퍼와 캐시에 사용해요(state.backend.rocksdb.memory.managed: true). 그 메커니즘이 어떻게 동작하는지에 대한 배경은 RocksDB Memory Management를 참고해요.
메모리 관련 성능 문제를 튜닝하려면 다음 단계가 도움이 될 수 있어요:
- 성능을 높이기 위해 시도할 첫 단계는 managed memory 양을 늘리는 것이에요. 이는 저수준 RocksDB 옵션 튜닝의 복잡성을 열지 않고도 상황을 크게 개선하는 경우가 많아요. 특히 큰 컨테이너/프로세스 크기에서는 애플리케이션 로직이 JVM 힙을 많이 필요로 하지 않는 한 총 메모리의 많은 부분을 일반적으로 RocksDB에 줄 수 있어요. 기본 managed memory 비율 *(0.4)*은 보수적이며, 다중 GB 프로세스 크기의 TaskManager를 사용할 때 종종 늘릴 수 있어요.
- RocksDB의 write buffer 수는 애플리케이션의 상태 수(파이프라인의 모든 연산자에 걸친 상태)에 따라 달라져요. 각 상태는 하나의 ColumnFamily에 해당하며, 각각 자체 write buffer가 필요해요. 따라서 상태가 많은 애플리케이션은 같은 성능을 위해 보통 더 많은 메모리가 필요해요.
state.backend.rocksdb.memory.managed: false를 설정해 managed memory가 있는 RocksDB와 column family별 메모리가 있는 RocksDB의 성능을 비교해볼 수 있어요. 특히 기준선에 대해 테스트하거나(컨테이너 메모리 제한이 없거나 충분하다고 가정) 이전 Flink 버전과 비교한 회귀 테스트를 할 때 유용할 수 있어요. managed memory 설정(상수 메모리 풀)과 비교해 managed memory를 사용하지 않으면 RocksDB가 애플리케이션의 상태 수에 비례해 메모리를 할당해요(메모리 풋프린트가 애플리케이션 변경에 따라 바뀜). 경험상 비-managed 모드는(ColumnFamily 옵션을 적용하지 않는 한) 대략 "140MB * 모든 작업의 상태 수 * 슬롯 수"의 상한을 가져요. 타이머도 상태로 간주돼요!- 애플리케이션에 상태가 많고 빈번한 MemTable 플러시(쓰기 측 병목)가 보이는데 더 많은 메모리를 줄 수 없다면 write buffer로 가는 메모리 비율을 늘릴 수 있어요(
state.backend.rocksdb.memory.write-buffer-ratio). 자세한 내용은 RocksDB Memory Management를 참고해요. - 상태가 많은 설정에서 MemTable 플러시 수를 줄이는 고급 옵션(expert mode)은 RocksDBOptionsFactory를 통해 RocksDB의 ColumnFamily 옵션(arena block size, max background flush threads 등)을 튜닝하는 것이에요:
public class MyOptionsFactory implements ConfigurableRocksDBOptionsFactory {
@Override
public DBOptions createDBOptions(DBOptions currentOptions, Collection<AutoCloseable> handlesToClose) {
// increase the max background flush threads when we have many states in one operator,
// which means we would have many column families in one DB instance.
return currentOptions.setMaxBackgroundFlushes(4);
}
@Override
public ColumnFamilyOptions createColumnOptions(
ColumnFamilyOptions currentOptions, Collection<AutoCloseable> handlesToClose) {
// decrease the arena block size from default 8MB to 1MB.
return currentOptions.setArenaBlockSize(1024 * 1024);
}
@Override
public OptionsFactory configure(ReadableConfig configuration) {
return this;
}
}
용량 계획 (Capacity Planning)
이 섹션은 Flink 작업을 안정적으로 실행하기 위해 얼마나 많은 리소스를 사용해야 할지 결정하는 방법을 논의해요. 용량 계획의 기본 경험 법칙은:
- 정상 운영은 상수 백프레셔 하에서 동작하지 않을 충분한 용량을 가져야 해요. 애플리케이션이 백프레셔 하에서 실행되는지 확인하는 방법은 백프레셔 모니터링을 참고해요.
- 실패 없는 시간 동안 프로그램을 백프레셔 없이 실행하는 데 필요한 리소스 위에 추가 리소스를 프로비저닝해요. 이 리소스들은 애플리케이션이 복구 중이던 동안 축적된 입력 데이터를 "따라잡는" 데 필요해요. 얼마나 필요한지는 복구 작업이 보통 얼마나 걸리는지(장애 조치 시 새 TaskManager에 로드해야 하는 상태 크기에 따라 달라짐)와 시나리오가 얼마나 빨리 실패를 복구해야 하는지에 따라 달라져요.
- 중요: 기준선은 체크포인팅을 활성화한 상태로 확립해야 해요. 체크포인팅은 일부 리소스(네트워크 대역폭 등)를 묶어두기 때문이에요.
- 일시적인 백프레셔는 보통 괜찮으며, 부하 급증, 따라잡기 단계, 또는 (싱크에서 쓰여지는) 외부 시스템이 일시적으로 느려질 때 실행 흐름 제어의 필수 부분이에요.
- 특정 연산(대규모 윈도우 등)은 다운스트림 연산자에 뾰족한(스파이크) 부하를 초래해요. 윈도우의 경우 윈도우가 만들어지는 동안 다운스트림 연산자는 할 일이 거의 없고 윈도우가 발행될 때 할 일이 생겨요. 다운스트림 병렬도 계획은 윈도우가 얼마나 발행하는지와 그러한 스파이크를 얼마나 빨리 처리해야 하는지 고려해야 해요.
중요: 나중에 리소스를 추가할 수 있도록 데이터 스트림 프로그램의 maximum parallelism을 합리적인 숫자로 설정해요. 최대 병렬도는 프로그램을(savepoint로) 재조정할 때 병렬도를 얼마나 높게 설정할 수 있는지 정의해요. Flink의 내부 부기는 병렬 상태를 max-parallelism개의 key group 단위로 추적해요. Flink의 설계는 프로그램을 낮은 병렬도로 실행하더라도 최대 병렬도를 매우 높은 값으로 설정하는 것을 효율적으로 만들기 위해 노력해요.
압축 (Compression)
Flink는 모든 체크포인트와 savepoint에 선택적 압축(기본값: 끔)을 제공해요. 현재 압축은 항상 snappy 압축 알고리즘(버전 1.1.10.x)을 사용하지만 향후 사용자 정의 압축 알고리즘을 지원할 계획이에요. 압축은 keyed state의 key-group 단위로 동작해요. 즉 각 key-group을 개별적으로 압축 해제할 수 있으며, 이는 재조정에 중요해요.
압축은 ExecutionConfig로 활성화할 수 있어요:
ExecutionConfig executionConfig = new ExecutionConfig();
executionConfig.setUseSnapshotCompression(true);
참고: 압축 옵션은 증분 스냅샷에 영향이 없어요. 증분 스냅샷은 항상 기본적으로 snappy 압축을 사용하는 RocksDB의 내부 포맷을 사용하기 때문이에요.
Task-Local 복구 (Task-Local Recovery)
동기 (Motivation)
Flink의 체크포인팅에서 각 태스크는 자신의 상태 스냅샷을 만들고 이를 분산 저장소에 써요. 각 태스크는 분산 저장소의 상태 위치를 설명하는 핸들을 보내 job manager에 상태 쓰기가 성공했음을 알려줘요. job manager는 모든 태스크의 핸들을 수집해 이를 체크포인트 객체로 묶어요.
복구 시 job manager는 최신 체크포인트 객체를 열고 핸들을 해당 태스크로 다시 보내며, 태스크는 분산 저장소에서 자신의 상태를 복원할 수 있어요. 분산 저장소를 사용해 상태를 저장하는 것은 두 가지 중요한 장점이 있어요. 첫째 저장소가 장애 허용적이고, 둘째 분산 저장소의 모든 상태가 모든 노드에서 접근 가능하며 쉽게 재분배될 수 있어요(예: 재조정용).
하지만 원격 분산 저장소를 사용하는 것은 큰 단점도 하나 있어요: 모든 태스크가 원격 위치에서 네트워크를 통해 자신의 상태를 읽어야 해요. 많은 시나리오에서 복구는 실패한 태스크를 이전 실행과 같은 task manager로 재스케줄링할 수 있지만(물론 머신 실패 같은 예외는 있음), 여전히 원격 상태를 읽어야 해요. 이는 단일 머신에 작은 실패만 있어도 대규모 상태에 대한 긴 복구 시간을 초래할 수 있어요.
접근 (Approach)
Task-local 상태 복구는 정확히 이 긴 복구 시간 문제를 대상으로 하며 핵심 아이디어는 다음과 같아요: 각 체크포인트에 대해 각 태스크는 태스크 상태를 분산 저장소에 쓸 뿐 아니라 태스크에 로컬인 저장소(예: 로컬 디스크 또는 메모리)에 상태 스냅샷의 보조 복사본을 유지해요. 스냅샷의 기본 저장소는 여전히 분산 저장소여야 한다는 점에 주의해요. 로컬 저장소는 노드 실패 시 내구성을 보장하지 않고 다른 노드에 상태 재분배 접근을 제공하지 않기 때문이에요. 이 기능은 여전히 기본 복사본을 필요로 해요.
하지만 복구를 위해 이전 위치로 재스케줄링할 수 있는 각 태스크에 대해 보조 로컬 복사본에서 상태를 복원하고 원격으로 상태를 읽는 비용을 피할 수 있어요. 많은 실패가 노드 실패가 아니고 노드 실패는 보통 한 번에 하나 또는 아주 소수의 노드에만 영향을 미치므로, 복구에서 대부분의 태스크가 이전 위치로 돌아가 자신의 로컬 상태가 온전함을 찾을 가능성이 매우 높아요. 이것이 로컬 복구가 복구 시간을 줄이는 데 효과적인 이유예요.
선택한 상태 백엔드와 체크포인팅 전략에 따라 보조 로컬 상태 복사본을 만들고 저장하는 데 체크포인트당 약간의 추가 비용이 들 수 있다는 점에 주의해요. 예를 들어 대부분의 경우 구현은 분산 저장소에 대한 쓰기를 로컬 파일로 단순히 중복해요.
기본(분산 저장소)과 보조(task-local) 상태 스냅샷의 관계 (Relationship of primary (distributed store) and secondary (task-local) state snapshots)
Task-local 상태는 항상 보조 복사본으로 간주되며, 체크포인트 상태의 진실(ground truth)은 분산 저장소의 기본 복사본이에요. 이는 체크포인팅과 복구 중 로컬 상태 문제에 대한 영향이 있어요:
- 체크포인팅의 경우 기본 복사본이 성공해야 하며 보조 로컬 복사본 생성 실패는 체크포인트를 실패시키지 않아요. 기본 복사본을 만들 수 없으면 보조 복사본이 성공적으로 만들어졌어도 체크포인트는 실패해요.
- 기본 복사본만 job manager가 인정하고 관리하며, 보조 복사본은 task manager가 소유하고 그 수명주기는 기본 복사본과 독립적일 수 있어요. 예를 들어 최근 체크포인트 3개의 히스토리를 기본 복사본으로 유지하고 최신 체크포인트의 task-local 상태만 유지할 수 있어요.
- 복구의 경우 일치하는 보조 복사본이 있으면 Flink는 항상 먼저 task-local 상태에서 복원을 시도해요. 보조 복사본에서 복구 중 문제가 발생하면 Flink는 기본 복사본에서 태스크 복구를 투명하게 재시도해요. 기본과 (선택적) 보조 복사본이 모두 실패해야 복구가 실패해요. 이 경우 구성에 따라 Flink가 더 오래된 체크포인트로 대체(fall back)할 수 있어요.
- task-local 복사본이 전체 태스크 상태의 일부만 포함할 수 있어요(예: 로컬 파일 하나를 쓰는 중 예외). 이 경우 Flink는 먼저 로컬 부분을 로컬에서 복구하려 하고, 비로컬 상태는 기본 복사본에서 복원돼요. 기본 상태는 항상 완전해야 하며 *task-local 상태의 상위 집합(superset)*이에요.
- task-local 상태는 기본 상태와 다른 포맷일 수 있으며, 바이트 단위로 동일할 필요가 없어요. 예를 들어 task-local 상태가 어떤 파일에도 저장되지 않고 힙 객체로 구성된 인메모리일 수도 있어요.
- task manager가 사라지면 그 모든 태스크의 로컬 상태도 사라져요.
task-local 복구 구성 (Configuring task-local recovery)
Task-local 복구는 기본적으로 비활성화되어 있으며 CheckpointingOptions.LOCAL_RECOVERY에 지정된 대로 state.backend.local-recovery 키로 Flink 구성을 통해 활성화할 수 있어요. 이 설정의 값은 로컬 복구를 활성화하는 true 또는 비활성화하는 false(기본값)일 수 있어요.
정렬되지 않은 체크포인트는 현재 task-local 복구를 지원하지 않는다는 점에 주의해요.
서로 다른 상태 백엔드의 task-local 복구 세부사항 (Details on task-local recovery for different state backends)
제한: 현재 task-local 복구는 keyed 상태 백엔드만 다뤄요. Keyed 상태는 일반적으로 상태의 압도적으로 큰 부분이에요. 가까운 미래에는 연산자 상태와 타이머도 다룰 것입니다.
다음 상태 백엔드가 task-local 복구를 지원할 수 있어요.
- HashMapStateBackend: keyed 상태에 대해 task-local 복구가 지원돼요. 구현은 상태를 로컬 파일로 중복해요. 이는 추가 쓰기 비용을 도입하고 로컬 디스크 공간을 차지할 수 있어요. 미래에는 task-local 상태를 메모리에 유지하는 구현도 제공할 수 있어요.
- EmbeddedRocksDBStateBackend: keyed 상태에 대해 task-local 복구가 지원돼요. 전체 체크포인트의 경우 상태가 로컬 파일로 중복돼요. 이는 추가 쓰기 비용을 도입하고 로컬 디스크 공간을 차지할 수 있어요. 증분 스냅샷의 경우 로컬 상태는 RocksDB의 네이티브 체크포인팅 메커니즘에 기반해요. 이 메커니즘은 기본 복사본을 만드는 첫 단계로도 사용돼서, 이 경우 보조 복사본 생성에 추가 비용이 들지 않아요. 분산 저장소에 업로드한 후 그 네이티브 체크포인트 디렉터리를 삭제하는 대신 그냥 유지해요. 이 로컬 복사본은 RocksDB의 작업 디렉터리와 활성 파일을(하드 링크로) 공유할 수 있어, 활성 파일의 경우 증분 스냅샷의 task-local 복구로 추가 디스크 공간이 소비되지 않아요. 하드 링크를 사용한다는 것은 RocksDB 디렉터리가 로컬 상태를 저장하는 데 사용할 수 있는 모든 구성된 로컬 복구 디렉터리와 같은 물리 장치에 있어야 함을 의미하며, 그렇지 않으면 하드 링크 생성이 실패할 수 있어요(FLINK-10954 참고). 현재는 또한 RocksDB 디렉터리가 둘 이상의 물리 장치에 위치하도록 구성될 때 로컬 복구 사용을 막아요.
할당 보존 스케줄링 (Allocation-preserving scheduling)
Task-local 복구는 실패 시 할당 보존 태스크 스케줄링을 가정하며, 이는 다음과 같이 동작해요. 각 태스크는 이전 할당을 기억하고 복구에서 재시작하기 위해 정확히 같은 슬롯을 요청해요. 그 슬롯을 사용할 수 없으면 태스크는 리소스 매니저에서 새롭고 비어 있는 슬롯을 요청해요. 이렇게 하면 task manager를 더 이상 사용할 수 없을 때 이전 위치로 돌아갈 수 없는 태스크가 다른 복구 중 태스크를 이전 슬롯에서 밀어내지 않아요. 우리의 추론은 이전 슬롯은 task manager를 더 이상 사용할 수 없을 때만 사라질 수 있고, 이 경우 어떤 태스크든 어쨌든 새 슬롯을 요청해야 한다는 것이에요. 이 스케줄링 전략으로 우리는 최대한 많은 태스크에 로컬 상태에서 복구할 기회를 주고 태스크가 서로의 이전 슬롯을 훔치는 계단식 효과를 피해요.