이 페이지에서는 Apache Beam 버전 2.75.0을 기준으로 Apache Kafka에서 읽고 Apache Iceberg 테이블에 쓰는 Dataflow 스트리밍 작업의 성능 특성을 설명합니다. 직접 Apache Iceberg 쓰기와 관리형 BigQuery API를 통해 라우팅된 쓰기 간의 성능 차이를 평가하고 이러한 결과를 Kafka에서 BigQuery 파이프라인의 기준 벤치마크와 비교합니다. Apache Iceberg I/O 최적화가 진행 중이므로 이러한 성능 측정항목은 변경될 수 있습니다.
세 가지 기본 상태 비저장 매핑 구성에서 벤치마크 비교를 사용할 수 있습니다. 즉, 소스에서 읽고, 메시지를 레코드로 변환하고, 상태를 추적하거나 복잡한 비즈니스 로직을 적용하지 않고 싱크에 씁니다. 벤치마크에서 map_only 또는 mapping이라고 함.
- Kafka에서 BigQuery로 (
map_only) (기준 Kafka에서 BigQuery 성능) - Kafka에서 Iceberg로 직접 (
map_only,autosharding=false) - 관리형 BigQuery API를 사용하여 Kafka에서 Iceberg로(
map_only)
또한 이 가이드에서는 groupbykey를 사용하는 상태 저장 일괄 처리와 같은 직접 Apache Iceberg 스트리밍 패턴을 평가하고 파일 크기 분포, 자동 샤딩 동작, 읽기 측 쿼리 지연 시간과 관련된 중요한 다운스트림 고려사항을 자세히 설명합니다.
테스트 방법
벤치마크는 다음 리소스를 사용하여 수행되었습니다.
- Managed Service for Apache Kafka 클러스터: 트래픽은 Dataflow 스트리밍 데이터 생성기 템플릿을 사용하여 생성되었습니다.
- 입력 처리량: 1GBps
- 메시지 비율: 초당 약 1,000,000개의 메시지
- 메시지 형식: 고정 스키마가 있는 JSON 텍스트 (메시지당 약 1KB)
- 파티션: 1,000개의 Kafka 파티션
- 대상 싱크:
- BigQuery: BigQuery Storage Write API를 사용하여 작성된 표준 테이블 (파티션 없음)
- Apache Iceberg: Cloud Storage에서 지원하는 카탈로그 직접 싱크 는
bucket(id, 64)(기본 키에서 64개의 샤드로 버킷화됨)를 사용하여 파티션을 나누고hash배포 모드를 사용합니다.
수평 자동 확장이 안정화된 후 각 파이프라인 구성은 24시간 동안 정상 상태로 실행되었습니다. 각 파이프라인 사례의 벤치마크는 3번 별도로 실행되었으며 보고된 모든 값은 지속적이고 안정적인 성능 측정항목을 보장하기 위해 이러한 실행에서 계산된 평균을 나타냅니다.
수집 성능: 매핑 워크로드
상태 비저장 매핑 파이프라인은 소스에서 읽고, 메시지 형식을 레코드로 변환하고, 레코드 간에 상태를 추적하지 않고 싱크에 씁니다. 다음 섹션에서는 1GBps로 실행되는 참조 아키텍처를 분석합니다.
작업 구성
| 설정 | Kafka에서 BigQuery로 (map_only) |
Kafka에서 Iceberg로 직접 (autosharding=false) |
관리형 BigQuery API를 사용하여 Kafka에서 Iceberg로 |
|---|---|---|---|
| 작업자 머신 유형 | e2-standard-2 |
e2-standard-4 |
e2-standard-4 |
| 작업자당 vCPU | 2 | 4 | 4 |
| 작업자당 RAM | 8 GB | 16 GB | 16 GB |
| Streaming Engine | 사용 설정됨 | 사용 설정됨 | 사용 설정됨 |
| 수평 자동 확장 | 사용 설정됨 | 사용 설정됨 | 사용 설정됨 |
| 트리거 빈도 | 5초 | 60초 | 60초 |
처리량 및 리소스 사용량
객체 스토리지의 실제 Parquet 파일에 직접 쓰는 것은 BigQuery 스트리밍 수집보다 I/O 오버헤드가 더 큽니다. 직접 Iceberg 쓰기와 비교하여 관리형 BigQuery API를 통해 쓰기를 라우팅하면 작업자 CPU 사용률이 향상되고 (~70% 대 ~60%) Streaming Engine 소비가 약간 감소하지만 (~180 SECU/h 대 ~200 SECU/h) 전반적인 작업자 컴퓨팅 요구사항은 비슷하게 유지됩니다 (~440 vCPU 대 ~450 vCPU).
| 측정항목 | Kafka에서 BigQuery로 (map_only) |
Kafka에서 Iceberg로 직접 (autosharding=false) |
관리형 BigQuery API를 사용하여 Kafka에서 Iceberg로 |
|---|---|---|---|
| 작업자당 평균 입력 처리량 | 약 15MBps | 약 9MBps | 약 9MBps |
| 평균 CPU 사용률 | 약 70% | 약 60% | 약 70% |
| 1GBps 입력의 예상 vCPU | 약 126 vCPU | 약 450 vCPU | 약 440 vCPU |
| 1GBps 입력의 예상 작업자 | 약 63명의 작업자 | 약 110명의 작업자 | 약 110명의 작업자 |
| 1GBps의 시간당 예상 SECU | 약 58 SECU/h | 약 200 SECU/h | 약 180 SECU/h |
쓰기 지연 시간 프로필
직접 Iceberg 쓰기는 객체 스토리지 메타데이터 커밋 제약 조건으로 인해 심각한 꼬리 지연 시간 (P99)을 나타냅니다. 관리형 BigQuery API를 사용하면 중간 지연 시간을 낮게 유지하면서 꼬리 지연 시간 급증을 없앨 수 있습니다.
| 엔드 투 엔드 쓰기 지연 시간 | Kafka에서 BigQuery로 | Kafka에서 Iceberg로 직접 (autosharding=false) |
관리형 BigQuery API를 사용하여 Kafka에서 Iceberg로 |
|---|---|---|---|
| P50 (중간값) | 약 1,200ms | 약 1,000ms | 약 1,000ms |
| P95 | 약 3,000ms | 약 7,400ms | 약 1,900ms |
| P99 (꼬리) | 약 5,400ms | 약 14,000ms | 약 2,700ms |
자동 샤딩 고려사항 및 설계 선택
이 섹션에서는 Apache Iceberg에 쓸 때 자동 샤딩이 파일 크기 및 파이프라인 지연 시간에 미치는 영향을 설명합니다.
autosharding=false가 기준으로 선택된 이유
초기 테스트에서 자동 샤딩을 사용 설정하면 일정한 집계 입력 부하에서도 로컬 스레드 수준 부하 급증으로 인해 트리거되는 동적 샤드 분할로 인해 파일 크기가 작은 청크로 축소되고 임의로 변동되었습니다.
안정적이고 예측 가능한 Parquet 파일 레이아웃 (평균 약 800KB)을 유지하고 조기 플러시 없이 공정한 기준을 보장하기 위해 autosharding=false 가 직접 싱크 구성에 선택되었습니다.
자동 샤딩을 사용 중지하면 어떻게 되나요? 아니면 사용 설정된 상태로 유지하면 어떻게 되나요?
autosharding=false(기준): 자동 샤딩에 비해 초기 파일 크기가 더 큽니다 (평균 약 800KB). 이는 이상적인 Iceberg 파일 크기 (128~512MB)에 비해 여전히 작지만 다운스트림 압축이 훨씬 적게 필요합니다. 그러나 객체 스토리지 메타데이터 병목 현상으로 인해 쓰기 꼬리 지연 시간이 길다는 단점이 있습니다 (P99가 약 14.0초에 도달).- 자동 샤딩이 사용 설정된 경우: Dataflow는 쓰기 꼬리 지연 시간을 줄이는 로컬 처리량 급증을 흡수하도록 작성기 스레드를 동적으로 확장합니다. 그러나 작은 Parquet 파일 (약 100KB 이하)을 대량으로 생성하여 스토리지 레이어를 손상시킵니다. 이러한 파일 크기는 높은 분산을 나타내며 실행 간에 임의로 변동하여 (평균 약 39KB~약 100KB) 적극적인 다운스트림 압축 유지보수의 필요성을 높입니다.
파티션 조정 및 권장사항
평가 중에 대상 테이블의 다양한 고정 파티션 값을 실험하여 최적의 균형을 찾았습니다. 대상 테이블 파티셔닝에 64개의 버킷 (예: bucket(id, 64))을 사용하면 적절한 사용률과 처리량을 유지하면서 대상 파일 크기가 생성되는 것으로 확인되었습니다. 이 접근 방식을 사용하면 완전 동적 확장에 연결된 임의 파일 크기 조각화 문제를 방지하면서 자동 샤딩의 성능 이점을 일치시킬 수 있었습니다.
실무자를 위한 권장사항: 고객은 Parquet 파일 크기를 손상시키지 않고 파이프라인 병렬 처리를 극대화하는 최적의 지점을 찾기 위해 대상 파티션 설정으로 유사한 예비 테스트를 수행하는 것이 좋습니다.
다운스트림 읽기 영향: 파일 크기 및 압축
쓰기 측 측정항목은 Iceberg 수집을 위한 관리형 BigQuery API를 선호하지만 전반적인 파이프라인 효율성은 다운스트림 읽기 성능에 크게 좌우됩니다.
- 관리형 BigQuery API의 작은 파일 생성: 관리형 BigQuery API는 쓰기 지연 시간을 낮게 유지하기 위해 데이터를 자주 플러시합니다. 이 동작으로 인해 대상 Iceberg 카탈로그에 작은 Parquet 파일이 대량으로 작성됩니다.
- 읽기 쿼리 지연 시간 영향: 수백만 개의 작은 Parquet 파일이 있는 테이블을 읽는 쿼리 엔진 (예: Starburst/Trino, Apache Spark, BigQuery, Dremio)은 과도한 메타데이터 파싱 오버헤드와 파티션 스캔 페널티를 발생시킵니다.
- 압축 요구사항: 관리형 BigQuery API를 사용할 때 읽기 성능 저하를 방지하려면 (또는 직접 쓰기에서 자동 샤딩이 사용 설정된 경우) 정기적인 Iceberg 압축 유지보수 작업을 실행합니다 (예:
REWRITE DATA FILES). 압축의 컴퓨팅 오버헤드는 전반적인 아키텍처 설계에 반영되어야 합니다. - 직접 쓰기 (
autosharding=false) 파일 배포: 고정 샤딩을 사용하는 직접 Iceberg 쓰기는 더 큰 평균 Parquet 파일 (약 800KB)을 생성하여 즉각적인 압축 요구 없이 즉각적인 쿼리 액세스를 위한 조각화되지 않은 레이아웃을 생성합니다 (여전히 이상적인 범위 미만).
상태 저장 직접 Iceberg 파이프라인 (groupbykey)
수동 일괄 처리 전략을 평가하기 위해 상태 저장 키 그룹화 (groupbykey)가 기준 Kafka에서 Iceberg로 직접 (map_only, autosharding=false) 파이프라인에 대해 테스트되었습니다. 두 구성 모두 Parquet 파일을 객체 스토리지에 직접 씁니다.
벤치마크 비교
| 측정항목 / 기능 | 직접 싱크 기준 (autosharding=false) |
상태 저장 직접 싱크 (groupbykey) |
성능 영향 |
|---|---|---|---|
| 1GBps의 예상 vCPU | 약 450 vCPU | 약 520 vCPU | 약 +16% 컴퓨팅 필요 |
| 평균 CPU 사용률 | 약 60% | 약 50% | 약 -17% 작업자 효율성 |
| 1GBps의 예상 SECU/h | 약 200 SECU/h | 약 300 SECU/h | 약 +50% Streaming Engine 부하 |
| 평균 파일 크기 | 약 800KB | 약 100KB | 더 작은 파일 일괄 처리 생성 |
| P50 지연 시간 | 약 1,000ms | 약 1,200ms | 약 +20% 더 느린 중간값 |
| P95 지연 시간 | 약 7,400ms | 약 5,500ms | 약 -26% 더 낮은 지연 시간 |
| P99 지연 시간 | 약 14,000ms | 약 13,000ms | 미미한 꼬리 지연 시간 변경 |
균형 분석
- Streaming Engine 오버헤드: 상태 저장
groupbykey단계를 추가하려면 Beam이 창 경계 간에 중간 상태를 저장해야 합니다. 이렇게 하면 Streaming Engine 컴퓨팅 단위 소비가 약 50% 증가합니다 (약 200 SECU/h 에서 약 300 SECU/h 로). - 버퍼링 지연 시간: 수동 키 집계는 필수 창 버퍼링을 도입하여 중간 쓰기 지연 시간 (P50)을 약 1, 200ms 로,P95 지연 시간을 약 5.5초 로 늘립니다.
역방향 파이프라인: Iceberg에서 Kafka로 스트리밍
양방향 레이크하우스 기능을 평가하기 위해 Apache Iceberg 테이블에서 추가 전용 스트림을 읽고 Apache Kafka에 다시 게시하는 역방향으로 흐르는 스트리밍 데이터에 대한 벤치마크도 수행되었습니다.
작업 구성 및 효율성
과도한 객체 저장소 파일 쓰기 또는 메타데이터 커밋 병목 현상을 처리해야 하는 수집 파이프라인과 달리 Iceberg에서 변경사항을 읽고 스트리밍하는 것은 매우 효율적으로 작동합니다.
| 측정항목 | Iceberg에서 Kafka로 (추가 전용, 단 한 번) |
|---|---|
| 작업자 머신 유형 | e2-standard-4 |
| 1GBps 입력의 예상 vCPU | 약 30 vCPU |
| 1GBps 입력의 예상 작업자 | 약 7명의 작업자 |
| 1GBps의 시간당 예상 SECU | 약 0.2 SECU/h |
역방향 파이프라인의 주요 요점
- 훨씬 낮은 컴퓨팅 오버헤드: Iceberg에서 CDC 스트림을 읽고 프로젝션하는 데는 훨씬 적은 컴퓨팅 리소스 (직접 쓰기의 경우 약 30 vCPU 대 약 450 vCPU)가 필요합니다. 이는 대량의 Parquet 파일을 객체 스토리지에 파티셔닝, 인코딩, 커밋하는 과도한 작업을 방지하기 때문입니다.
- 리소스 효율성: 레이크하우스 형식에서 스트리밍 레이어로의 다운스트림 이벤트 기반 소비 또는 복제는 인바운드 수집 경로에 비해 매우 효율적입니다.
아키텍처 권장사항 요약
| 아키텍처 패턴 | P99 쓰기 지연 시간 | 파일 레이아웃 | 다운스트림 읽기 고려사항 |
|---|---|---|---|
Kafka에서 BigQuery로 (map_only) |
약 5.4초 | 해당 사항 없음 | 최적 (관리형 BigQuery 스토리지 엔진) |
| 관리형 BigQuery API를 사용하여 Kafka에서 Iceberg로 | 약 2.7초 | 임의로 작은 파일 | 대량 읽기의 경우 정기적인 압축이 필요 합니다. |
Kafka에서 Iceberg로 직접 (autosharding=false) |
약 14.0초 | 약 800KB | 양호 (더 큰 초기 파일 크기, 더 낮은 압축 요구) |
Kafka에서 Iceberg로 직접 (groupbykey) |
약 13.0초 | 약 100KB | 보통 (더 높은 컴퓨팅 및 상태 오버헤드) |
예상 비용
Google Cloud 가격 계산기를 사용하여 리소스 기반 결제로 자체의 유사한 파이프라인의 기준 비용을 다음과 같이 추정할 수 있습니다.
- 가격 계산기를 엽니다.
- 합산하여 추정 을 클릭합니다.
- Dataflow를 선택합니다.
- 서비스 유형으로 'Dataflow Classic'을 선택합니다.
- 고급 설정 을 선택하여 전체 옵션 세트를 표시합니다.
- 작업이 실행되는 위치를 선택합니다.
- 작업 유형으로 "스트리밍"을 선택합니다.
- Streaming Engine 사용 설정 을 선택합니다.
- 작업 실행 시간, 작업자 노드, 작업자 머신, Persistent Disk 스토리지에 대한 정보를 입력합니다.
- Streaming Engine 컴퓨팅 단위의 예상 개수를 입력합니다.
리소스 사용량과 비용은 입력 처리량에 따라 대략 선형적으로 확장되지만 작업자가 몇 명에 불과한 소규모 작업의 경우 총비용은 고정 비용이 지배합니다. 시작점으로 벤치마크 결과에서 작업자 노드 수와 리소스 소비를 추정할 수 있습니다.
예를 들어 입력 데이터 전송률이 100MBps인 Kafka에서 Iceberg로 직접 (autosharding=false) 아키텍처를 사용하여 파이프라인을 실행한다고 가정해 보겠습니다. 1GBps 파이프라인의 벤치마크 결과를 기반으로 리소스 요구사항을 다음과 같이 추정할 수 있습니다.
- 확장 요소: (100MBps) / (1024MBps) = 약 0.1
- 예상 작업자 노드: 작업자 110명 × 0.1 = 약 11명의 작업자
- 시간당 예상 Streaming Engine 컴퓨팅 단위 수: 200 × 0.1 = 시간당 약 20개 단위
이 값은 초기 추정치로만 사용해야 합니다. 실제 처리량과 비용은 머신 유형, 메시지 크기 분포, 사용자 코드, 집계 유형, 키 병렬 처리, 창 크기와 같은 요인에 따라 크게 달라질 수 있습니다. 자세한 내용은 Dataflow 비용 최적화 권장사항을 참고하세요.
테스트 파이프라인 실행
Dataflow Flex
템플릿을 사용하여 Apache Iceberg 스트리밍 작업을 배포하려면
gcloud dataflow flex-template run
명령어를 사용합니다.
gcloud dataflow flex-template run JOB_NAME \
--project=PROJECT_ID \
--region=REGION \
--template-file-gcs-location=gs://dataflow-templates-us-central1/latest/flex/Kafka_To_Iceberg_Yaml \
--enable-streaming-engine \
--parameters ^@^bootstrapServers="KAFKA_BOOTSTRAP_ADDRESS"\
@topic="KAFKA_TOPIC"\
@table="ICEBERG_TABLE_IDENTIFIER"\
@catalogName="CATALOG_NAME"\
@catalogProperties='{"type":"CATALOG_TYPE","warehouse":"gs://BUCKET_NAME/warehouse/"}'\
@triggeringFrequencySeconds=60\
@schema='SCHEMA_DEFINITION'
다음을 바꿉니다.
JOB_NAME: Dataflow 작업의 이름PROJECT_ID: Google Cloud 프로젝트 IDREGION: 작업이 실행되는 Google Cloud 리전 (예:us-central1)KAFKA_BOOTSTRAP_ADDRESS: Apache Kafka 클러스터의 부트스트랩 주소KAFKA_TOPIC: Kafka 주제의 이름ICEBERG_TABLE_IDENTIFIER: 대상 Iceberg 테이블의 식별자CATALOG_NAME: Iceberg 카탈로그의 이름CATALOG_TYPE: 사용할 카탈로그 유형 (예:hadoop또는bigquery)BUCKET_NAME: 창고 위치의 Cloud Storage 버킷 이름SCHEMA_DEFINITION: Kafka 주제 데이터의 스키마 정의 (예:{"type": "record", "name": "Record", "fields": [{"name": "id", "type": "string"}]})