탄력적 스케일링

탄력적 스케일링 (Elastic Scaling)

역사적으로 작업의 병렬도는 수명주기 전체에 걸쳐 정적이었고 제출 시점에 한 번 정의됐어요. 배치 작업은 전혀 재조정할 수 없었고, 스트리밍 작업은 savepoint로 중지한 뒤 다른 병렬도로 재시작할 수 있었어요. 이 페이지는 Flink가 런타임에 작업의 병렬도를 조정할 수 있게 하는 새로운 종류의 스케줄러를 설명하며, 이는 Flink를 진정한 클라우드 네이티브 스트림 프로세서에 한 걸음 더 가깝게 만들어요. 새 스케줄러는 Adaptive Scheduler(스트리밍)와 Adaptive Batch Scheduler(배치)예요.

출처: 문서

본문

Adaptive Scheduler

Adaptive Scheduler는 사용 가능한 슬롯에 기반해 작업의 병렬도를 조정할 수 있어요. 원래 구성된 병렬도로 작업을 실행하기에 슬롯이 충분하지 않으면 자동으로 병렬도를 줄여요. 이는 제출 시점에 리소스가 충분하지 않거나 작업 실행 중 TaskManager가 중단된 경우 때문일 수 있어요. 새 슬롯이 생기면 작업은 구성된 병렬도까지 다시 확장돼요.

Reactive Mode(아래 참고)에서는 구성된 병렬도가 무시되고 무한대로 설정된 것처럼 취급되어 작업이 항상 가능한 한 많은 리소스를 사용하게 돼요.

Adaptive Scheduler가 기본 스케줄러보다 좋은 점 하나는 TaskManager 손실을 우아하게 처리할 수 있다는 것이에요. 그런 경우 그냥 축소하면 되기 때문이에요.

Adaptive Scheduler는 선언적 리소스 관리(Declarative Resource Management)라는 기능 위에 구축돼요. 정확한 슬롯 수를 요청하는 대신 JobMaster가 원하는 리소스를(reactive 모드에서는 최대가 무한대로 설정됨) ResourceManager에 선언하고, ResourceManager가 그 리소스를 충족하려고 해요.

JobMaster가 런타임 중 더 많은 리소스를 얻으면 사용 가능한 최신 savepoint로 작업을 자동으로 재조정하여 외부 오케스트레이션의 필요를 없애요. Flink 1.18.x부터 Externalized Declarative Resource Management를 사용해 실행 중인 작업의 리소스 요구사항을 다시 선언할 수 있어요. 그렇지 않으면 Adaptive Scheduler는 입력 속도 변화나 워크로드 성능 변화로 작업을 재조정해야 하는 경우를 처리할 수 없어요.

Externalized Declarative Resource Management

Externalized Declarative Resource Management는 MVP(“minimum viable product”) 기능이에요. Flink 커뮤니티는 메일링 리스트를 통해 사용자 피드백을 적극적으로 구하고 있어요. 이 페이지에 나열된 제한 사항을 확인해주세요.

Apache Flink Kubernetes operator와 함께 Externalized Declarative Resource Management를 사용해 완전한 자동 스케일링 경험을 얻을 수 있어요.

Externalized Declarative Resource Management는 두 가지 배포 시나리오를 다루는 것을 목표로 해요:

  • Session Cluster의 Adaptive Scheduler: 여러 작업이 리소스를 놓고 경쟁할 수 있으며, 작업 간 리소스 분배를 더 세밀하게 제어해야 하는 경우
  • Application Cluster의 Adaptive Scheduler + Active Resource Manager(예: Native Kubernetes): Flink가 "탐욕적으로" 새 TaskManager를 만들도록 의존하면서도 Reactive Mode처럼 재조정 기능을 활용하려는 경우

REST API 엔드포인트를 도입해 vertex별 병렬도 경계를 설정함으로써 실행 중인 작업의 리소스 요구사항을 다시 선언할 수 있어요.

PUT /jobs/<job-id>/resource-requirements

REQUEST BODY:
{
"<first-vertex-id>": {
"parallelism": {
"lowerBound": 3,
"upperBound": 5
}
},
"<second-vertex-id>": {
"parallelism": {
"lowerBound": 2,
"upperBound": 3
}
}
}

