적응형 배치 실행

적응형 배치 실행 (Adaptive Batch Execution)

이 문서는 적응형 배치 실행(adaptive batch execution)의 배경, 사용법, 제한 사항을 설명합니다.

출처: 문서

본문

배경 (Background)

전통적인 Flink 배치 작업 실행 과정에서는 작업의 실행 계획이 제출 전에 결정됩니다. 실행 계획을 최적화하려면 사용자와 Flink 의 정적 실행 계획 최적화 도구가 작업 로직을 이해하고 각 노드가 처리하는 데이터 특성과 연결 에지의 데이터 분포를 포함해 작업이 어떻게 실행될지 정확히 평가해야 합니다.

그러나 실제 시나리오에서 이러한 데이터 특성은 작업이 실행되기 전에 예측할 수 없습니다. 입력 데이터에 대한 풍부한 통계 정보가 있다면 사용자와 Flink 의 정적 실행 계획 최적화 도구는 이 통계를 실행 계획의 각 operator 특성과 결합해 어느 정도의 추론적 최적화를 수행할 수 있습니다. 하지만 실제 프로덕션 환경에서는 입력 데이터에 대한 통계 정보가 불완전하거나 부정확한 경우가 많아 Flink 작업의 중간 데이터를 추정하기 어렵습니다.

이 문제를 해결하기 위해 Flink 는 AdaptiveBatchScheduler, 즉 실행 계획을 자동으로 조정할 수 있는 배치 작업 스케줄러를 도입했습니다. 이 스케줄러는 작업이 실행되는 동안 작업 실행 계획을 점진적으로 결정하고 결정된 실행 계획에 따라 JobGraph 를 증분 생성합니다. 아직 결정되지 않은 실행 계획 부분에 대해서는 Flink 가 특정 최적화 전략과 중간 데이터의 특성에 따라 런타임에 실행 계획을 동적으로 조정할 수 있습니다. 현재 스케줄러가 지원하는 최적화 전략은 다음과 같습니다:

operator 의 병렬도 자동 결정

Adaptive Batch Scheduler 는 배치 작업의 operator 병렬도를 자동으로 결정하는 것을 지원합니다. operator 에 병렬도가 설정되어 있지 않으면 스케줄러는 operator 가 소비하는 데이터셋의 크기에 따라 병렬도를 결정합니다. 이는 많은 이점을 제공합니다:

  • 배치 작업 사용자는 병렬도 튜닝에서 해방될 수 있습니다.
  • 자동으로 조정된 병렬도는 매일 크기가 달라지는 소비 데이터셋에 더 잘 맞을 수 있습니다.
  • SQL 배치 작업의 operator 에는 자동으로 조정된 서로 다른 병렬도를 할당할 수 있습니다.

사용법

Adaptive Batch Scheduler 로 operator 의 병렬도를 자동으로 결정하려면 다음이 필요합니다:

1. 기능 켜기: Adaptive Batch Scheduler 는 기본적으로 자동 병렬도 파생을 활성화합니다. execution.batch.adaptive.auto-parallelism.enabled 를 구성해 이 기능을 전환할 수 있습니다. 또한 Adaptive Batch Scheduler 로 operator 병렬도를 자동 결정할 때 조정이 필요할 수 있는 관련 구성 옵션이 몇 가지 있습니다:

2. operator 의 병렬도 설정 피하기: Adaptive Batch Scheduler 는 병렬도가 설정되지 않은 operator 의 병렬도만 결정합니다. 따라서 operator 의 병렬도가 자동으로 결정되게 하려면 'setParallelism()' 메서드를 통해 operator 의 병렬도를 설정하지 않아야 합니다.

Source 에 대한 동적 병렬도 추론 지원 활성화

새로운 SourceDynamicParallelismInference 인터페이스를 구현해 동적 병렬도 추론을 활성화할 수 있습니다:

public interface DynamicParallelismInference {
    int inferParallelism(Context context);
}

Context 는 추론된 병렬도의 상한, 각 task 가 처리할 것으로 예상되는 평균 데이터 크기, 병렬도 추론에 도움이 되는 동적 필터링 정보를 제공합니다.

