배치 셔플

배치 셔플 (Batch Shuffle)

Flink는 유한(bounded) 입력에서 실행되는 작업을 위해 DataStream APITable / SQL 양쪽 모두에서 배치 실행 모드를 지원합니다. 배치 실행 모드에서 Flink는 네트워크 교환을 위해 Blocking ShuffleHybrid Shuffle 두 가지 모드를 제공합니다.

출처: 문서

본문

개요 (Overview)

Flink는 유한 입력에서 실행되는 작업을 위해 DataStream APITable / SQL에서 배치 실행 모드를 지원합니다. 배치 실행 모드에서 Flink는 네트워크 교환을 위해 두 가지 모드를 제공합니다.

  • Blocking Shuffle는 배치 실행의 기본 데이터 교환 모드입니다. 모든 중간 데이터를 영속화하며, 완전히 생성된 후에만 소비될 수 있습니다.
  • Hybrid Shuffle는 배치 실행의 차세대 데이터 교환 모드입니다. 데이터를 더 스마트하게 영속화하고, 생성되는 동안 소비를 허용합니다. 이 기능은 아직 실험적이며 몇 가지 알려진 제한 사항이 있습니다.

Blocking Shuffle

스트리밍 애플리케이션에 사용되는 파이프라인 셔플과 달리, blocking exchange는 데이터를 일부 저장소에 영속화합니다. 그리고 다운스트림 작업들이 네트워크를 통해 이 값들을 가져옵니다. 이러한 교환은 업스트림과 다운스트림 작업이 동시에 실행될 필요가 없으므로 작업 실행에 필요한 리소스를 줄여줍니다.

전반적으로 Flink는 Hash shuffleSort shuffle 두 가지 유형의 blocking shuffle을 제공합니다.

Hash Shuffle

1.14 이하의 기본 blocking shuffle 구현인 Hash Shuffle는 각 업스트림 작업이 결과를 TaskManager의 로컬 디스크에 각 다운스트림 작업별 별도 파일로 영속화합니다. 다운스트림 작업이 실행되면 업스트림 TaskManager에 파티션을 요청하며, 이들은 파일을 읽고 네트워크를 통해 데이터를 전송합니다.

Hash Shuffle는 파일을 쓰고 읽는 여러 메커니즘을 제공합니다.

  • file: 일반 File IO로 파일을 쓰고, Netty FileRegion으로 파일을 읽고 전송합니다. FileRegionsendfile 시스템 호출에 의존해 데이터 복사 횟수와 메모리 소비를 줄입니다.
  • mmap: mmap 시스템 호출로 파일을 쓰고 읽습니다.
  • auto: 일반 File IO로 파일을 씁니다. 파일 읽기의 경우 32비트 머신에서는 일반 file 옵션으로, 64비트 머신에서는 mmap을 사용합니다. 이는 32비트 머신에서 Java mmap 구현의 파일 크기 제한을 피하기 위함입니다.

다른 메커니즘은 TaskManager 구성으로 선택할 수 있습니다.

이 옵션은 실험적이며 향후 변경될 수 있습니다.

SSL이 활성화되면 file 메커니즘은 FileRegion을 사용할 수 없고 대신 전송 전에 데이터를 캐시하기 위해 un-pooled 버퍼를 사용합니다. 이는 직접 메모리 OOM을 유발할 수 있습니다. 또한 동기 파일 읽기가 잠시 Netty 스레드를 차단할 수 있으므로 SSL 핸드셰이크 타임아웃을 늘려 연결 재설정 오류를 피해야 합니다.

mmap의 메모리 사용은 구성된 메모리 한도에 포함되지 않지만, Yarn 같은 일부 리소스 프레임워크는 이 메모리 사용을 추적해 메모리가 일부 임계값을 초과하면 컨테이너를 종료할 수 있습니다.

Hash Shuffle는 SSD를 사용하는 소규모 작업에 잘 동작하지만 몇 가지 단점도 있습니다.

  • 작업 규모가 크면 너무 많은 파일을 만들 수 있고, 이 파일들을 동시에 쓰기 위한 큰 쓰기 버퍼가 필요합니다.
  • HDD에서 여러 다운스트림 작업이 동시에 데이터를 가져올 때 임의 IO(random IO) 문제가 발생할 수 있습니다.

Sort Shuffle

