Flink 아키텍처
Flink 아키텍처 (Flink Architecture)
Flink는 분산 시스템으로, 스트리밍 애플리케이션을 실행하려면 컴퓨팅 리소스의 효과적인 할당과 관리가 필요해요. Hadoop YARN과 Kubernetes 같은 모든 일반적인 클러스터 리소스 매니저와 통합되지만, standalone 클러스터로 또는 심지어 라이브러리로 실행되도록 설정할 수도 있어요.
출처: 문서
본문
이 섹션은 Flink 아키텍처의 개요를 포함하고, 핵심 컴포넌트들이 애플리케이션을 실행하고 장애에서 복구하기 위해 어떻게 상호작용하는지 설명해요.
Flink 클러스터의 구조 (Anatomy of a Flink Cluster)
Flink 런타임은 두 가지 유형의 프로세스로 구성돼요: JobManager와 하나 이상의 TaskManager.
Client는 런타임과 프로그램 실행의 일부가 아니지만, 데이터플로우를 준비해 JobManager로 보내는 데 사용돼요. 그 후 클라이언트는 연결을 끊거나(detached mode), 연결된 상태로 진행 상황 보고를 받을 수 있어요(attached mode). 클라이언트는 실행을 트리거하는 Java 프로그램의 일부로, 또는 명령줄 프로세스 ./bin/flink run ...으로 실행돼요.
JobManager와 TaskManager는 다양한 방식으로 시작될 수 있어요: 머신에 직접 standalone 클러스터로, 컨테이너 안에서, 또는 YARN 같은 리소스 프레임워크로 관리되면서요. TaskManager는 JobManager에 연결해 자신이 사용 가능하다고 알리고, 작업을 할당받아요.
JobManager
JobManager는 Flink 애플리케이션의 분산 실행을 조정하는 데 관련된 여러 책임을 가져요: 다음 태스크(또는 태스크 집합)를 언제 스케줄링할지 결정하고, 완료된 태스크나 실행 실패에 반응하며, 체크포인트를 조정하고, 장애 시 복구를 조정해요. 이 프로세스는 세 가지 다른 컴포넌트로 구성돼요:
-
ResourceManager
ResourceManager는 Flink 클러스터에서 리소스의 할당/해제와 프로비저닝을 담당해요 — Flink 클러스터에서 리소스 스케줄링의 단위인 **태스크 슬롯(task slots)**을 관리해요 (TaskManagers 참고). Flink는 YARN, Kubernetes, standalone 배포 같은 다양한 환경과 리소스 프로바이더에 대해 여러 ResourceManager를 구현해요. standalone 설정에서 ResourceManager는 사용 가능한 TaskManager의 슬롯만 분배할 수 있고, 스스로 새 TaskManager를 시작할 수 없어요.
-
Dispatcher
Dispatcher는 Flink 애플리케이션을 실행을 위해 제출하는 REST 인터페이스를 제공하고, 제출된 각 작업에 대해 새 JobMaster를 시작해요. 또한 작업 실행에 대한 정보를 제공하는 Flink WebUI도 실행해요.
-
JobMaster
JobMaster는 단일 JobGraph의 실행을 관리하는 책임을 가져요. 여러 작업이 Flink 클러스터에서 동시에 실행될 수 있으며, 각각 고유한 JobMaster를 가져요.
항상 최소 하나의 JobManager가 있어요. 고가용성(HA) 설정은 여러 JobManager를 가질 수 있는데, 그중 하나는 항상 *리더(leader)*이고 나머지는 대기(standby) 상태예요 (High Availability (HA) 참고).
TaskManagers
TaskManagers(또는 workers)는 데이터플로우의 태스크를 실행하고, 데이터 스트림을 버퍼링하고 교환해요.
항상 최소 하나의 TaskManager가 있어야 해요. TaskManager에서 리소스 스케줄링의 가장 작은 단위는 태스크 *슬롯(slot)*이에요. TaskManager의 태스크 슬롯 수는 동시 처리 태스크 수를 나타내요. 여러 연산자가 하나의 태스크 슬롯에서 실행될 수 있다는 점에 유의하세요 (Tasks and Operator Chains 참고).
태스크와 연산자 체인 (Tasks and Operator Chains)
분산 실행을 위해 Flink는 연산자 서브태스크를 *태스크(task)*로 *체이닝(chains)*해요. 각 태스크는 하나의 스레드로 실행돼요. 연산자를 태스크로 체이닝하는 것은 유용한 최적화예요: 스레드 간 핸드오버와 버퍼링 오버헤드를 줄이고, 지연 시간을 낮추면서 전체 처리량을 높여요. 체이닝 동작은 구성할 수 있어요; 자세한 내용은 chaining docs를 참조하세요.
아래 그림의 샘플 데이터플로우는 다섯 개의 서브태스크로, 따라서 다섯 개의 병렬 스레드로 실행돼요.
태스크 슬롯과 리소스 (Task Slots and Resources)
각 워커(TaskManager)는 JVM 프로세스이며, 별도의 스레드에서 하나 이상의 서브태스크를 실행할 수 있어요. TaskManager가 수용하는 태스크 수를 제어하기 위해, 소위 태스크 슬롯(최소 하나)을 가져요.
각 태스크 슬롯은 TaskManager 리소스의 고정된 부분 집합을 나타내요. 예를 들어 세 개의 슬롯을 가진 TaskManager는 매니지드 메모리의 1/3을 각 슬롯에 할당해요. 리소스를 슬롯으로 나누는 것은 서브태스크가 다른 작업의 서브태스크와 매니지드 메모리를 놓고 경쟁하지 않고, 일정량의 예약된 매니지드 메모리를 갖는다는 것을 의미해요. 여기서 CPU 격리는 일어나지 않는다는 점에 유의하세요; 현재 슬롯은 태스크의 매니지드 메모리만 분리해요.
태스크 슬롯 수를 조정해 사용자는 서브태스크가 서로 어떻게 격리되는지 정의할 수 있어요. TaskManager당 하나의 슬롯을 가지는 것은 각 태스크 그룹이 별도의 JVM(예를 들어 별도의 컨테이너에서 시작될 수 있는)에서 실행된다는 뜻이에요. 여러 슬롯을 가지는 것은 더 많은 서브태스크가 같은 JVM을 공유한다는 뜻이에요. 같은 JVM의 태스크들은 TCP 연결(멀티플렉싱을 통해)과 하트비트 메시지를 공유해요. 데이터 집합과 데이터 구조도 공유할 수 있어 태스크당 오버헤드를 줄여요.
기본적으로 Flink는 서브태스크가 다른 태스크의 서브태스크라도, 같은 작업에서 왔다면 슬롯을 공유하도록 허용해요. 결과적으로 하나의 슬롯이 작업의 전체 파이프라인을 담을 수 있어요. 이 *슬롯 공유(slot sharing)*를 허용하는 데는 두 가지 주요 이점이 있어요:
- Flink 클러스터는 작업에서 사용되는 가장 높은 병렬도만큼 정확히 많은 태스크 슬롯이 필요해요. 프로그램이 (병렬도가 다양한) 태스크를 총 몇 개 포함하는지 계산할 필요가 없어요.
- 더 나은 리소스 활용을 얻기 쉽다. 슬롯 공유가 없으면, 리소스를 많이 쓰지 않는 source/map() 서브태스크가 리소스 집약적인 window 서브태스크만큼 많은 리소스를 차단해요. 슬롯 공유가 있으면, 예시에서 기본 병렬도를 2에서 6으로 늘리면 무거운 서브태스크가 TaskManager 사이에 공정하게 분배되도록 하면서 슬롯 리소스를 완전히 활용해요.
Flink 애플리케이션 실행 (Flink Application Execution)
Flink Application은 main() 메서드에서 하나 또는 여러 Flink 작업을 생성하는 모든 사용자 프로그램이에요. 이 작업들의 실행은 로컬 JVM(LocalEnvironment) 또는 여러 머신으로 이루어진 클러스터의 원격 설정(RemoteEnvironment)에서 일어날 수 있어요. 각 프로그램에 대해 ExecutionEnvironment는 작업 실행을 제어하는(예: 병렬도 설정) 메서드와 외부 세계와 상호작용하는 메서드를 제공해요 (Flink Program의 구조 참고).
Flink Application의 작업들은 장기 실행되는 Flink Session Cluster, 전용 Flink Job Cluster (deprecated), 또는 Flink Application Cluster에 제출될 수 있어요. 이 옵션들의 차이는 주로 클러스터의 수명 주기(lifecycle)와 리소스 격리 보장과 관련돼요.
Flink Application Cluster
- 클러스터 수명 주기: Flink Application Cluster는 하나의 Flink Application에서만 작업을 실행하고,
main()메서드가 클라이언트가 아닌 클러스터에서 실행되는 전용 Flink 클러스터예요. 작업 제출은 원스텝 프로세스예요: 먼저 Flink 클러스터를 시작한 다음 기존 클러스터 세션에 작업을 제출할 필요가 없어요; 대신 애플리케이션 로직과 의존성을 실행 가능한 작업 JAR에 패키징하고, 클러스터 엔트리포인트(ApplicationClusterEntryPoint)가main()메서드를 호출해 JobGraph를 추출하는 책임을 져요. 이것은 예를 들어 Kubernetes에서 Flink Application을 다른 애플리케이션처럼 배포할 수 있게 해줘요. 따라서 Flink Application Cluster의 수명은 Flink Application의 수명에 묶여 있어요. - 리소스 격리: Flink Application Cluster에서 ResourceManager와 Dispatcher는 단일 Flink Application에 범위가 지정돼, Flink Session Cluster보다 관심사의 분리가 더 잘 돼요.
Flink Session Cluster
- 클러스터 수명 주기: Flink Session Cluster에서 클라이언트는 사전 존재하는 장기 실행 클러스터에 연결하며, 이 클러스터는 여러 애플리케이션 제출을 수용할 수 있어요. 모든 애플리케이션이 끝난 후에도 클러스터(그리고 JobManager)는 세션이 수동으로 중지될 때까지 계속 실행돼요. 따라서 Flink Session Cluster의 수명은 어떤 Flink Application이나 Job의 수명에도 묶여 있지 않아요.
- 리소스 격리: TaskManager 슬롯은 작업 제출 시 ResourceManager가 할당하고, 작업이 끝나면 해제돼요. 모든 작업이 같은 클러스터를 공유하므로, 클러스터 리소스에 대한 경쟁이 어느 정도 있어요 — 제출 작업 단계의 네트워크 대역폭 같은 것들이요. 이 공유 설정의 한계는 TaskManager 하나가 충돌하면 그 TaskManager에서 실행 중인 태스크가 있는 모든 작업이 실패한다는 것이에요; 비슷하게, JobManager에서 치명적 오류가 발생하면 클러스터에서 실행 중인 모든 작업에 영향을 주게 돼요.
- 기타 고려 사항: 사전 존재하는 클러스터를 가지는 것은 리소스를 요청하고 TaskManager를 시작하는 데 상당한 시간을 절약해줘요. 이것은 작업의 실행 시간이 매우 짧고 높은 시작 시간이 종단 간 사용자 경험에 부정적 영향을 주는 시나리오에서 중요해요 — 짧은 쿼리의 대화형 분석의 경우처럼, 작업이 기존 리소스를 사용해 신속하게 계산을 수행할 수 있는 것이 바람직해요.
이전에 Flink Session Cluster는
session mode의 Flink Cluster로도 알려져 있었어요.