Adaptive Batch Scheduler 는 소스 vertex 를 스케줄링하기 전에 이 인터페이스를 호출하며, 구현은 시간이 오래 걸리는 연산을 가능한 한 피해야 함에 유의하세요.

Source 가 이 인터페이스를 구현하지 않으면 구성 설정 execution.batch.adaptive.auto-parallelism.default-source-parallelism 이 소스 vertex 의 병렬도로 사용됩니다.

동적 소스 병렬도 추론은 병렬도가 이미 지정되지 않은 소스 operator 의 병렬도만 결정한다는 점에 유의하세요.

성능 튜닝
  • Sort Shuffle 을 사용하고 taskmanager.network.memory.buffers-per-channel0 으로 설정하는 것이 권장됩니다. 이는 필요한 네트워크 메모리를 병렬도와 분리해, 대규모 작업에서 "Insufficient number of network buffers" 오류가 발생할 가능성을 줄여줍니다.
  • execution.batch.adaptive.auto-parallelism.max-parallelism 을 최악의 경우에 필요할 것으로 기대하는 병렬도로 설정하는 것이 권장됩니다. 그보다 큰 값은 성능에 영향을 줄 수 있으므로 권장되지 않습니다. 이 옵션은 업스트림 task 가 생성하는 subpartition 수에 영향을 줄 수 있으며, subpartition 수가 많으면 hash shuffle 의 성능과 작은 패킷으로 인한 네트워크 전송 성능이 저하될 수 있습니다.

데이터 분포의 자동 로드 밸런싱

Adaptive Batch Scheduler 는 데이터 분포의 자동 로드 밸런싱을 지원합니다. 스케줄러는 데이터를 다운스트림 subtask 에 균등하게 분배하려 시도해 각 다운스트림 subtask 가 소비하는 데이터 양이 대략 동일하도록 보장합니다. 이 최적화는 사용자의 수동 구성이 필요 없으며 point-wise 연결(예: Rescale) 과 all-to-all 연결(예: Hash, Rebalance, Custom) 을 포함한 다양한 연결 에지에 적용됩니다.

제한 사항

  • 현재 데이터 분포의 자동 로드 밸런싱은 operator 의 병렬도 자동 결정이 있는 operator 만 지원합니다. 따라서 사용자는 operator 의 병렬도 자동 결정을 활성화하고 operator 의 병렬도를 수동으로 설정하는 것을 피해야 자동 데이터 분포 밸런싱 최적화의 혜택을 받을 수 있습니다.
  • 이 최적화는 단일 키 핫스팟(hotspot)이 있는 시나리오를 완전히 해결하지는 못합니다. 단일 키의 데이터가 다른 키의 데이터를 훨씬 초과하면 여전히 핫스팟이 발생할 수 있습니다. 하지만 데이터의 정확성을 위해 이 키의 데이터를 분할해 서로 다른 subtask 에 할당할 수는 없습니다. 그러나 특정 상황에서는 적응형 Skewed Join 최적화 에서 볼 수 있듯 단일 키 문제를 해결할 수 있습니다.

적응형 Broadcast Join

분산 데이터 처리에서 broadcast join 은 일반적인 최적화입니다. Broadcast join 은 두 테이블 중 하나가 매우 작아 단일 컴퓨팅 노드의 메모리에 들어갈 수 있을 만큼 작다면, 그 전체 데이터셋을 각 분산 컴퓨팅 노드에 브로드캐스트할 수 있다는 원리로 작동합니다. 이 join 연산은 메모리에서 직접 수행할 수 있습니다. Broadcast join 최적화는 큰 테이블의 셔플링과 정렬을 크게 줄일 수 있습니다. 그러나 정적 최적화 도구는 이 최적화가 언제 효과적인지 결정하는 조건을 종종 잘못 판단하며, 이는 다음과 같은 이유로 프로덕션에서의 적용을 제한합니다:

  • 소스 테이블에 대한 통계 정보의 완전성과 정확성 이 프로덕션 환경에서 종종 부족합니다.
  • 소스 테이블에서 오지 않은 입력 에 대한 정확한 판단은 어렵습니다. 작업이 실행되는 동안까지 중간 데이터의 크기를 정확히 결정할 수 없기 때문입니다. join 연산이 소스 테이블 데이터에서 멀리 떨어져 있을 때, 이러한 조기 평가는 보통 불가능합니다.
  • 정적 최적화가 효과적인 조건을 잘못 평가하면 심각한 문제를 일으킬 수 있습니다. 예를 들어 큰 테이블이 실수로 작다고 분류되어 단일 노드 메모리에 들어가지 못하면, broadcast join operator 는 메모리에 해시 테이블을 만들려다 Out of Memory 오류로 실패할 수 있으며, 그 결과 task 를 재시작 해야 할 수 있습니다.