Sort Shuffle는 1.13에서 도입된 또 다른 blocking shuffle 구현이며 1.15에서 기본 blocking shuffle 구현이 되었습니다. Hash Shuffle과 달리 Sort Shuffle는 각 결과 파티션에 대해 하나의 파일만 씁니다. 여러 다운스트림 작업이 동시에 결과 파티션을 읽으면 데이터 파일은 한 번만 열리고 모든 리더가 공유합니다. 결과적으로 클러스터는 inode와 파일 디스크립터 같은 리소스를 더 적게 사용해 안정성이 향상됩니다. 또한 더 적은 파일을 쓰고 데이터를 순차적으로 읽도록 최선을 다하기 때문에 Sort Shuffle는 특히 HDD에서 Hash Shuffle보다 더 나은 성능을 얻을 수 있습니다. 게다가 Sort Shuffle은 추가 관리 메모리를 데이터 읽기 버퍼로 사용하고 sendfile이나 mmap 메커니즘에 의존하지 않으므로 SSL에서도 잘 동작합니다. Sort Shuffle에 대한 자세한 내용은 FLINK-19582FLINK-19614를 참조하세요.

sort blocking shuffle을 사용할 때 조정이 필요할 수 있는 몇 가지 구성 옵션은 다음과 같습니다.

현재 Sort Shuffle은 레코드 자체가 아니라 파티션 인덱스로만 레코드를 정렬합니다. 즉, sort는 데이터 클러스터링 알고리즘으로만 사용됩니다.

Blocking Shuffle 선택 (Choices of Blocking Shuffle)

요약하면,

  • SSD에서 실행되는 소규모 작업에는 두 구현 모두 동작해야 합니다.
  • 대규모 작업이나 HDD에서 실행되는 작업에는 Sort Shuffle이 더 적합해야 합니다.

Sort ShuffleHash Shuffle 사이를 전환하려면 taskmanager.network.sort-shuffle.min-parallelism 구성 옵션을 조정해야 합니다. 이는 다운스트림 작업의 병렬도에 따라 어떤 shuffle 구현을 사용할지 제어합니다. 병렬도가 구성된 값보다 낮으면 Hash Shuffle이, 그렇지 않으면 Sort Shuffle이 사용됩니다. 1.15보다 낮은 버전에서는 기본값이 Integer.MAX_VALUE이므로 기본적으로 Hash Shuffle이 사용됩니다. 1.15부터 기본값은 1이므로 기본적으로 Sort Shuffle이 사용됩니다.

Hybrid Shuffle

이 기능은 아직 실험적이며 몇 가지 알려진 제한 사항이 있습니다.

Hybrid shuffle은 차세대 배치 데이터 교환입니다. blocking shuffle과 (스트리밍 모드의) 파이프라인 shuffle의 장점을 결합합니다.

  • blocking shuffle처럼 업스트림과 다운스트림 작업이 동시에 실행될 필요가 없으므로 적은 리소스로 작업을 실행할 수 있습니다.
  • 파이프라인 shuffle처럼 다운스트림 작업이 업스트림 작업이 끝난 후에 실행될 필요가 없으므로 충분한 리소스가 주어졌을 때 작업의 전체 실행 시간을 줄입니다.
  • 서로 다른 spilling 전략을 제공함으로써 데이터를 덜 영속화하는 것과 실패 시 작업을 덜 다시 시작하는 것 사이의 사용자 선호에 적응합니다.

hybrid shuffle 모드를 사용하려면 execution.batch-shuffle-modeALL_EXCHANGES_HYBRID_FULL(full spilling 전략) 또는 ALL_EXCHANGES_HYBRID_SELECTIVE(selective spilling 전략)로 구성해야 합니다.

Spilling 전략

Hybrid shuffle은 두 가지 spilling 전략을 제공합니다.

  • Selective Spilling Strategy는 데이터가 다운스트림 작업에 적시에 소비되지 않을 때만 영속화합니다. 이는 영속화할 데이터 양을 줄이지만, 실패 시 업스트림 작업을 다시 시작해 완전한 중간 결과를 재생성해야 한다는 대가가 있습니다.
  • Full Spilling Strategy는 다운스트림 작업이 소비하는지 여부와 무관하게 모든 데이터를 영속화합니다. 실패 시 영속화된 완전한 중간 결과를 다시 소비할 수 있어 업스트림 작업을 다시 시작할 필요가 없습니다.

데이터 소비 제약 (Data Consumption Constraints)

Hybrid shuffle은 생산자와 소비자 사이의 파티션 데이터 소비 제약을 다음 세 가지 경우로 나눕니다.

  • ALL_PRODUCERS_FINISHED: 모든 생산자가 끝났을 때만 hybrid 파티션 데이터를 소비할 수 있습니다.
  • ONLY_FINISHED_PRODUCERS: hybrid 파티션은 끝난 생산자의 데이터만 소비할 수 있습니다.
  • UNFINISHED_PRODUCERS: hybrid 파티션은 끝나지 않은 생산자의 데이터도 소비할 수 있습니다.