어느 정도 위 엔드포인트는 "재조정 엔드포인트"로 생각할 수 있으며 Flink를 위한 자동 스케일링 경험을 구축하는 중요한 구성 요소를 도입해요. Flink UI의 Job 개요를 탐색하고 작업 목록의 up-scale/down-scale 버튼을 사용해 이 기능을 수동으로 시도할 수 있어요.

사용법 (Usage)

session cluster에서 Adaptive Scheduler를 사용할 때, 클러스터에 리소스가 충분하지 않으면 같은 세션의 여러 실행 작업 간 슬롯 분배에 대한 보장이 없어요. External Declarative Resource Management가 이 문제를 부분적으로 완화할 수 있지만, application cluster에서 Adaptive Scheduler를 사용하는 것이 여전히 권장돼요.

적응형 스케줄러가 기본 스케줄러 대신 사용되도록 jobmanager.scheduler를 클러스터 수준에서 설정해야 해요.

jobmanager.scheduler: adaptive

Adaptive Scheduler의 동작은 이름에 jobmanager.adaptive-scheduler 접두사가 붙은 모든 구성 옵션으로 구성돼요.

제한 사항 (Limitations)

  • 스트리밍 작업만 지원: Adaptive Scheduler는 스트리밍 작업에서만 동작해요. 배치 작업을 제출하면 Flink는 배치 작업의 기본 스케줄러, 즉 Adaptive Batch Scheduler를 사용해요.
  • 부분 장애 조치 미지원: 부분 장애 조치는 스케줄러가 실패한 작업의 일부(Flink 내부의 "regions")를 전체 작업 대신 재시작할 수 있음을 의미해요. 이 제한은 거의 완전 병렬(embarrassingly parallel) 작업의 복구 시간에만 영향을 줘요. Flink의 기본 스케줄러는 실패한 부분을 재시작할 수 있지만 Adaptive Scheduler는 전체 작업을 재시작해요.
  • 스케일링 이벤트는 작업·태스크 재시작을 트리거해 Task 시도 수를 늘려요.

Reactive Mode

Reactive Mode는 Adaptive Scheduler의 특수 모드로, 클러스터당 단일 작업을 가정해요(Application Mode가 강제). Reactive Mode는 작업이 항상 클러스터의 모든 리소스를 사용하도록 구성해요. TaskManager를 추가하면 작업이 확장되고, 리소스를 제거하면 축소돼요. Flink는 작업의 병렬도를 관리하며 항상 가능한 한 가장 높은 값으로 설정해요.

Reactive Mode는 재조정 이벤트에서 작업을 재시작하고 최신 완료 체크포인트에서 복원해요. 이는 savepoint를 만들기 위한 오버헤드가 없음을 의미해요(수동 재조정에 필요). 또한 재조정 후 다시 처리되는 데이터 양은 체크포인팅 간격에 따라 달라지고, 복원 시간은 상태 크기에 따라 달라져요.

Reactive Mode는 외부 서비스가 consumer lag, 총 CPU 사용률, 처리량, 지연 같은 특정 메트릭을 모니터링하게 해 Flink 사용자가 강력한 자동 스케일링 메커니즘을 구현할 수 있게 해줘요. 이 메트릭이 특정 임계값 위/아래가 되면 Flink 클러스터에 TaskManager를 추가하거나 제거할 수 있어요. 이는 Kubernetes 배포의 replica factor를 바꾸거나 AWS의 autoscaling group으로 구현할 수 있어요. 이 외부 서비스는 리소스 할당·해제만 처리하면 돼요. Flink가 사용 가능한 리소스로 작업이 실행되도록 관리할게요.

시작하기 (Getting started)

Reactive Mode를 시도해보려면 다음 지침을 따르세요. 단일 머신에 Flink를 배포한다고 가정해요.

# these instructions assume you are in the root directory of a Flink distribution.