따라서 broadcast join 은 올바르게 사용하면 상당한 성능 향상을 제공할 수 있지만 실제 적용은 제한적입니다. 반면 적응형 Broadcast Join 은 Flink 가 실제 데이터 입력에 따라 런타임에 Join operator 를 Broadcast Join 으로 동적으로 변환할 수 있게 해줍니다.

Join 의미론의 정확성을 보장하기 위해 적응형 Broadcast Join 은 join 유형에 따라 입력 에지가 브로드캐스트될 수 있는지 결정합니다. 브로드캐스팅이 가능한 조건은 다음과 같습니다:

Join 유형 왼쪽 입력 오른쪽 입력
Inner
LeftOuter
RightOuter
FullOuter
Semi
Anti

사용법

Adaptive Batch Scheduler 는 기본적으로 컴파일 타임 정적 Broadcast Join 과 런타임 동적 적응형 Broadcast Join 을 동시에 활성화 합니다. table.optimizer.adaptive-broadcast-join.strategy 를 구성해 Broadcast Join 의 시점을 제어할 수 있습니다. 예를 들어 값을 RUNTIME_ONLY 로 설정하면 적응형 Broadcast Join 이 런타임에만 적용됩니다. 또한 작업의 특정 요구사항에 따라 다음 구성을 조정할 수 있습니다:

  • table.optimizer.join.broadcast-threshold: 브로드캐스트할 수 있는 데이터 양의 임계값. TaskManager 의 메모리가 크면 이 값을 적절히 늘릴 수 있고, 메모리가 제한적이면 줄일 수 있습니다.

제한 사항

  • 적응형 Broadcast Join 은 MultiInput operator 내에 포함된 Join operator 의 최적화를 지원하지 않습니다.
  • 적응형 Broadcast Join 은 Batch Job Recovery Progress 와 동시에 활성화되는 것을 지원하지 않습니다. 따라서 Batch Job Recovery Progress 가 활성화된 후에는 적응형 Broadcast Join 이 적용되지 않습니다.

적응형 Skewed Join 최적화

Join 쿼리에서 특정 키가 자주 나타나면 각 Join task 가 처리하는 데이터 양에 상당한 변동이 생길 수 있습니다. 이는 개별 처리 task 의 성능을 크게 저하시켜 작업의 전반적인 성능을 떨어뜨릴 수 있습니다. 그러나 Join operator 의 두 입력 측은 연관되어 있으므로 같은 keyGroup 은 같은 다운스트림 sub-task 에서 처리되어야 합니다. 따라서 Join 연산의 데이터 왜곡 문제를 해결하기 위해 자동 로드 밸런싱에만 의존하는 것으로는 충분하지 않습니다. 적응형 Skewed Join 최적화는 Join operator 가 런타임 입력 통계에 따라 왜곡되고 분할 가능한 데이터 파티션을 동적으로 분할 할 수 있게 해줘, 왜곡된 데이터로 인한 꼬리 지연(tail latency) 문제를 완화합니다.

Join 의미론의 정확성을 보장하기 위해 적응형 Skewed Join 최적화는 Join 유형에 따라 입력 에지를 동적으로 분할할 수 있는지 결정합니다. 분할이 가능한 시나리오는 다음과 같습니다:

Join 유형 왼쪽 입력 오른쪽 입력
Inner
LeftOuter
RightOuter
FullOuter
Semi
Anti