이들은 jobmanager.partition.hybrid.partition-data-consume-constraint로 구성할 수 있습니다.

  • AdaptiveBatchScheduler의 경우: 기본 제약은 파이프라인형 shuffle을 수행하기 위한 UNFINISHED_PRODUCERS입니다. 값이 ALL_PRODUCERS_FINISHED 또는 ONLY_FINISHED_PRODUCERS로 설정되면 성능이 저하될 수 있습니다.
  • SpeculativeExecution이 활성화된 경우: 기본 제약은 blocking shuffle과 비교해 일부 성능 최적화를 가져오는 ONLY_FINISHED_PRODUCERS입니다. 생산자와 소비자가 동시에 실행될 기회가 있으므로 더 많은 추측 실행(speculative execution) 작업이 생성될 수 있고 장애 조치 비용도 증가합니다. blocking shuffle과 동일한 동작으로 되돌리려면 이 값을 ALL_PRODUCERS_FINISHED로 구성할 수 있습니다. 또한 이 모드에서는 UNFINISHED_PRODUCERS가 지원되지 않는다는 점에 유의하세요.

원격 저장소 지원 (Remote Storage Support)

Hybrid shuffle은 shuffle 데이터를 원격 저장소에 저장하는 것을 지원합니다. 원격 저장소 경로는 taskmanager.network.hybrid-shuffle.remote.path로 구성할 수 있습니다. 이 기능은 OSS, HDFS, S3 등을 포함한 다양한 원격 저장소 시스템을 지원합니다. Flink가 지원하는 파일 시스템에 대한 자세한 내용은 Flink Filesystem을 참조하세요.

제한 사항 (Limitations)

Hybrid shuffle 모드는 여전히 실험적이며 Flink 커뮤니티가 제거하기 위해 노력 중인 몇 가지 알려진 제한 사항이 있습니다.

  • Slot Sharing 미지원. hybrid shuffle 모드에서 Flink는 현재 각 작업이 전용 슬롯에서 독점적으로 실행되도록 강제합니다. slot sharing이 명시적으로 지정되면 오류가 발생합니다.
  • 동적 그래프에 대한 파이프라인 실행 없음. auto-parallelism(동적 그래프)이 활성화되면 Adaptive Batch Scheduler가 업스트림 작업이 끝날 때까지 기다렸다가 다운스트림 작업의 병렬도를 결정합니다. 이는 hybrid shuffle이 효과적으로 blocking shuffle(ALL_PRODUCERS_FINISHED 제약)로 폴백한다는 뜻입니다.

성능 튜닝 (Performance Tuning)

다음 지침은 특히 대규모 배치 작업에서 더 나은 성능을 얻는 데 도움이 될 수 있습니다.

Blocking Shuffle

  • HDD에서는 항상 Sort Shuffle을 사용하세요. Sort Shuffle은 안정성과 IO 성능을 크게 향상시킬 수 있습니다. 1.15부터 Sort Shuffle은 이미 기본 blocking shuffle 구현이며, 1.14 이하 버전에서는 taskmanager.network.sort-shuffle.min-parallelism을 1로 설정해 수동으로 활성화해야 합니다.
  • 두 blocking shuffle 구현 모두에서 데이터가 압축하기 어렵지 않다면 데이터 압축 활성화를 고려하세요. 1.15부터 데이터 압축은 기본적으로 활성화되어 있으며, 1.14 이하 버전에서는 수동으로 활성화해야 합니다.
  • Sort Shuffle을 사용할 때 채널당 전용 버퍼 수를 줄이고 게이트당 플로팅 버퍼 수를 늘리면 도움이 됩니다. 1.14 이상 버전에서는 taskmanager.network.memory.buffers-per-channel을 0으로, taskmanager.network.memory.floating-buffers-per-gate를 더 큰 값(예: 4096)으로 설정하는 것이 좋습니다. 이 설정의 두 가지 주요 장점: 1) 네트워크 메모리 소비를 병렬도에서 분리해 대규모 작업에서 "Insufficient number of network buffers" 오류 가능성을 낮춥니다; 2) 네트워크 버퍼를 필요에 따라 다른 채널에 분배해 네트워크 버퍼 활용도를 높이고 성능도 향상시킵니다.
  • 네트워크 메모리 총 크기를 늘리세요. 현재 기본 네트워크 메모리 크기는 상당히 보수적입니다. 대규모 작업에서는 더 나은 성능을 위해 총 네트워크 메모리 fraction을 최소 0.2로 늘리는 것이 좋습니다. 동시에 네트워크 메모리 크기의 하한상한도 조정해야 할 수 있습니다. 자세한 내용은 메모리 구성 문서를 참조하세요.
  • shuffle 데이터 쓰기 메모리 크기를 늘리세요. 위 섹션에서 언급했듯이, 대규모 작업에서는 메모리가 충분하다면 결과 파티션당 쓰기 버퍼 수를 최소 (2 * parallelism)으로 늘리는 것이 좋습니다. 이 구성 값을 늘린 후에는 "Insufficient number of network buffers" 오류를 피하기 위해 네트워크 메모리 총 크기도 늘려야 할 수 있음에 유의하세요.
  • shuffle 데이터 읽기 메모리 크기를 늘리세요. 위 섹션에서 언급했듯이, 대규모 작업에서는 공유 읽기 메모리 크기를 더 큰 값(예: 256M 또는 512M)으로 늘리는 것이 좋습니다. 이 메모리는 framework off-heap 메모리에서 잘라내므로 직접 메모리 OOM 오류를 피하려면 taskmanager.memory.framework.off-heap.size도 같은 크기만큼 늘려야 합니다.

