세분화된 리소스 관리

세분화된 리소스 관리 (Fine-Grained Resource Management)

Apache Flink는 모든 애플리케이션에 대해 기본적으로 합리적인 기본 리소스 요구사항을 자동으로 도출하기 위해 많은 노력을 기울입니다. 자신의 특정 시나리오에 대한 지식을 바탕으로 리소스 소비를 미세 조정하려는 사용자를 위해 Flink는 **세분화된 리소스 관리(fine-grained resource management)**를 제공합니다. 이 페이지는 세분화된 리소스 관리의 사용법, 적용 시나리오, 그리고 동작 방식을 설명합니다.

출처: 문서

본문

참고: 이 기능은 현재 MVP("minimum viable product") 기능이며 DataStream API에서만 사용할 수 있습니다.

적용 시나리오 (Applicable Scenarios)

세분화된 리소스 관리로부터 이점을 얻을 수 있는 일반적인 시나리오는 다음과 같습니다.

  • 작업들이 병렬도가 크게 다른 경우.
  • 전체 파이프라인에 필요한 리소스가 단일 슬롯/TaskManager에 들어가기에는 너무 많은 경우.
  • 서로 다른 스테이지의 작업에 필요한 리소스가 크게 다른 배치 작업.

위 시나리오에서 세분화된 리소스 관리가 왜 리소스 효율을 개선할 수 있는지에 대한 심층 논의는 리소스 효율을 어떻게 개선하는가에서 제시됩니다.

동작 방식 (How it works)

Flink 아키텍처에서 설명한 것처럼, TaskManager의 작업 실행 리소스는 많은 슬롯으로 나뉩니다. 슬롯은 Flink 런타임에서 리소스 스케줄링과 리소스 요구사항의 기본 단위입니다.

세분화된 리소스 관리에서는 슬롯 요청이 사용자가 지정할 수 있는 특정 리소스 프로필을 포함합니다. Flink는 사용자가 지정한 리소스 요구사항을 존중하고 TaskManager의 사용 가능한 리소스에서 정확히 일치하는 슬롯을 동적으로 잘라냅니다. 0.25 Core와 1GB 메모리의 슬롯 요구사항이 있고 Flink가 이를 위해 Slot 1을 할당합니다.

이전에는 Flink에서 리소스 요구사항이 세분화된 리소스 프로필 없이 필요한 슬롯만 포함했는데, 이를 **조잡한 리소스 관리(coarse-grained resource management)**라고 합니다. TaskManager는 이러한 요구사항을 충족하기 위한 고정된 개수의 동일한 슬롯을 가졌습니다.

리소스 프로필이 지정되지 않은 리소스 요구사항에 대해 Flink는 리소스 프로필을 자동으로 결정합니다. 현재 그 리소스 프로필은 조잡한 리소스 관리에서와 같이 TaskManager의 총 리소스taskmanager.numberOfTaskSlots에서 계산됩니다. TaskManager의 총 리소스가 1 Core와 4 GB 메모리이고 task slots 수가 2로 설정되어 있다면, 리소스 프로필이 지정되지 않은 요구사항을 위해 0.5 Core와 2 GB 메모리의 Slot 2가 생성됩니다.

Slot 1Slot 2 할당 후 TaskManager에는 0.25 Core와 1 GB 메모리가 여유 리소스로 남습니다. 이 여유 리소스는 이후의 리소스 요구사항을 충족시키기 위해 더 파티셔닝될 수 있습니다.

자세한 내용은 리소스 할당 전략을 참조하세요.

사용법 (Usage)

세분화된 리소스 관리를 사용하려면 리소스 요구사항을 지정해야 합니다.

세분화된 리소스 요구사항은 슬롯 공유 그룹(slot sharing group)에 정의됩니다. 슬롯 공유 그룹은 JobManager에 그 안의 오퍼레이터/작업들이 같은 슬롯에 들어갈 수 있다고 알려주는 힌트입니다.

리소스 요구사항을 지정하려면 다음을 수행해야 합니다.

  • 슬롯 공유 그룹과 그 안에 포함된 오퍼레이터를 정의합니다.
  • 슬롯 공유 그룹의 리소스를 지정합니다.

슬롯 공유 그룹과 그 안에 포함된 오퍼레이터를 정의하는 방법은 두 가지가 있습니다.

  • 이름만으로 슬롯 공유 그룹을 정의하고 slotSharingGroup(String name)을 통해 오퍼레이터에 붙일 수 있습니다.
  • 이름과 슬롯 공유 그룹의 선택적 리소스 프로필을 포함하는 SlotSharingGroup 인스턴스를 구성할 수 있습니다. SlotSharingGroupslotSharingGroup(SlotSharingGroup ssg)를 통해 오퍼레이터에 붙일 수 있습니다.