# Put Job into lib/ directory
cp ./examples/streaming/TopSpeedWindowing.jar lib/
# Submit Job in Reactive Mode
./bin/standalone-job.sh start -Dscheduler-mode=reactive -Dexecution.checkpointing.interval="10s" -j org.apache.flink.streaming.examples.windowing.TopSpeedWindowing
# Start first TaskManager
./bin/taskmanager.sh start

사용된 제출 명령을 빠르게 살펴볼게요:

  • ./bin/standalone-job.sh start는 Flink를 Application Mode로 배포
  • -Dscheduler-mode=reactive는 Reactive Mode 활성화
  • -Dexecution.checkpointing.interval="10s"는 체크포인팅과 재시작 전략 구성
  • 마지막 인자는 Job의 main 클래스 이름 전달

이제 Reactive Mode로 Flink 작업을 시작했어요. 웹 인터페이스에서 하나의 TaskManager에서 작업이 실행 중인 것을 보여줘요. 작업을 확장하려면 클러스터에 TaskManager를 추가하면 돼요:

# Start additional TaskManager
./bin/taskmanager.sh start

축소하려면 TaskManager 인스턴스를 제거해요:

# Remove a TaskManager
./bin/taskmanager.sh stop

사용법 (Usage)

구성 (Configuration)

Reactive Mode를 활성화하려면 scheduler-mode를 reactive로 구성해야 해요.

작업의 개별 연산자의 병렬도는 스케줄러가 결정해요. 구성할 수 없으며, 개별 연산자나 전체 작업에 명시적으로 설정해도 무시돼요. 병렬도에 영향을 줄 수 있는 유일한 방법은 연산자에 max parallelism을 설정하는 것(스케줄러가 존중)이에요. maxParallelism은 2^15(32768)로 제한돼요. 개별 연산자나 전체 작업에 max parallelism을 설정하지 않으면 기본 병렬도 규칙이 적용되어 최대 가능 값보다 낮은 하한을 적용할 수 있어요. 기본 스케줄링 모드와 마찬가지로 병렬도 모범 사례를 고려해주세요.

이렇게 높은 max parallelism은 Flink의 일부 내부 구조를 유지하는 데 더 많은 내부 구조가 필요하므로 작업 성능에 영향을 줄 수 있다는 점에 주의해요.

Reactive Mode를 활성화하면 jobmanager.adaptive-scheduler.resource-wait-timeout 구성 키가 기본값 -1이 돼요. 이는 JobManager가 충분한 리소스를 기다리며 영원히 실행됨을 의미해요. 작업 실행에 충분한 TaskManager가 없으면 일정 시간 후 JobManager가 멈추게 하려면 jobmanager.adaptive-scheduler.resource-wait-timeout을 구성해요.

Reactive Mode가 활성화되면 jobmanager.adaptive-scheduler.resource-stabilization-timeout 구성 키가 기본값 0이 돼요: Flink는 충분한 리소스가 있으면 즉시 작업을 실행하기 시작해요. TaskManager가 동시에 연결되지 않고 하나씩 천천히 연결되는 시나리오에서는 이 동작이 TaskManager가 연결될 때마다 작업 재시작으로 이어져요. 작업을 스케줄링하기 전에 리소스가 안정화되기를 기다리려면 이 구성 값을 늘리세요.

추가로 jobmanager.adaptive-scheduler.min-parallelism-increase를 구성할 수 있어요: 이 구성 옵션은 확장을 트리거하기 전에 필요한 추가 집계 병렬도 증가의 최소량을 지정해요. 예를 들어 source(병렬도 2)와 sink(병렬도 2)가 있는 작업의 경우 집계 병렬도는 4예요. 기본적으로 구성 키는 1로 설정되어 있어 집계 병렬도의 어떤 증가든 재시작을 트리거해요.

jobmanager.adaptive-scheduler.scaling-interval.max를 설정해 스케일링 연산을 강제로 발생시킬 수 있어요. 기본적으로 비활성화돼요. 설정하면 클러스터에 새 리소스가 추가될 때 jobmanager.adaptive-scheduler.min-parallelism-increase가 충족되지 않더라도 jobmanager.adaptive-scheduler.scaling-interval.max 후에 재조정이 스케줄링돼요.

