태스크 수명주기
태스크 수명주기 (Task Lifecycle)
Flink에서 태스크는 실행의 기본 단위입니다. 연산자의 각 병렬 인스턴스가 실행되는 곳입니다. 예를 들어 병렬도가 5인 연산자는 각 인스턴스가 별도의 태스크로 실행됩니다.
StreamTask는 Flink 스트리밍 엔진의 모든 다양한 태스크 하위 유형의 기반입니다. 이 문서는 StreamTask 수명주기의 다양한 단계를 살펴보고 각 단계를 나타내는 주요 메서드를 설명합니다.
출처: 문서
본문
한눈에 보는 연산자 수명주기
태스크는 연산자의 병렬 인스턴스를 실행하는 엔티티이므로 그 수명주기는 연산자의 수명주기와 밀접하게 통합되어 있습니다. 따라서 StreamTask 자체에 들어가기 전에 연산자의 수명주기를 나타내는 기본 메서드를 간략히 언급하겠습니다. 목록은 각 메서드가 호출되는 순서대로 아래에 제시됩니다. 연산자는 사용자 정의 함수(UDF)를 가질 수 있으므로, 각 연산자 메서드 아래에 (들여쓰기하여) 그것이 호출하는 UDF 수명주기의 메서드도 제시합니다. 이러한 메서드는 연산자가 UDF를 실행하는 모든 연산자의 기본 클래스인 AbstractUdfStreamOperator를 확장하는 경우 사용할 수 있습니다.
// initialization phase
OPERATOR::setup
UDF::setRuntimeContext
OPERATOR::initializeState
OPERATOR::open
UDF::open
// processing phase (called on every element/watermark)
OPERATOR::processElement
UDF::run
OPERATOR::processWatermark
// checkpointing phase (called asynchronously on every checkpoint)
OPERATOR::snapshotState
// notify the operator about the end of processing records
OPERATOR::finish
// termination phase
OPERATOR::close
UDF::close
간단히 말해 setup()은 RuntimeContext와 메트릭 수집 데이터 구조 같은 연산자 특정 메커니즘을 초기화하기 위해 호출됩니다. 그 후 initializeState()는 연산자에게 초기 상태를 제공하고, open() 메서드는 AbstractUdfStreamOperator의 경우 사용자 정의 함수를 여는 것 같은 연산자 특정 초기화를 실행합니다.
initializeState()는 초기 실행 중 연산자 상태를 초기화하는 로직(예: keyed state 등록)과 실패 후 checkpoint에서 상태를 검색하는 로직을 모두 포함합니다. 이 페이지의 나머지 부분에서 더 자세히 설명합니다.
이제 모든 것이 설정되었으니 연산자는 들어오는 데이터를 처리할 준비가 되었습니다. 들어오는 요소는 입력 요소, 워터마크, checkpoint barrier 중 하나일 수 있습니다. 각각을 처리하는 특별한 요소가 있습니다. 요소는 processElement() 메서드로, 워터마크는 processWatermark()로 처리되며, checkpoint barrier는 아래에서 설명할 snapshotState() 메서드를 (비동기적으로) 호출하는 checkpoint를 트리거합니다. 들어오는 각 요소에 대해 그 타입에 따라 앞서 언급한 메서드 중 하나가 호출됩니다. processElement()는 UDF의 로직이 호출되는 곳이기도 합니다. 예: MapFunction의 map() 메서드.
마지막으로 연산자의 정상적이고 실패 없는 종료(예: 스트림이 유한하고 끝에 도달한 경우)의 경우 finish() 메서드가 호출되어 연산자 로직이 요구하는 최종 부기(bookkeeping) 동작(예: 버퍼링된 데이터 플러시, 처리 종료를 표시하기 위한 데이터 발행)을 수행하고, 그 후 close()가 호출되어 연산자가 보유한 리소스(예: 열린 네트워크 연결, io 스트림, 연산자 데이터가 보유한 네이티브 메모리)를 해제합니다.
실패 또는 수동 취소로 인한 종료의 경우 실행은 close()로 직접 점프하며, 실패가 발생했을 때 연산자가 있던 단계와 close() 사이의 중간 단계를 건너뜁니다.
Checkpoints: 연산자의 snapshotState() 메서드는 checkpoint barrier가 수신될 때마다 위에 설명된 다른 메서드들과 비동기적으로 호출됩니다. Checkpoint는 처리 단계 동안, 즉 연산자가 열린 후 닫히기 전에 수행됩니다. 이 메서드의 책임은 연산자의 현재 상태를 지정된 state backend에 저장하는 것이며, 실패 후 작업이 실행을 재개할 때 거기에서 검색됩니다. 아래에 Flink의 checkpointing 메커니즘에 대한 간략한 설명을 포함하며, Flink의 checkpointing 원리에 대한 더 자세한 논의는 해당 문서를 읽어보세요: Data Streaming Fault Tolerance.
태스크 수명주기
연산자의 주요 단계에 대한 간략한 소개에 이어, 이 섹션은 태스크가 클러스터에서 실행되는 동안 관련 메서드를 어떻게 호출하는지 자세히 설명합니다. 여기 설명된 단계의 시퀀스는 주로 StreamTask 클래스의 invoke() 메서드에 포함되어 있습니다. 이 문서의 나머지는 두 개의 하위 섹션으로 나뉩니다. 하나는 태스크의 정상적이고 실패 없는 실행 중 단계를 설명하고(Normal Execution 참조), 다른 (더 짧은) 하나는 태스크가 수동으로 또는 실행 중 예외 같은 다른 이유로 취소되는 경우에 따르는 다른 시퀀스를 설명합니다(Interrupted Execution 참조).
정상 실행
중단 없이 완료될 때까지 실행되는 태스크가 거치는 단계는 아래에 설명되어 있습니다.
TASK::setInitialState
TASK::invoke
create basic utils (config, etc) and load the chain of operators
setup-operators
task-specific-init
initialize-operator-states
open-operators
run
finish-operators
wait for the final checkpoint completed (if enabled)
close-operators
task-specific-cleanup
common-cleanup
위와 같이 태스크 구성 복구와 일부 중요한 런타임 매개변수 초기화 후, 태스크의 가장 첫 단계는 초기 태스크 전체 상태를 검색하는 것입니다. 이는 setInitialState()에서 수행되며, 두 경우에 특히 중요합니다.
- 태스크가 실패에서 복구되어 마지막 성공한 checkpoint에서 다시 시작할 때
- savepoint에서 재개할 때
태스크가 처음 실행되는 경우 초기 태스크 상태는 비어 있습니다.
초기 상태를 복구한 후 태스크는 invoke() 메서드로 들어갑니다. 여기서 먼저 각각의 setup() 메서드를 호출하여 로컬 계산에 관여하는 연산자를 초기화하고, 로컬 init() 메서드를 호출하여 태스크 특정 초기화를 수행합니다. 태스크 특정이란 태스크 타입(SourceTask, OneInputStreamTask 또는 TwoInputStreamTask 등)에 따라 이 단계가 다를 수 있음을 의미하지만, 어떤 경우든 필요한 태스크 전체 리소스를 여기서 획득합니다. 예를 들어 단일 입력 스트림을 기대하는 태스크를 나타내는 OneInputStreamTask는 로컬 태스크와 관련된 입력 스트림의 서로 다른 파티션 위치에 대한 연결을 초기화합니다.
필요한 리소스를 획득했다면 이제 서로 다른 연산자와 사용자 정의 함수가 위에서 검색한 태스크 전체 상태에서 각자의 개별 상태를 획득할 차례입니다. 이는 각 개별 연산자의 initializeState()를 호출하는 initializeState() 메서드에서 수행됩니다. 이 메서드는 모든 상태 기반 연산자가 재정의해야 하며, 작업이 처음 실행될 때와 태스크가 실패에서 복구되거나 savepoint를 사용할 때 모두에 대한 상태 초기화 로직을 포함해야 합니다.
이제 태스크의 모든 연산자가 초기화되었으므로 각 개별 연산자의 open() 메서드는 StreamTask의 openAllOperators() 메서드에 의해 호출됩니다. 이 메서드는 타이머 서비스에 검색된 타이머를 등록하는 것 같은 모든 운영 초기화를 수행합니다. 단일 태스크는 하나가 전임자의 출력을 소비하는 여러 연산자를 실행할 수 있습니다. 이 경우 open() 메서드는 마지막 연산자, 즉 출력이 태스크 자체의 출력이기도 한 연산자부터 첫 번째 연산자까지 호출됩니다. 이는 첫 번째 연산자가 태스크 입력 처리를 시작할 때 모든 다운스트림 연산자가 출력을 받을 준비가 되도록 하기 위해서입니다.
태스크의 연속된 연산자는 마지막부터 첫 번째까지 열립니다.
이제 태스크는 실행을 재개할 수 있고 연산자는 새로운 입력 데이터 처리를 시작할 수 있습니다. 이곳이 태스크 특정 run() 메서드가 호출되는 곳입니다. 이 메서드는 더 이상 입력 데이터가 없거나(유한 스트림) 태스크가 취소될 때까지(수동이든 아니든) 실행됩니다. 여기서 연산자 특정의 processElement()와 processWatermark() 메서드가 호출됩니다.
완료까지 실행되는 경우, 즉 처리할 입력 데이터가 더 이상 없을 때, run() 메서드에서 나온 후 태스크는 종료 프로세스에 들어갑니다. 처음에 타이머 서비스는 새 타이머 등록(예: 실행 중인 실행된 타이머에서)을 중지하고, 아직 시작되지 않은 타이머를 모두 지우며, 현재 실행 중인 타이머의 완료를 기다립니다. 그런 다음 finishAllOperators()는 각 연산자의 finish() 메서드를 호출하여 계산에 관여하는 연산자에게 알립니다. 그런 다음 버퍼링된 출력 데이터가 플러시되어 다운스트림 태스크가 처리할 수 있습니다. 그런 다음 최종 checkpoint가 활성화되면 태스크는 최종 checkpoint 완료를 기다려 2단계 커밋을 사용하는 연산자가 모든 레코드를 커밋했는지 확인합니다. 마지막으로 태스크는 각각의 close() 메서드를 호출하여 연산자가 보유한 모든 리소스를 정리하려고 시도합니다. 서로 다른 연산자를 열 때 순서가 마지막부터 첫 번째라고 언급했습니다. 닫기는 반대 방식인 첫 번째부터 마지막으로 발생합니다.
태스크의 연속된 연산자는 첫 번째부터 마지막까지 닫힙니다.
마지막으로 모든 연산자가 닫히고 모든 리소스가 해제되면 태스크는 타이머 서비스를 종료하고 태스크 특정 정리를 수행합니다. 예: 모든 내부 버퍼 정리. 그런 다음 모든 출력 채널을 닫고 출력 버퍼를 정리하는 일반 태스크 정리를 수행합니다.
Checkpoints: 앞서 initializeState() 동안 실패에서 복구하는 경우 태스크와 모든 연산자 및 함수는 실패 전 마지막 성공 checkpoint 동안 안정적인 저장소에 유지된 상태를 검색한다는 것을 보았습니다. Flink의 checkpoint는 사용자 지정 간격에 기반해 주기적으로 수행되며, 메인 태스크 스레드와 다른 스레드에 의해 수행됩니다. 그래서 태스크 수명주기의 주요 단계에 포함되지 않습니다. 간단히 말해 CheckpointBarriers라는 특수 요소가 작업의 소스 태스크에 의해 입력 데이터 스트림에 주기적으로 주입되며, 실제 데이터와 함께 소스에서 싱크로 이동합니다. 소스 태스크는 실행 모드에 들어간 후 CheckpointCoordinator도 실행 중이라고 가정할 때 이 barrier를 주입합니다. 태스크가 그러한 barrier를 받을 때마다 checkpoint 스레드가 수행할 작업을 예약하며, 이 스레드는 태스크의 연산자의 snapshotState()를 호출합니다. checkpoint가 수행되는 동안 태스크는 입력 데이터를 계속 받을 수 있지만 데이터는 버퍼링되며 checkpoint가 성공적으로 완료된 후에만 처리되고 다운스트림으로 출력됩니다.
중단된 실행
이전 섹션에서 완료까지 실행되는 태스크의 수명주기를 설명했습니다. 어느 시점에 태스크가 취소되면 정상 실행이 중단되고 그 시점부터 수행되는 유일한 작업은 위에서 설명한 대로 타이머 서비스 종료, 태스크 특정 정리, 연산자 닫기, 일반 태스크 정리입니다.