작업과 스케줄링
작업과 스케줄링 (Jobs and Scheduling)
이 문서는 Flink가 작업을 어떻게 스케줄링하고 JobManager 에서 작업 상태를 어떻게 표현하고 추적하는지 간략하게 설명합니다.
출처: 문서
본문
이 문서는 Flink가 작업을 어떻게 스케줄링하고 JobManager 에서 작업 상태를 어떻게 표현하고 추적하는지 간략하게 설명합니다.
스케줄링 (Scheduling)
Flink의 실행 리소스는 Task Slots 로 정의됩니다. 각 TaskManager 는 하나 이상의 task slot 을 가지며, 각 slot 은 하나의 병렬 작업 파이프라인을 실행할 수 있습니다. 파이프라인은 n번째 병렬 MapFunction 인스턴스와 n번째 병렬 ReduceFunction 인스턴스처럼 여러 개의 연속된 task 로 구성됩니다. Flink는 종종 연속된 task 를 동시에 실행한다는 점에 유의하세요: 스트리밍 프로그램에서는 어떤 경우든 그렇게 되며, 배치 프로그램에서도 자주 발생합니다.
아래 그림은 이를 보여줍니다. 데이터 소스, MapFunction, ReduceFunction 이 있는 프로그램을 생각해 보세요. source 와 MapFunction 은 병렬도 4로 실행되고, ReduceFunction 은 병렬도 3으로 실행됩니다. 파이프라인은 Source - Map - Reduce 시퀀스로 구성됩니다. 각각 3개의 slot 을 가진 TaskManager 2개가 있는 클러스터에서 프로그램은 아래에 설명된 대로 실행됩니다.
내부적으로 Flink는 SlotSharingGroup 과 CoLocationGroup 을 통해 어떤 task 가 slot 을 공유할 수 있는지(관대하게), 각각 어떤 task 가 반드시 같은 slot 에 엄격히 배치되어야 하는지를 정의합니다.
JobManager 데이터 구조
작업 실행 중 JobManager 는 분산 task 를 추적하고 다음 task(또는 task 집합)를 언제 스케줄링할지 결정하며, 완료된 task 나 실행 실패에 대응합니다.
JobManager 는 operator(JobVertex)와 중간 결과(IntermediateDataSet)로 구성된 데이터 흐름의 표현인 JobGraph 을 수신합니다. 각 operator 는 병렬도와 실행하는 코드 같은 속성을 가집니다. 또한 JobGraph 에는 operator 의 코드를 실행하는 데 필요한 라이브러리 집합이 첨부됩니다.
JobManager 는 JobGraph 를 ExecutionGraph 으로 변환합니다. ExecutionGraph 는 JobGraph 의 병렬 버전입니다: 각 JobVertex 에 대해 병렬 subtask 마다 하나의 ExecutionVertex 를 포함합니다. 병렬도 100의 operator 는 하나의 JobVertex 와 100개의 ExecutionVertex 를 가집니다. ExecutionVertex 는 특정 subtask 의 실행 상태를 추적합니다. 한 JobVertex 의 모든 ExecutionVertex 는 operator 전체의 상태를 추적하는 ExecutionJobVertex 에 보관됩니다. vertex 외에도 ExecutionGraph 는 IntermediateResult 와 IntermediateResultPartition 을 포함합니다. 전자는 IntermediateDataSet 의 상태를, 후자는 각 파티션의 상태를 추적합니다.
각 ExecutionGraph 에는 작업 상태(job status)가 연관됩니다. 이 작업 상태는 작업 실행의 현재 상태를 나타냅니다.
Flink 작업은 먼저 created 상태가 되고, 그 다음 running 으로 전환되며, 모든 작업이 완료되면 finished 로 전환됩니다. 실패의 경우 작업은 먼저 모든 실행 중 task 를 취소하는 failing 상태로 전환됩니다. 모든 job vertex 가 최종 상태에 도달했고 작업을 재시작할 수 없다면 작업은 failed 상태로 전환됩니다. 작업을 재시작할 수 있다면 restarting 상태로 들어갑니다. 작업이 완전히 재시작되면 created 상태에 도달합니다.
사용자가 작업을 취소하면 cancelling 상태로 들어갑니다. 이는 또한 현재 실행 중인 모든 task 의 취소를 수반합니다. 실행 중인 모든 task 가 최종 상태에 도달하면 작업은 cancelled 상태로 전환됩니다.
전역적으로 종료 상태를 나타내어 작업의 정리를 유발하는 finished, canceled, failed 상태와 달리 suspended 상태는 로컬로만 종료 상태입니다. 로컬로만 종료 상태라는 것은 작업의 실행이 해당 JobManager 에서 종료되었지만 Flink 클러스터의 다른 JobManager 가 영구 HA 저장소에서 작업을 검색해 재시작할 수 있다는 뜻입니다. 결과적으로 suspended 상태에 도달한 작업은 완전히 정리되지 않습니다.
ExecutionGraph 실행 중 각 병렬 task 는 created 부터 finished 또는 failed 까지 여러 단계를 거칩니다. 아래 다이어그램은 상태와 그 사이의 가능한 전환을 보여줍니다. task 는 여러 번 실행될 수 있습니다(예: 장애 복구 과정에서). 그 이유로 ExecutionVertex 의 실행은 Execution 에서 추적됩니다. 각 ExecutionVertex 는 현재 Execution 과 이전 Execution 을 가집니다.