슬롯 공유 그룹에 대한 리소스 프로필을 지정할 수 있습니다.

  • slotSharingGroup(SlotSharingGroup ssg)로 슬롯 공유 그룹을 설정했다면, SlotSharingGroup 인스턴스를 구성할 때 리소스 프로필을 지정할 수 있습니다.
  • slotSharingGroup(String name)으로만 슬롯 공유 그룹의 이름을 설정했다면, 같은 이름과 리소스 프로필로 SlotSharingGroup 인스턴스를 구성하고 StreamExecutionEnvironment#registerSlotSharingGroup(SlotSharingGroup ssg)로 그 리소스를 등록할 수 있습니다.

Java

final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

SlotSharingGroup ssgA = SlotSharingGroup.newBuilder("a")
  .setCpuCores(1.0)
  .setTaskHeapMemoryMB(100)
  .build();

SlotSharingGroup ssgB = SlotSharingGroup.newBuilder("b")
  .setCpuCores(0.5)
  .setTaskHeapMemoryMB(100)
  .build();

someStream.filter(...).slotSharingGroup("a") // Set the slot sharing group with name "a"
.map(...).slotSharingGroup(ssgB); // Directly set the slot sharing group with name and resource.

env.registerSlotSharingGroup(ssgA); // Then register the resource of group "a"

Python

env = StreamExecutionEnvironment.get_execution_environment()

ssg_a = SlotSharingGroup.builder('a') \
            .set_cpu_cores(1.0) \
            .set_task_heap_memory_mb(100) \
            .build()
ssg_b = SlotSharingGroup.builder('b') \
            .set_cpu_cores(0.5) \
            .set_task_heap_memory_mb(100) \
            .build()

some_stream.filter(...).slot_sharing_group('a') # Set the slot sharing group with name "a"
.map(...).slot_sharing_group(ssg_b) # Directly set the slot sharing group with name and resource.

env.register_slot_sharing_group(ssg_a) # Then register the resource of group "a"

참고: 각 슬롯 공유 그룹은 하나의 지정된 리소스에만 붙을 수 있으며, 어떤 충돌도 작업 컴파일을 실패시킵니다.

SlotSharingGroup을 구성할 때 슬롯 공유 그룹에 대해 다음 리소스 컴포넌트를 설정할 수 있습니다.

  • CPU Cores. 필요한 CPU 코어 수를 정의합니다. 양수 값으로 명시적으로 구성해야 합니다.
  • Task Heap Memory. 필요한 task heap 메모리 양을 정의합니다. 양수 값으로 명시적으로 구성해야 합니다.
  • Task Off-Heap Memory. 필요한 task off-heap 메모리 양을 정의하며, 0일 수 있습니다.
  • Managed Memory. 필요한 task managed 메모리 양을 정의하며, 0일 수 있습니다.
  • External Resources. 필요한 외부 리소스를 정의하며, 비어 있을 수 있습니다.

Java

// Directly build a slot sharing group with specific resource
SlotSharingGroup ssgWithResource =
    SlotSharingGroup.newBuilder("ssg")
        .setCpuCores(1.0) // required
        .setTaskHeapMemoryMB(100) // required
        .setTaskOffHeapMemoryMB(50)
        .setManagedMemory(MemorySize.ofMebiBytes(200))
        .setExternalResource("gpu", 1.0)
        .build();

// Build a slot sharing group without specific resource and then register the resource of it in StreamExecutionEnvironment
SlotSharingGroup ssgWithName = SlotSharingGroup.newBuilder("ssg").build();
env.registerSlotSharingGroup(ssgWithResource);

Python

# Directly build a slot sharing group with specific resource
ssg_with_resource = SlotSharingGroup.builder('ssg') \
            .set_cpu_cores(1.0) \
            .set_task_heap_memory_mb(100) \
            .set_task_off_heap_memory_mb(50) \
            .set_managed_memory(MemorySize.of_mebi_bytes(200)) \
            .set_external_resource('gpu', 1.0) \
            .build()

# Build a slot sharing group without specific resource and then register the resource of it in StreamExecutionEnvironment
ssg_with_name = SlotSharingGroup.builder('ssg').build()
env.register_slot_sharing_group(ssg_with_resource)

