Dataflow는 관리형 컴퓨팅 인스턴스 풀에 작업을 분산하여 대규모 데이터 처리 파이프라인을 실행하도록 설계되었습니다. Dataflow가 처리를 병렬화하는 방식을 이해하면 효율적인 파이프라인을 설계하고, 성능 병목 현상을 방지하고, 리소스 비용을 최적화하는 데 도움이 됩니다.
이 페이지에서는 Dataflow가 데이터 처리를 병렬화하는 방법, 실행을 관리하고 확장하는 방법, 병렬 처리를 제한하는 일반적인 요인, 파이프라인 처리량을 최적화하는 데 사용할 수 있는 기법을 설명합니다.
병렬 처리 모델: 수평적 모델과 수직적 모델
Dataflow는 다음과 같은 두 가지 상호 보완적인 전략을 사용하여 병렬 처리를 달성합니다.
수평적 동시 처리: 파이프라인 데이터가 파티셔닝되고 여러 작업자 인스턴스 (가상 머신)에서 동시에 처리됩니다. Dataflow는 수평 자동 확장을 통해 워크로드 수요에 따라 작업자 풀 크기를 자동으로 조정할 수 있습니다. 기본적으로 Dataflow는 작업당 리소스 한도를 4,000개의 작업자로 설정하며, 이는 할당량 요청을 사용하여 조정할 수 있습니다.
수직적 동시 처리: 단일 작업자 인스턴스 내의 여러 CPU 코어와 스레드가 파이프라인 데이터를 동시에 처리합니다. 각 작업자 VM은 사용 가능한 컴퓨팅 리소스를 사용하기 위해 작업자 프로세스와 하네스 스레드를 실행합니다. 동적 스레드 확장을 사용하면 Dataflow가 CPU 사용률과 메모리 여유 공간을 기반으로 일괄 파이프라인의 작업자당 활성 스레드 수를 조정할 수 있습니다. Dataflow Prime에서 수직 자동 확장은 작업자에게 할당된 메모리와 컴퓨팅을 동적으로 확장합니다.
작업 단위 및 실행 계층 구조
작업자와 스레드에 처리를 분산하기 위해 Dataflow는 Apache Beam 파이프라인을 개별 작업 단위로 나눕니다.
- PCollection 및 파티션:
PCollection는 분산 데이터 세트를 나타냅니다. 바운드 데이터 (일괄 파이프라인)의 경우 Dataflow는 데이터 세트를 분할 또는 샤드로 나눕니다. 바운드가 없는 데이터 (스트리밍 파이프라인)의 경우 데이터가 지속적으로 도착하고 메시지 또는 스트림 파티션으로 수집됩니다. - 번들: Dataflow는
DoFn에서 처리할 수 있도록 요소를 임의 번들로 그룹화합니다. 번들은 실패 및 재시도의 단위입니다. 요소 처리 시 처리되지 않은 예외가 발생하면 전체 번들이 재시도됩니다. 메모리 소비가 많은 작업은 작업자 메모리 부족을 야기하고 메모리 부족 오류를 발생시킬 수 있습니다. - 단계 및 단계 융합: 그래프 최적화 중에 Dataflow는 인접한 변환을 융합된 실행 단계로 결합하여 중간 데이터 구체화의 오버헤드를 제거합니다. 융합된 단계 내에서 요소는 다음 단계 또는 셔플 경계로 전달되기 전에 단일 스레드에서 긴밀한 실행 루프로 처리됩니다.
파이프라인 변환 및 그래프 생성에 대한 자세한 내용은 파이프라인 수명 주기를 참고하세요.
관리형 동시 처리 및 자동 확장
기본적으로 Dataflow는 다음과 같은 방법으로 수동 파티션 조정 없이 파이프라인 병렬 처리를 자동으로 관리합니다.
- 수평 자동 확장:
- 일괄 파이프라인: 총 예상 남은 작업, 소스 백로그, CPU 사용량을 평가하여 작업을 빠르고 비용 효율적으로 완료할 수 있도록 작업자 풀을 확장하거나 축소합니다.
- 스트리밍 파이프라인: 시스템 지연 시간, 백로그 크기, CPU 사용률을 분석하여 처리량 급증 시 작업자를 확장하고 트래픽이 적은 기간에는 축소합니다. 자세한 내용은 스트리밍 수평 자동 확장 조정을 참고하세요.
- 동적 작업 재분배 (DWR): 일괄 파이프라인에서 Dataflow는 개별 작업자 작업의 진행 상황을 모니터링합니다. 작업자가 일찍 완료되거나 데이터 기울기(지연자)로 인해 다른 작업자가 뒤처지는 경우 Dataflow는 느린 작업자에서 처리되지 않은 잔여 작업을 동적으로 분할하여 유휴 작업자에게 다시 할당합니다. 자세한 내용은 동적 작업 재분배를 참고하세요.
- 동적 스레드 확장: 휴대용 러너를 사용하는 일괄 파이프라인에서 CPU 사용률 및 메모리 여유 공간에 따라 작업자당 동시 처리 스레드 수를 자동으로 조정합니다. 자세한 내용은 동적 스레드 확장을 참고하세요.
- 수직 자동 확장: Dataflow Prime에서 Dataflow는 메모리 부족 오류를 방지하고 리소스 사용률을 최적화하기 위해 작업자 메모리와 컴퓨팅 리소스를 동적으로 확장합니다. 자세한 내용은 수직 자동 확장을 참고하세요.
병렬 처리를 제한하는 요인
다음 데이터 특성 또는 파이프라인 그래프 설계로 인해 파이프라인이 예상되는 병렬 처리를 달성하지 못할 수 있습니다.
분할할 수 없는 입력 소스
입력 소스를 독립적인 범위로 분할할 수 없는 경우 Dataflow는 단일 작업자 스레드로 소스를 순차적으로 읽어야 합니다.
- 분할할 수 없는 파일 압축:
.gz(gzip) 또는.bzip2(색인 미포함)과 같은 형식은 임의 바이트 오프셋에서 병렬로 읽을 수 없습니다. 단일 대형 압축 파일을 읽으면 데이터가 압축 해제되고 재분배될 때까지 수집 단계가 단일 스레드로 제한됩니다. - 해결 방법: 분할 가능한 파일 형식 (예: Parquet, Avro 또는 Snappy 압축 형식)으로 데이터를 저장하거나 Cloud Storage에서 입력 데이터를 여러 개의 작은 파일로 분할합니다.
단계 병합 및 높은 팬아웃
단계 융합은 직렬화 오버헤드를 줄여 성능을 개선하지만 병렬 처리가 낮은 단계에서 많은 수의 출력 요소('높은 팬아웃' 작업)를 생성하는 경우 실수로 병렬 처리를 제한하고 메모리 압력을 높일 수 있습니다.
- 예: 소스가 5개의 파일을 읽고 1,000,000개의 출력 요소를 생성하는
FlatMap변환과 융합됩니다.FlatMap변환이 다운스트림 변환과 융합되면 1,000,000개의 요소가 최대 5개의 작업자 스레드에서 계속 실행되므로 다운스트림 처리량이 심각하게 제한됩니다. 또한 커밋하기 전에 중간 변환이 메모리에서 크게 확장되면 큰 번들로 인해 사용 가능한 작업자 메모리가 소진될 수 있습니다. - 해결 방법: 융합을 중단하고 작업자 풀 전체에 작업을 재분배하려면 팬아웃이 많은 단계와 다운스트림 변환 사이에
Redistribute(또는 기존Reshuffle) 변환을 삽입하세요. 메모리 관련 문제를 디버깅하려면 메모리 부족 오류 문제 해결을 참고하세요.
키 기울기 및 단축키
집계 작업 (GroupByKey, CoGroupByKey, Combine.PerKey)은 연결된 키를 기준으로 요소를 그룹화합니다.
- 핫 키 병목 현상: Dataflow는 집계를 위해 동일한 키를 가진 모든 요소를 단일 작업자 스레드로 라우팅합니다. 단일 키에 전체 데이터 세트의 큰 비율이 포함된 경우 해당 작업자는 지연자가 되고 업스트림 작업자에게 역압이 발생할 수 있습니다. 예를 들어 기본
null키 또는 매우 인기 있는 카테고리 키입니다. - 해결 방법:
- 가능한 경우
GroupByKey대신 Combiner(CombineFn또는Combine.PerKey)를 사용하여 Dataflow가 셔플 전에 부분 로컬 조합을 실행하도록 합니다. - 작업자 간에 키 공간을 분산하기 위해 단축키에 무작위 정수 접두사 또는 접미사를 추가하고 (키 솔팅) 2단계 집계를 통해 솔팅된 결과를 병합합니다.
- 가능한 경우
다운스트림 싱크 제한
데이터베이스나 서드 파티 API와 같은 외부 서비스에 파이프라인 출력을 쓸 때 병렬 처리가 높으면 대상 시스템이 포화될 수 있습니다.
- 제한: 동시 쓰기 호출을 실행하는 수백 개의 작업자 스레드는 비율 제한 오류, 연결 시간 초과 또는 데이터베이스 성능 저하로 이어질 수 있습니다.
- 해결 방법:
GroupByKey로 요소를 그룹화하거나 제어된 병렬 처리를 사용하여 일괄 처리 싱크를 사용하여 쓰기 병렬 처리를 제한합니다.- 싱크
DoFn구현에서 클라이언트 측 지수 백오프 및 재시도 로직을 구현합니다.
최적화 전략
Dataflow 작업의 병렬 처리를 최적화하려면 다음 방법을 고려하세요.
Redistribute로 원치 않는 융합 방지:Redistribute.arbitrarily(): 단계 병합을 중단하고 사용 가능한 모든 작업자에게 요소를 균등하게 재분배합니다.Redistribute.byKey(): 키 지역성을 유지하면서 작업자 스레드 간에 키-값 쌍의 균형을 재조정합니다.- 구현 예는 융합 방지를 참고하세요.
- 지연 및 병목 현상 모니터링: Google Cloud 콘솔 실행 세부정보를 사용하여 지연 수가 많거나 진행이 중단된 단계를 식별합니다.
다음 단계
- 파이프라인 수명 주기에 대해 알아봅니다.
- 수평 자동 확장을 살펴봅니다.
- 동적 작업 재분배를 이해합니다.
- Dataflow 파이프라인 권장사항을 검토합니다.
- 메모리 부족 오류 문제 해결 방법을 알아보세요.