너무 빈번한 스케일링 연산을 피하려면 jobmanager.adaptive-scheduler.scaling-interval.min을 구성해 두 스케일링 연산 사이의 최소 시간을 설정할 수 있어요. 기본값은 30초예요.

권장 사항 (Recommendations)

  • 상태 기반 작업에는 주기적 체크포인팅 구성: Reactive mode는 재조정 이벤트에서 최신 완료 체크포인트에서 복원해요. 주기적 체크포인팅이 없으면 프로그램이 상태를 잃어요. 체크포인팅은 또한 재시작 전략을 구성해요. Reactive Mode는 구성된 재시작 전략을 존중해요. 재시작 전략이 구성되지 않으면 reactive mode는 작업을 스케일링하는 대신 실패시켜요.
  • TaskManager가 제대로 종료되지 않으면(i.e. SIGTERM 신호 대신 SIGKILL 신호를 사용하면) Reactive Mode의 축소가 더 오래 걸릴 수 있어요. 이 경우 Flink는 JobManager와 중단된 TaskManager 사이의 하트비트가 타임아웃되기를 기다려요. 작업을 더 낮은 병렬도로 재배포하기 전에 Flink 작업이 약 50초 동안 멈춰 있는 것을 볼 수 있어요.

기본 타임아웃은 50초로 구성돼요. 인프라가 허용하면 heartbeat.timeout 구성을 더 낮은 값으로 조정해요. 하트비트 타임아웃을 낮게 설정하면 TaskManager가 네트워크 혼잡이나 긴 가비지 컬렉션 일시정지 같은 이유로 하트비트에 응답하지 못하면 실패로 이어질 수 있어요. heartbeat.interval은 항상 타임아웃보다 낮아야 한다는 점에 주의해요.

제한 사항 (Limitations)

Reactive Mode는 새롭고 실험적인 기능이므로 기본 스케줄러가 지원하는 모든 기능이 Reactive Mode(그리고 그 적응형 스케줄러)에서도 사용 가능한 것은 아니에요. Flink 커뮤니티는 이러한 제한을 해결하기 위해 작업 중이에요.

  • 배포는 standalone 애플리케이션 배포로만 지원. 활성 리소스 제공자(네이티브 Kubernetes, YARN 등)는 명시적으로 지원되지 않아요. Standalone 세션 클러스터도 지원되지 않아요. 애플리케이션 배포는 단일 작업 애플리케이션으로 제한돼요.

지원되는 유일한 배포 옵션은 Application Mode의 Standalone(이 페이지의 getting started 참고), Application Mode의 Docker, Standalone Kubernetes Application Cluster예요.

Adaptive Scheduler의 제한 사항도 Reactive Mode에 적용돼요.

재조정 히스토리 (Rescale History)

Flink 2.3 이전에는 사용자와 개발자가 AdaptiveScheduler 재조정 히스토리의 내부 세부사항을 조사할 수 없어 운영상 불편함이 있었어요. 예를 들어 사용자는 특정 리소스 변화, 병렬도 조정, 재조정 과정 중 각 내부 상태 전이에 걸린 시간에 대한 가시성이 필요해요. 이 정보는 재조정에서 더 낮은 지연과 더 높은 안정성을 얻도록 파라미터를 튜닝하는 데 중요해요.

따라서 Flink 커뮤니티는 재조정 히스토리 기록·저장을 지원하는 FLIP-495와 REST API를 통한 조회와 Web UI 표시를 활성화하는 FLIP-487을 도입했어요.

AdaptiveScheduler가 활성화된 스트림 작업에 대해 다음 구성 항목을 양의 정수로 설정해 재조정 히스토리를 활성화할 수 있어요. 이 값은 작업에 보존되는 최근 재조정 레코드 수를 나타내요.

구성 옵션의 기본값은 0이에요. 구성 값이 0보다 작거나 같으면 이 기능은 비활성화돼요.

재조정 히스토리의 정보와 스타일 (The Information and Style About Rescale History)

