백프레셔 하의 체크포인팅
백프레셔 하의 체크포인팅 (Checkpointing under backpressure)
일반적으로 정렬된(aligned) 체크포인팅 시간은 체크포인팅 과정의 동기 및 비동기 부분이 지배합니다. 그러나 Flink 작업이 심한 백프레셔 하에서 실행 중이면 체크포인트의 end-to-end 시간에서 지배 요인은 체크포인트 배리어를 모든 연산자/서브태스크에 전파하는 시간이 될 수 있습니다.
출처: 문서
본문
일반적으로 정렬된 체크포인팅 시간은 체크포인팅 과정의 동기 및 비동기 부분이 지배합니다. 그러나 Flink 작업이 심한 백프레셔 하에서 실행 중이면 체크포인트의 end-to-end 시간에서 지배 요인은 체크포인트 배리어를 모든 연산자/서브태스크에 전파하는 시간일 수 있습니다. 이는 체크포인팅 과정의 개요에서 설명되며, 높은 alignment time과 start delay 메트릭으로 관찰할 수 있습니다. 이런 일이 발생하여 문제가 되면 문제를 해결하는 세 가지 방법이 있습니다:
- Flink 작업 최적화, Flink 또는 JVM 구성 조정, 또는 스케일 업으로 백프레셔 원인 제거.
- Flink 작업에서 버퍼링된 in-flight 데이터 양 줄이기.
- 비정렬(unaligned) 체크포인트 활성화.
이 옵션들은 상호 배타적이지 않으며 함께 결합할 수 있습니다. 이 문서는 후자의 두 옵션에 초점을 맞춥니다.
버퍼 디블로팅 (Buffer debloating)
Flink 1.14는 Flink 연산자/서브태스크 사이에 버퍼링된 in-flight 데이터 양을 자동으로 제어하는 새 도구를 도입했습니다. 버퍼 디블로팅 메커니즘은 taskmanager.network.memory.buffer-debloat.enabled 속성을 true로 설정하여 활성화할 수 있습니다.
이 기능은 정렬 및 비정렬 체크포인트 모두에서 동작하며 두 경우 모두 체크포인팅 시간을 개선할 수 있지만, 디블로팅 효과는 정렬된 체크포인트에서 가장 두드러집니다. 비정렬 체크포인트와 함께 버퍼 디블로팅을 사용하면 더 작은 체크포인트 크기와 더 빠른 복구 시간(유지·복구할 in-flight 데이터가 줄어듦)이라는 추가 이점이 있습니다.
버퍼 디블로팅 기능이 어떻게 동작하고 어떻게 구성하는지에 대한 자세한 내용은 네트워크 메모리 튜닝 가이드를 참고하세요. 버퍼링된 in-flight 데이터 양을 수동으로 줄일 수도 있는데, 이는 앞서 언급한 튜닝 가이드에 설명되어 있습니다.
비정렬 체크포인트 (Unaligned checkpoints)
Flink 1.11부터 체크포인트는 비정렬일 수 있습니다. 비정렬 체크포인트는 체크포인트 상태의 일부로 in-flight 데이터(즉, 버퍼에 저장된 데이터)를 포함하여 체크포인트 배리어가 이 버퍼들을 추월할 수 있게 합니다. 따라서 체크포인트 배리어가 더 이상 데이터 스트림에 효과적으로 내장되지 않으므로 체크포인트 기간은 현재 처리량과 무관해집니다.
백프레셔로 인해 체크포인팅 기간이 매우 높다면 비정렬 체크포인트를 사용해야 합니다. 그러면 체크포인팅 시간이 end-to-end 지연과 대부분 무관해집니다. 비정렬 체크포인팅은 상태 저장소에 I/O를 추가하므로, 체크포인팅 중 상태 저장소에 대한 I/O가 실제로 병목이라면 사용해서는 안 됩니다.
비정렬 체크포인트를 활성화하려면:
Java:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// enables the unaligned checkpoints
env.getCheckpointConfig().enableUnalignedCheckpoints();
Python:
env = StreamExecutionEnvironment.get_execution_environment()
# enables the unaligned checkpoints
env.get_checkpoint_config().enable_unaligned_checkpoints()
또는 flink-conf.yml 구성 파일에서:
execution.checkpointing.unaligned: true
정렬 체크포인트 타임아웃 (Aligned checkpoint timeout)
비정렬 체크포인트를 활성화한 후 정렬 체크포인트 타임아웃을 프로그래밍 방식으로 지정할 수도 있습니다:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.getCheckpointConfig().setAlignedCheckpointTimeout(Duration.ofSeconds(30));
또는 flink-conf.yml 구성 파일에서:
execution.checkpointing.aligned-checkpoint-timeout: 30 s
활성화되면 각 체크포인트는 여전히 정렬 체크포인트로 시작되지만, 전역 체크포인트 기간이 aligned-checkpoint-timeout을 초과하면 정렬 체크포인트가 완료되지 않은 경우 체크포인트가 비정렬 체크포인트로 진행됩니다.
제한 사항 (Limitations)
동시 체크포인트 (Concurrent checkpoints)
Flink는 현재 동시 비정렬 체크포인트를 지원하지 않습니다. 그러나 더 예측 가능하고 짧은 체크포인팅 시간 때문에 동시 체크포인트가 전혀 필요하지 않을 수도 있습니다. 그러나 savepoint도 비정렬 체크포인트와 동시에 발생할 수 없으므로 savepoint가 약간 더 오래 걸립니다.
워터마크와의 상호작용 (Interplay with watermarks)
비정렬 체크포인트는 복구 중 워터마크와 관련된 암묵적 보장을 깨뜨립니다. 현재 Flink는 리스케일링을 쉽게 하기 위해 연산자에 최신 워터마크를 저장하는 대신 복구의 첫 단계로 워터마크를 생성합니다. 비정렬 체크포인트에서는 복구 시 Flink가 in-flight 데이터를 복원한 후에 워터마크를 생성한다는 뜻입니다. 파이프라인이 각 레코드에 최신 워터마크를 적용하는 연산자를 사용하면 정렬 체크포인트와는 다른 결과를 생성합니다. 연산자가 최신 워터마크가 항상 사용 가능한 것에 의존한다면, 워터마크를 연산자 상태에 저장하는 것이 해결책입니다. 그 경우 리스케일링을 지원하려면 워터마크를 유니온 상태(union state)의 키 그룹마다 저장해야 합니다.
장기 실행 레코드 처리와의 상호작용 (Interplay with long-running record processing)
비정렬 체크포인트 배리어가 큐의 다른 모든 레코드를 추월할 수 있음에도 불구하고, 현재 레코드를 처리하는 데 오래 걸리면 이 배리어의 처리는 여전히 지연될 수 있습니다. 이 상황은 예를 들어 윈도우 연산에서 많은 타이머를 한 번에 발화시킬 때 발생할 수 있습니다. 두 번째 문제 시나리오는 단일 입력 레코드를 처리할 때 둘 이상의 네트워크 버퍼 가용성을 기다리며 시스템이 차단되는 경우입니다. Flink는 단일 입력 레코드의 처리를 중단할 수 없으며, 비정렬 체크포인트는 현재 처리 중인 레코드가 완전히 처리될 때까지 기다려야 합니다. 이는 두 시나리오에서 문제를 일으킬 수 있습니다. 단일 네트워크 버퍼에 맞지 않는 큰 레코드의 직렬화 결과이거나, 하나의 입력 레코드에 대해 많은 출력 레코드를 생성하는 flatMap 연산에서입니다. 이러한 시나리오에서 백프레셔는 단일 입력 레코드를 처리하는 데 필요한 모든 네트워크 버퍼가 사용 가능해질 때까지 비정렬 체크포인트를 차단할 수 있습니다. 단일 레코드 처리가 잠시 걸리는 다른 상황에서도 발생할 수 있습니다. 결과적으로 체크포인트 시간이 예상보다 높아지거나 변동할 수 있습니다.
특정 데이터 분포 패턴은 체크포인트되지 않음
체크포인트에 저장된 채널 데이터로 유지할 수 없는 속성을 가진 연결 유형이 있습니다. 이러한 특성을 보존하고 상태 손상이나 예기치 않은 동작을 방지하기 위해 그러한 연결에 대해서는 비정렬 체크포인트가 비활성화됩니다. 다른 모든 교환(exchange)은 여전히 비정렬 체크포인트를 수행합니다.
Pointwise 연결
현재 pointwise 연결에 대해 데이터 순서성에 대한 하드 보장은 없습니다. 그러나 데이터가 이전 소스나 keyby와 같은 방식으로 암묵적으로 구조화되었기 때문에 일부 사용자는 순서성 보장에 의존하면서 계산 집약적 태스크를 더 작은 덩어리로 나누는 데 이 동작에 의존했습니다.
병렬도가 변하지 않는 한 비정렬 체크포인트(UC)는 이러한 속성을 유지합니다. UC의 리스케일링이 추가되면서 그것은 바뀌었습니다.
병렬도 p = 2에서 p = 3으로 리스케일하려 한다면 갑자기 keyby 채널 안의 레코드가 키 그룹에 따라 세 개의 채널로 나뉘어야 합니다. 이는 연산자의 키 그룹 범위와 레코드의 키(그룹)를 결정하는 방법을 사용하면 쉽게 가능합니다(실제 접근 방식과 무관). 순방향(forward) 채널의 경우 키 컨텍스트가 전혀 없습니다. 순방향 채널의 어떤 레코드도 키 그룹이 할당되어 있지 않습니다. 키가 여전히 존재한다는 보장이 없으므로 계산하는 것도 불가능합니다.
Broadcast 연결
Broadcast 연결은 또 다른 문제를 제기합니다. 모든 채널에서 레코드가 같은 속도로 소비된다는 보장이 없습니다. 이로 인해 일부 태스크는 특정 broadcast 이벤트에 해당하는 상태 변경을 적용하는 반면 다른 태스크는 적용하지 않을 수 있습니다.
Broadcast 파티셔닝은 모든 연산자에 걸쳐 동일해야 하는 broadcast 상태를 구현하는 데 자주 사용됩니다. Flink는 상태 저장 연산자의 서브태스크 0에서 상태의 단일 복사본만 체크포인팅하여 broadcast 상태를 구현합니다. 복원 시 우리는 그 복사본을 모든 연산자에 보냅니다. 따라서 연산자가 체크포인트된 채널에서 곧 소비할 레코드에 대한 변경이 적용된 상태를 얻게 될 수 있습니다.
트러블슈팅 (Troubleshooting)
손상된 in-flight 데이터 (Corrupted in-flight data)
경고: 아래 설명된 조치는 최후의 수단으로, 데이터 손실로 이어집니다.
in-flight 데이터가 손상되었거나 다른 이유로 작업을 in-flight 데이터 없이 복원해야 하는 경우, recover-without-channel-state.checkpoint-id 속성을 사용할 수 있습니다. 이 속성은 in-flight 데이터가 무시될 체크포인트 id를 지정해야 합니다. 영속된 in-flight 데이터 내부의 손상이 달리 복구 불가능한 상황을 초래하지 않는 한 이 속성을 설정하지 마세요. 이 속성은 작업이 재배포된 후에만 적용될 수 있으며, 이는 externalized checkpoint가 활성화된 경우에만 이 작업이 의미가 있다는 뜻입니다.