참고: 리소스 프로필을 지정하거나 지정하지 않고 SlotSharingGroup을 구성할 수 있습니다. 리소스 프로필을 지정하면 CPU coresTask Heap Memory를 양수 값으로 명시적으로 설정해야 하며, 다른 컴포넌트는 선택 사항입니다.

제한 사항 (Limitations)

세분화된 리소스 관리는 새롭고 실험적인 기능이므로, 기본 스케줄러가 지원하는 모든 기능이 여기에서도 사용 가능한 것은 아닙니다. Flink 커뮤니티가 이러한 제한을 해결하기 위해 노력 중입니다.

  • Elastic Scaling 미지원. elastic scaling은 현재 리소스가 지정되지 않은 슬롯 요청만 지원합니다.
  • Flink Web UI와의 제한된 통합. 세분화된 리소스 관리의 슬롯은 서로 다른 리소스 사양을 가질 수 있습니다. 웹 UI는 현재 세부 정보 없이 슬롯 번호만 표시합니다.
  • 배치 작업과의 제한된 통합. 현재 세분화된 리소스 관리는 배치 워크로드가 모든 엣지 타입이 BLOCKING인 것으로 실행되도록 요구합니다. 이를 위해 fine-grained.shuffle-mode.all-blockingtrue로 구성해야 합니다. 이는 성능에 영향을 줄 수 있음에 유의하세요. 자세한 내용은 FLINK-20865를 참조하세요.
  • 혼합 리소스 요구사항은 권장되지 않음. 작업의 일부 부분에만 리소스 요구사항을 지정하고 나머지는 지정하지 않는 것은 권장되지 않습니다. 현재 지정되지 않은 요구사항은 어떤 리소스의 슬롯으로도 충족될 수 있습니다. 실제로 획득한 리소스는 작업 실행이나 장애 조치마다 일치하지 않을 수 있습니다.
  • 슬롯 할당 결과가 최적이 아닐 수 있음. 슬롯 요구사항은 여러 차원의 리소스를 포함하므로 슬롯 할당은 사실상 NP-hard인 다차원 패킹(packing) 문제입니다. 기본 리소스 할당 전략은 최적의 슬롯 할당을 달성하지 못할 수 있고, 일부 시나리오에서 리소스 단편화나 리소스 할당 실패를 초래할 수 있습니다.

주의사항 (Notice)

  • 슬롯 공유 그룹을 설정하면 성능이 변경될 수 있음. 체인 가능한(chain-able) 오퍼레이터를 다른 슬롯 공유 그룹에 설정하면 오퍼레이터 체인이 깨질 수 있어 성능이 변경됩니다.
  • 슬롯 공유 그룹은 오퍼레이터의 스케줄링을 제한하지 않음. 슬롯 공유 그룹은 스케줄러에게 그룹화된 오퍼레이터들이 공유 슬롯에 배포될 수 있다는 것만 힌트로 알려줍니다. 스케줄러가 항상 그룹화된 오퍼레이터를 함께 배포한다는 보장은 없습니다. 그룹화된 오퍼레이터가 별도의 슬롯에 배포되는 경우, 슬롯 리소스는 지정된 그룹 요구사항에서 파생됩니다.

심층 분석 (Deep Dive)

리소스 효율을 어떻게 개선하는가 (How it improves resource efficiency)

이 섹션에서는 세분화된 리소스 관리가 리소스 효율을 어떻게 개선하는지 깊이 살펴봅니다. 이는 여러분의 작업에 이점이 될 수 있는지 이해하는 데 도움이 됩니다.

이전에 Flink는 작업을 미리 정의된, 보통 동일한 슬롯에 배포하는 조잡한 리소스 관리 방식을 채택했으며, 각 슬롯이 얼마나 많은 리소스를 포함하는지에 대한 개념이 없었습니다. 많은 작업에서 조잡한 리소스 관리를 사용하고 모든 작업을 하나의 슬롯 공유 그룹에 넣는 것은 리소스 활용 측면에서 충분히 잘 동작합니다.

  • 모든 작업이 같은 병렬도를 가진 많은 스트리밍 작업에서 각 슬롯은 전체 파이프라인을 포함합니다. 이상적으로 모든 파이프라인은 대략 같은 리소스를 사용해야 하며, 이는 동일한 슬롯의 리소스를 조정함으로써 쉽게 충족될 수 있습니다.
  • 작업의 리소스 소비는 시간에 따라 변합니다. 작업의 소비가 감소하면, 여분의 리소스는 소비가 증가하는 다른 작업이 사용할 수 있습니다. 이를 peak shaving과 valley filling 효과라고 하며, 전체적으로 필요한 리소스를 줄입니다.