Flink 2.3부터 Web UI에 Rescales 표시용 페이지가 도입됐어요. Checkpoints 페이지와 같은 계층 수준에 위치하며 비슷한 스타일을 가져요. 주로 다음 하위 페이지를 포함해요:

  • Overview: 이 하위 페이지는 다양한 재조정 종료 상태에 걸친 최근 재조정 레코드와, 작업 시작 이후 총 재조정 수, 실패·성공 수 같은 기본 작업 재조정 통계를 표시해요. 또한 페이지는 상세 재조정 정보 표시를 지원해요.
  • History: 이 하위 페이지는 가장 최근 재조정 레코드(구성된 web.adaptive-scheduler.rescale-history.size 한도까지)의 요약 정보를 표시해요. 또한 아래에 설명된 상세 재조정 정보 표시를 지원해요:
    • 재조정의 기본 정보
      • Rescale UUID: 재조정의 고유 ID로 32개의 16진 문자로 구성(아래 UUID 정의는 여기와 동일)
      • Attempt ID: 같은 작업 리소스 요구사항에서 트리거된 재조정 시도 수
      • Requirements ID: 리소스 요구사항의 고유 UUID
      • Trigger Cause: 재조정을 트리거한 이유
      • Terminal State: 재조정의 종료 상태
      • Terminated Reason: 재조정 수명주기 종료 이유
      • Start Time: 재조정의 시작 시간
      • Duration: 재조정 시작부터 완료까지의 기간, 또는 재조정 연산이 아직 완료되지 않았다면 현재까지의 기간
      • End Time: 재조정이 종료됐다면 재조정의 종료 시간, 그렇지 않으면 현재 시간
    • Job Vertex별 기본 속성과 재조정 변경
      • ID: 대상 Job Vertex의 고유 UUID
      • Name: 대상 vertex의 짧은 이름
      • Slot Sharing Group ID: 대상 Slot Sharing Group의 고유 UUID
      • Previous Parallelism: 현재 재조정 전의 대상 vertex 병렬도
      • Acquired Parallelism: 현재 재조정 후의 대상 vertex 병렬도
      • Sufficient Parallelism: 원하는 병렬도에 도달할 수 없어도 재조정 연산이 진행될 수 있게 했을 vertex의 최소 병렬도
      • Desired Parallelism: 재조정 연산을 트리거한 초기 변경 요청에서 지정된 Job Vertex의 원하는 병렬도
    • Slot Sharing Group별 기본 속성과 재조정 변경
      • Slot Sharing Group ID: 대상 슬롯이 속한 Slot Sharing Group의 UUID
      • Slot Sharing Group Name: 슬롯이 속한 Slot Sharing Group의 이름
      • Previous Slot: 재조정 전의 슬롯 수
      • Acquired Slot: 재조정 후의 슬롯 수
      • Desired Slot: 재조정의 원하는 슬롯 수
      • Sufficient Slot: 재조정에서 작업을 배포하는 최소 슬롯 수
      • Request Profile: 재조정에서 Slot Sharing Group의 요청 리소스 프로파일
      • Acquired Profile: 재조정에서 Slot Sharing Group의 획득 리소스 프로파일
    • 재조정 내 AdaptiveScheduler의 내부 Scheduler State History(자세한 내용은 FLIP-160의 AdaptiveScheduler 상태 참고)
      • State: 스케줄러 상태 이름
      • Enter Time: 상태에 들어간 시간
      • Leave Time: 상태를 떠난 시간
      • Duration: 상태에서 보낸 시간(Leave Time − Enter Time)
      • Exception: 상태 내 현재 재조정에 대한 예외 정보
  • Summary: 이 하위 페이지는 작업 시작 이후 발생한 총 재조정 이벤트 수와 각각의 실패·성공 수를 표시해요. 또한 Min, Max, Avg, P50 메트릭 등을 포함해 재조정 상태별로 분류된 재조정 기간 통계 같은 재조정 히스토리의 통계 요약을 제공해요.
  • Configuration: 이 하위 페이지는 현재 스트리밍 작업의 재조정 연산 중 AdaptiveScheduler가 사용하는 관련 파라미터 값을 표시해요.

더 많은 세부사항 (More details)

자세한 내용은 FLIP-495FLIP-487을 참고해요.

더 알아보기 (Learn more)