Job Master 장애로부터의 배치 작업 진행 복구
Job Master 장애로부터의 배치 작업 진행 복구 (Batch jobs progress recovery from job master failures)
JobMaster가 실패한 뒤에도 배치 작업이 가능한 한 많은 진행 상황을 복구할 수 있게 하는 메커니즘입니다. JobEventStore를 도입하고 TaskManager가 중간 결과 데이터를 보존하도록 하여, 이미 끝난 작업을 다시 실행하지 않도록 합니다.
출처: 문서
본문
배경 (Background)
이전에는 JobMaster가 실패하여 종료되면 다음 두 가지 상황 중 하나가 발생했습니다:
- 고가용성(HA)이 비활성화되어 있으면 작업이 실패합니다.
- HA가 활성화되어 있으면 JobMaster 장애 조치가 발생하고 작업이 다시 시작됩니다. 스트리밍 작업은 가장 최근의 성공한 체크포인트부터 재개할 수 있습니다. 그러나 배치 작업은 체크포인트가 없어 처음부터 다시 시작해야 하며, 이전에 만든 모든 진행 상황을 잃게 됩니다. 이는 장기 실행 배치 작업에 상당한 회귀(regression)입니다.
이 문제를 해결하기 위해, JobMaster 장애 조치 후 배치 작업이 가능한 한 많은 진행 상황을 복구하고 이미 끝난 작업을 다시 실행하지 않도록 하는 배치 작업 복구 메커니즘이 도입되었습니다.
이 기능을 구현하기 위해 JobMaster의 상태 변경 이벤트(예: ExecutionGraph, OperatorCoordinator 등)를 외부 파일시스템에 기록하는 JobEventStore 구성 요소가 도입되었습니다. JobMaster가 크래시한 후 재시작되는 동안 TaskManager는 작업이 생산한 중간 결과 데이터를 보존하고 계속해서 재연결을 시도합니다. JobMaster가 재시작되면 TaskManager와의 연결을 다시 수립하고, 보존된 중간 결과와 JobEventStore에 이전에 기록된 이벤트를 기반으로 작업 상태를 복구하여 작업의 실행 진행을 재개합니다.
사용법 (Usage)
이 섹션은 JobMaster 장애로부터 배치 작업 복구를 활성화하는 방법, 튜닝하는 방법, 그리고 배치 작업 진행 복구와 함께 동작하도록 소스를 개발하는 방법을 설명합니다.
Job master 장애로부터 배치 작업 진행 복구를 활성화하는 방법
-
클러스터 고가용성 활성화:
JobMaster 장애로부터 배치 작업 복구를 활성화하려면 먼저 클러스터 고가용성(HA)이 활성화되어 있는지 확인하는 것이 필수적입니다. Flink는 ZooKeeper 또는 Kubernetes로 백업되는 HA 서비스를 지원합니다. 구성에 대한 자세한 내용은 High Availability 페이지에서 확인할 수 있습니다.
-
execution.batch.job-recovery.enabled 설정: true
현재 Adaptive Batch Scheduler만이 이 기능을 지원합니다. 그리고 Flink 배치 작업은 다른 스케줄러가 명시적으로 구성되지 않는 한 기본적으로 이 스케줄러를 사용합니다.
최적화 (Optimization)
JobMaster 장애 조치 후 배치 작업이 가능한 한 많은 진행 상황을 복구하고 이미 끝난 작업을 다시 실행하지 않도록 하려면 다음 옵션을 구성해 최적화할 수 있습니다:
- execution.batch.job-recovery.snapshot.min-pause: OperatorCoordinator와 ShuffleMaster의 스냅샷 사이에 허용되는 최소 일시 정지 시간을 결정합니다. 이 파라미터는 클러스터의 예상 I/O 부하와 허용 가능한 상태 회귀의 양에 따라 조정할 수 있습니다. 더 작은 상태 회귀를 선호하고 더 높은 I/O 부하가 허용된다면 이 간격을 줄이세요.
- execution.batch.job-recovery.previous-worker.recovery.timeout: Shuffle worker가 재연결하는 데 허용되는 타임아웃 기간을 결정합니다. 복구 과정에서 Flink는 Shuffle Master에게 보존된 중간 결과 데이터 정보를 요청합니다. 타임아웃에 도달하면 Flink는 획득한 모든 중간 결과 데이터를 사용해 상태를 복구합니다.
- job-event.store.write-buffer.flush-interval: JobEventStore의 쓰기 버퍼 플러시 간격을 결정합니다.
- job-event.store.write-buffer.size: JobEventStore의 쓰기 버퍼 크기를 결정합니다. 버퍼가 가득 차면 그 내용이 외부 파일시스템으로 플러시됩니다.
소스에 대한 배치 작업 진행 복구 활성화
현재 새 소스(FLIP-27)만 배치 작업의 진행 복구를 지원합니다. 이 기능을 달성하려면 새 소스(FLIP-27)의 SplitEnumerator가 배치 처리 시나리오(checkpointId가 -1로 설정된 경우)에서 상태 스냅샷을 찍을 수 있고 SupportsBatchSnapshot 인터페이스를 구현할 수 있어야 합니다. 그래야 job master 장애 이전의 진행 상황으로 복구할 수 있습니다. 그렇지 않으면 데이터 정확성을 보장하기 위해 job master 장애 조치 후 다음 두 상황 중 하나가 발생합니다:
- 이 소스의 모든 태스크가 끝나지 않았다면, 모든 태스크를 리셋하고 다시 실행합니다.
- 이 소스의 모든 태스크가 끝났다면 추가 조치가 필요 없고 작업은 계속 실행될 수 있습니다. 그러나 이 태스크들 중 어떤 것이 나중에 (예:
PartitionNotFound예외로 인해) 다시 시작되어야 한다면, 이 소스의 모든 서브태스크를 리셋하고 다시 실행해야 합니다.
제한 사항 (Limitations)
- 새 소스(FLIP-27)에서만 동작: 레거시 소스는 더 이상 사용되지 않으므로 이 기능은 새 소스만 지원합니다.
- Adaptive Batch Scheduler 전용: 현재 Adaptive Batch Scheduler만 JobMaster 장애 조치 후 배치 작업 복구를 지원합니다. 결과적으로 이 기능은 Adaptive Batch Scheduler의 한계를 모두 상속합니다.
- 원격 셔플 서비스(remote shuffle services)를 사용할 때는 동작하지 않습니다.