하지만 조잡한 리소스 관리가 잘 동작하지 않는 경우도 있습니다.

  • 작업은 서로 다른 병렬도를 가질 수 있습니다. 때로는 이러한 병렬도 차이를 피할 수 없습니다. 예: source/sink/lookup 작업의 병렬도는 외부 업스트림/다운스트림 시스템의 파티션과 IO 부하에 의해 제한될 수 있습니다. 이런 경우 작업이 적은 슬롯은 전체 파이프라인의 작업을 가진 슬롯보다 적은 리소스가 필요합니다.
  • 때로는 전체 파이프라인에 필요한 리소스가 단일 슬롯/TaskManager에 넣기에는 너무 클 수 있습니다. 이런 경우 파이프라인을 여러 SSG로 나눠야 하며, 각각의 리소스 요구사항이 항상 같지는 않을 수 있습니다.
  • 배치 작업의 경우 모든 작업을 동시에 실행할 수는 없습니다. 따라서 파이프라인의 순간 리소스 요구사항은 시간에 따라 변합니다.

모든 작업을 동일한 슬롯으로 실행하려고 하면 비최적의 리소스 활용을 초래할 수 있습니다. 동일한 슬롯의 리소스는 가장 높은 리소스 요구사항을 충족할 수 있어야 하며, 이는 다른 요구사항에 대해 낭비가 될 수 있습니다. GPU 같은 비싼 외부 리소스가 관련되면 이러한 낭비는 감당하기 더 어려워질 수 있습니다. 세분화된 리소스 관리는 서로 다른 리소스의 슬롯을 활용해 이러한 시나리오에서 리소스 활용을 개선합니다.

리소스 할당 전략 (Resource Allocation Strategy)

이 섹션에서는 Flink 런타임의 슬롯 파티셔닝 메커니즘과 리소스 할당 전략에 대해 다룹니다. 여기에는 Flink 런타임이 슬롯을 잘라낼 TaskManager를 어떻게 선택하고, Native KubernetesYARN에서 TaskManager를 어떻게 할당하는지가 포함됩니다. 리소스 할당 전략은 Flink 런타임에서 플러그형이며, 여기서는 세분화된 리소스 관리의 첫 단계에서의 기본 구현을 소개합니다. 미래에는 사용자가 다양한 시나리오에 대해 선택할 수 있는 여러 전략이 있을 수 있습니다.

동작 방식 섹션에서 설명한 것처럼, Flink는 지정된 리소스가 있는 슬롯 요청을 위해 TaskManager에서 정확히 일치하는 슬롯을 잘라냅니다. 내부 과정은 위에 표시되어 있습니다. TaskManager는 총 리소스로 시작되지만 미리 정의된 슬롯 없이 시작됩니다. 0.25 Core와 1GB 메모리의 슬롯 요청이 도착하면 Flink는 충분한 여유 리소스가 있는 TaskManager를 선택하고 요청된 리소스로 새 슬롯을 만듭니다. 슬롯이 해제되면 그 리소스는 TaskManager의 사용 가능한 리소스로 반환됩니다.

현재 리소스 할당 전략에서 Flink는 등록된 모든 TaskManager를 순회하며 슬롯 요청을 충족할 충분한 여유 리소스가 있는 첫 번째 TaskManager를 선택합니다. 충분한 여유 리소스가 있는 TaskManager가 없으면, Native KubernetesYARN에 배포할 때 Flink는 새 TaskManager를 할당하려 시도합니다. 현재 전략에서 Flink는 사용자 구성에 따라 동일한 TaskManager를 할당합니다. TaskManager의 리소스 사양이 미리 정의되어 있으므로:

  • 클러스터에 리소스 단편화가 있을 수 있습니다. 예: 3 GB heap 메모리의 슬롯 요청이 두 개인데 TaskManager의 총 heap 메모리가 4 GB라면, Flink는 두 개의 TaskManager를 시작하고 각 TaskManager에서 1 GB heap 메모리가 낭비됩니다. 미래에는 작업의 슬롯 요청에 따라 이종 TaskManager를 할당할 수 있는 리소스 할당 전략이 있을 수 있어 리소스 단편화를 완화할 수 있습니다.
  • 슬롯 공유 그룹에 구성된 리소스 컴포넌트가 TaskManager의 총 리소스보다 크지 않도록 해야 합니다. 그렇지 않으면 작업이 예외와 함께 실패합니다.

더 알아보기 (Learn more)