Hybrid Shuffle

  • 네트워크 메모리 총 크기를 늘리세요. 현재 기본 네트워크 메모리 크기는 상당히 보수적입니다. 대규모 작업에서는 더 나은 성능을 위해 총 네트워크 메모리 fraction을 최소 0.2로 늘리는 것이 좋습니다. 동시에 네트워크 메모리 크기의 하한상한도 조정해야 할 수 있습니다. 자세한 내용은 메모리 구성 문서를 참조하세요.
  • shuffle 데이터 쓰기 메모리 크기를 늘리세요. 대규모 작업에서는 네트워크 메모리 총 크기를 늘리는 것이 좋습니다. shuffle 쓰기 단계에서 더 많은 메모리를 사용할수록 다운스트림이 메모리에서 직접 데이터를 읽을 기회가 많아집니다. 레거시 Hybrid shuffle 모드를 사용한다면 각 Result Partition에 최소 numSubpartition + 1개의 버퍼를 할당할 수 있도록 해야 하며, 그렇지 않으면 "Insufficient number of network buffers"가 발생할 수 있습니다.
  • shuffle 데이터 읽기 메모리 크기를 늘리세요. 대규모 작업에서는 공유 읽기 메모리 크기를 더 큰 값(예: 256M 또는 512M)으로 늘리는 것이 좋습니다. 이 메모리는 framework off-heap 메모리에서 잘라내므로 taskmanager.memory.framework.off-heap.size도 같은 크기만큼 늘려 직접 메모리 OOM 오류를 피해야 합니다.
  • 레거시 Hybrid shuffle 모드를 사용할 때 채널당 전용 버퍼 수를 줄이면 성능이 심각하게 저하됩니다. 따라서 이 값은 0으로 설정하면 안 되며, 대규모 작업에서는 적절히 늘릴 수 있습니다. 또한 hybrid shuffle의 경우 taskmanager.network.memory.read-buffer.required-per-gate.max가 기본적으로 Integer.MAX_VALUE로 설정되어 있습니다. 이 값을 조정하지 않는 것이 좋습니다. 그렇지 않으면 성능 저하의 위험이 있습니다.

문제 해결 (Trouble Shooting)

다음은 (드물게) 발생할 수 있는 예외와 도움이 될 수 있는 해결 방법입니다.

Blocking Shuffle