사용법

Adaptive Batch Scheduler 는 기본적으로 Skewed Join 최적화를 활성화합니다. table.optimizer.skewed-join-optimization.strategy 를 구성해 Skewed Join 의 최적화 전략을 제어할 수 있습니다. 이 구성의 가능한 값:

  • none: Skewed Join 최적화를 비활성화합니다.
  • auto: Skewed Join 최적화를 허용합니다. 그러나 일부 경우 Skewed Join 최적화는 데이터 정확성을 보장하기 위해 추가 Shuffle 이 필요할 수 있습니다. 이 경우 추가 Shuffle 로 인한 오버헤드를 피하기 위해 Skewed Join 최적화는 적용되지 않습니다.
  • forced: Skewed Join 최적화를 허용하며, 추가 Shuffle 이 도입되는 시나리오에서도 적용됩니다.

또한 작업의 특성에 따라 다음 구성을 조정할 수 있습니다:

  • table.optimizer.skewed-join-optimization.skewed-threshold: Skewed Join 최적화를 트리거하는 최소 데이터 양. Join 단계에서 subtask 가 처리하는 최대 데이터 양이 이 skew 임계값을 초과하면 Flink 는 subtask 가 처리하는 최대 데이터 양과 중앙값 데이터 양의 비율을 skew factor 미만으로 자동 낮출 수 있습니다.
  • table.optimizer.skewed-join-optimization.skewed-factor: Skew factor. Join 단계에서 Flink 는 subtask 가 처리하는 최대 데이터 양과 중앙값 데이터 양의 비율을 이 skew factor 미만으로 자동 낮춰 더 균형 잡힌 데이터 분포를 달성합니다.

제한 사항

  • 적응형 Skewed Join 최적화는 Join operator 의 병렬도에 영향을 줄 수 있으므로, 적용되려면 operator 의 병렬도 자동 결정이 필요합니다.
  • 적응형 Skewed Join 최적화는 MultiInput operator 내에 포함된 Join operator 의 최적화를 지원하지 않습니다.
  • 적응형 Skewed Join 최적화는 Batch Job Recovery Progress 와 동시에 활성화되는 것을 지원하지 않습니다. 따라서 Batch Job Recovery Progress 가 활성화된 후에는 적응형 Skewed Join 최적화가 적용되지 않습니다.

적응형 배치 실행의 제한 사항

  • AdaptiveBatchScheduler 전용: AdaptiveBatchScheduler 를 사용할 때만 적용됩니다. Adaptive Batch Scheduler 는 Flink 의 기본 배치 작업 스케줄러이므로, 사용자가 명시적으로 다른 스케줄러를 구성하지 않는 한(예: [jobmanager.scheduler: default]) 추가 구성이 필요하지 않습니다.
  • BLOCKING 또는 HYBRID 작업 전용: 현재 Adaptive Batch Scheduler 는 shuffle modeALL_EXCHANGES_BLOCKING / ALL_EXCHANGES_HYBRID_FULL / ALL_EXCHANGES_HYBRID_SELECTIVE 인 작업만 지원합니다.
  • FileInputFormat source 는 지원되지 않음: StreamExecutionEnvironment#readFile(...)StreamExecutionEnvironment#createInput(FileInputFormat, ...) 를 포함한 FileInputFormat source 는 지원되지 않습니다. Adaptive Batch Scheduler 를 사용할 때는 새 source(FileSystem DataStream 커넥터 또는 FileSystem SQL 커넥터)를 사용해 파일을 읽어야 합니다.
  • WebUI 의 일관되지 않은 broadcast 결과 메트릭: Adaptive Batch Scheduler 로 operator 병렬도를 자동 결정할 때, broadcast 결과의 경우 메트릭으로 집계된 업스트림 task 가 보낸 바이트/레코드 수가 다운스트림 task 가 받은 바이트/레코드 수와 같지 않아 Web UI 에 표시될 때 사용자를 혼동시킬 수 있습니다. 자세한 내용은 FLIP-187 을 참고하세요.

더 알아보기 (Learn more)