예외 (Exceptions) 잠재적 해결 방법 (Potential Solutions)
Insufficient number of network buffers 대상 작업을 실행하기에 네트워크 메모리 양이 충분하지 않다는 뜻이며 네트워크 메모리 총 크기를 늘려야 합니다. 1.15부터 Sort Shuffle이 기본 blocking shuffle 구현이 되었고, 어떤 경우에는 이전보다 더 많은 네트워크 메모리가 필요할 수 있습니다. 즉 1.15로 업그레이드한 후 배치 작업이 이 문제를 겪을 약간의 가능성이 있습니다. 이 경우 총 네트워크 메모리 크기만 늘리면 됩니다.
Too many open files 파일 디스크립터가 부족하다는 뜻입니다. Hash Shuffle을 사용 중이라면 Sort Shuffle으로 전환하세요. 이미 Sort Shuffle을 사용 중이라면 파일 디스크립터의 시스템 제한을 늘리는 것을 고려하고 사용자 코드가 너무 많은 파일 디스크립터를 소비하는지 확인하세요.
Connection reset by peer 보통 네트워크가 불안정하거나 과부하 상태라는 뜻입니다. 위에서 언급한 SSL 핸드셰이크 타임아웃 같은 다른 문제도 이 문제를 일으킬 수 있습니다. Hash Shuffle을 사용 중이라면 Sort Shuffle으로 전환하세요. 이미 Sort Shuffle을 사용 중이라면 네트워크 백로그를 늘리면 도움이 될 수 있습니다.
Network connection timeout 보통 네트워크가 불안정하거나 과부하 상태라는 뜻이며 네트워크 연결 타임아웃을 늘리거나 연결 재시도를 활성화하면 도움이 될 수 있습니다.
Socket read/write timeout 네트워크가 느리거나 과부하 상태임을 나타낼 수 있으며 네트워크 송수신 버퍼 크기를 늘리면 도움이 될 수 있습니다. 작업이 Kubernetes 환경에서 실행 중이라면 호스트 네트워크를 사용하는 것도 도움이 될 수 있습니다.
Read buffer request timeout Sort Shuffle을 사용할 때만 발생할 수 있으며 shuffle 읽기 메모리의 심한 경쟁을 의미합니다. 문제를 해결하려면 taskmanager.memory.framework.off-heap.batch-shuffle.size와 함께 taskmanager.memory.framework.off-heap.size를 늘릴 수 있습니다.
No space left on device 보통 디스크 공간이나 inode가 소진되었음을 의미합니다. 저장 공간을 늘리거나 정리를 하세요.
Out of memory error Hash Shuffle을 사용 중이라면 Sort Shuffle으로 전환하세요. 이미 Sort Shuffle을 사용하고 위 지침을 따른다면 해당 메모리 크기를 늘리는 것을 고려하세요. 힙 메모리는 taskmanager.memory.task.heap.size를, 직접 메모리는 taskmanager.memory.task.off-heap.size를 늘릴 수 있습니다.
Container killed by external resource manager 컨테이너가 종료되는 데는 여러 이유가 있습니다. 예를 들어 우선순위가 낮은 컨테이너를 종료해 우선순위가 높은 컨테이너를 위한 공간을 만들거나, 컨테이너가 메모리와 디스크 공간 같은 리소스를 너무 많이 사용하는 경우입니다. 위 섹션에서 언급했듯이 Hash Shuffle은 너무 많은 메모리를 사용해 YARN에게 종료될 수 있습니다. 따라서 Hash Shuffle을 사용 중이라면 Sort Shuffle으로 전환하세요. 이미 Sort Shuffle을 사용한다면 근본 원인을 찾아 해결하기 위해 Flink 로그와 외부 리소스 매니저 로그를 모두 확인해야 할 수 있습니다.

Hybrid Shuffle

예외 (Exceptions) 잠재적 해결 방법 (Potential Solutions)
Insufficient number of network buffers 대상 작업을 실행하기에 네트워크 메모리 양이 충분하지 않다는 뜻이며 네트워크 메모리 총 크기를 늘려야 합니다.
Connection reset by peer 보통 네트워크가 불안정하거나 과부하 상태라는 뜻입니다. SSL 핸드셰이크 타임아웃 같은 다른 문제도 이 문제를 일으킬 수 있습니다. 네트워크 백로그를 늘리면 도움이 될 수 있습니다.
Network connection timeout 보통 네트워크가 불안정하거나 과부하 상태라는 뜻이며 네트워크 연결 타임아웃을 늘리거나 연결 재시도를 활성화하면 도움이 될 수 있습니다.
Socket read/write timeout 네트워크가 느리거나 과부하 상태임을 나타낼 수 있으며 네트워크 송수신 버퍼 크기를 늘리면 도움이 될 수 있습니다. 작업이 Kubernetes 환경에서 실행 중이라면 호스트 네트워크를 사용하는 것도 도움이 될 수 있습니다.
Read buffer request timeout shuffle 읽기 메모리의 심한 경쟁을 의미합니다. taskmanager.memory.framework.off-heap.batch-shuffle.size와 함께 taskmanager.memory.framework.off-heap.size를 늘려 문제를 해결할 수 있습니다.
No space left on device 보통 디스크 공간이나 inode가 소진되었음을 의미합니다. 저장 공간을 늘리거나 정리를 하세요.

더 알아보기 (Learn more)