Google PDE: Dataflow 및 Apache Beam을 사용한 스트림 처리 — 학습 가이드
다음의 일부입니다: Google Professional Data Engineer — 학습 가이드. 검증된 답안으로 연습하기: Google 시험 허브, 또는 다음에서 시간 제한 모의고사 풀기: ExamRoll.io.
개요
Google Cloud의 스트림 처리는 Dataflow runner에서 실행되는 Apache Beam의 통합 프로그래밍 모델을 중심으로 합니다. Beam은 PCollection에 대한 변환(transform) 파이프라인이라는 논리적 추상화를 제공하여, 병렬 처리, 자동 확장, 내결함성과 같은 실행 세부 정보로부터 코드를 분리합니다. 스트리밍에서 정확성은 시간 시맨틱(이벤트 시간 vs 처리 시간), 윈도잉(고정, 슬라이딩, 세션, 전역), 워터마크, 트리거, 지연 데이터 처리에 따라 결정됩니다. Dataflow에서 운영 우수성을 달성하려면 올바른 워커 크기 조정, 자동 확장 정책, 스트리밍 엔진, 셔플 선택, 멱등성 있는(idempotent) 싱크 설계, 데드-레터(dead-letter) 처리, 강력한 관측 가능성(observability)이 필요합니다.
Apache Beam 모델과 시간 시맨틱
파이프라인, 변환(transform), PCollection, runner:
- Beam 파이프라인은 PCollection(유한 또는 무한)에 PTransform으로 구성된 방향성 비순환 그래프(DAG)를 적용합니다.
- Runner(Dataflow, Spark, Flink, Direct)가 파이프라인을 실행합니다. Dataflow는 관리형 자동 확장, 체크포인팅, 운영 가시성을 제공합니다.
- 변환에는 요소별(ParDo), 그룹화 및 결합(GroupByKey, Combine), 조인(CoGroupByKey), IO(PubSubIO, BigQueryIO, FileIO)가 포함됩니다.
윈도우:
- 고정 윈도우(Fixed windows): 주기적인 집계를 위한 겹치지 않는 시간 조각(예: 1분 텀블링 윈도우).
- 슬라이딩 윈도우(Sliding windows): 매끄러운 롤링 메트릭을 위한 겹치는 윈도우(예: 1분마다 슬라이딩하는 5분 윈도우).
- 세션 윈도우(Session windows): 비활성 간격 후에 닫히는 동적 윈도우로, 사용자 세션이나 디바이스 버스트(burst)에 이상적입니다.
- 전역 윈도우(Global window): 전체 무한 스트림에 대한 기본 윈도우 없는 뷰로, 주기적인 구체화(materialization)를 위해 트리거와 함께 사용되는 경우가 많습니다.
이벤트 시간 vs 처리 시간:
- 이벤트 시간: 소스에서 이벤트가 발생한 시간. 가변적인 전송 지연 시간에도 불구하고 논리적으로 일관된 집계를 가능하게 합니다.
- 처리 시간: 파이프라인이 이벤트를 관찰한 시간. 운영 트리거에는 유용하지만 시맨틱 정확성에는 적합하지 않습니다.
워터마크:
- 워터마크는 이벤트 시간의 완전성을 추정합니다(runner가 T 시간까지의 모든 이벤트를 보았다는 추측).
- 워터마크는 역압(backpressure)이나 소스 지연으로 인해 불규칙하게 진행되거나 정체될 수 있습니다. 지연 데이터는 타임스탬프 < 워터마크인 상태로 도착하는 모든 데이터입니다.
트리거와 지연:
- 기본값: 워터마크가 윈도우 끝을 지나면 실행되는 AfterWatermark 트리거. 허용된 지연 시간(allowed lateness) = 0이면 지연 데이터는 삭제됩니다.
- 조기 실행(처리 시간 또는 카운트 기반)은 지연 시간이 짧은 예비 결과를 제공합니다.
- 지연 실행은 지연 데이터가 도착했을 때 수정을 허용합니다. 누적 모드(accumulation mode)는 창(pane)이 결과를 누적할지 이전 출력을 폐기할지를 결정합니다.
- 비즈니스 허용 범위와 스토리지/컴퓨팅 트레이드오프에 따라 허용된 지연 시간을 선택하세요. 지연 시간이 길어질수록 상태 유지 기간과 비용이 증가합니다.
상태 저장 처리, 타이머, 세션화, 중복 제거:
- 상태 저장 DoFn은 키별 상태(예: 마지막으로 본 이벤트, 실행 중인 집계)를 유지하고 타이머를 설정하여 상태를 내보내거나 지웁니다.
- 세션화는 SessionWindows를 통해 자연스럽게 표현됩니다. 커스텀 로직의 경우 키 지정 상태(keyed state)와 처리/이벤트 시간 타이머를 사용합니다.
- 중복 제거: 이벤트별로 안정적인 ID를 사용하고, 윈도우별로 Distinct/Combine을 사용하거나 키별 상태(예: Bloom filter 또는 TTL이 있는 집합)를 사용합니다. 메모리 및 거짓 양성(false positive)과 엄격한 정확성 사이의 트레이드오프가 있습니다.
실패 모드와 트레이드오프:
- 비즈니스 메트릭에 처리 시간 윈도우를 사용하면 스파이크나 재시도 시 편차(drift)가 발생하므로 이벤트 시간 윈도우를 사용하는 것이 좋습니다.
- 너무 작은 윈도우와 빈번한 조기 트리거는 과도한 창(pane) 출력과 싱크 쓰기 증폭을 유발합니다.
- 허용된 지연 시간을 무제한으로 설정하면 상태가 비대해질 수 있습니다. 항상 상태 TTL을 제한하고 타이머를 설정하여 휴면 키를 지워야 합니다.
스트리밍 워크로드를 위한 Dataflow 운영
작업자(worker) 크기 조정 및 자동 확장:
- 수평 자동 확장은 백로그, 워터마크 지연, CPU, 처리량을 기반으로 작업자를 추가/제거합니다. 급증하는 트래픽을 흡수할 수 있도록 합리적인 maxWorkers를 설정하세요.
- 병목 현상에 따라 머신 유형을 선택하세요: CPU 바운드(더 많은 vCPU), 메모리 바운드(고용량 메모리 유형), 네트워크 바운드(더 큰 VM은 셔플 오버헤드를 줄임).
- 셔플이 많거나 파일 기반 싱크를 사용하는 경우 부팅 디스크를 늘리세요. 시스템 지연(system lag)과 백로그 시간(backlog seconds)을 모니터링하세요.
스트리밍 엔진 및 셔플:
- 스트리밍 엔진은 상태와 셔플을 서비스 백엔드로 외부화하여 탄력성을 개선하고, 작업자 메모리 부담을 줄이며, 더 빠른 업데이트를 가능하게 합니다.
- 배치 중심의 단계나 대규모 키 그룹화의 경우, Dataflow Shuffle을 사용하여 작업자의 셔플 I/O 부담을 덜어주세요. 두 방법 모두 핫 워커(hot-worker) 장애와 디스크 스래싱(disk thrash)을 줄여줍니다.
역압(Backpressure), 핫 키(hot key), 스큐(skew):
- Dataflow는 동적 작업 재조정을 통해 역압(backpressure)을 관리합니다. 그럼에도 불구하고, 해당하는 경우 소스의 흐름 제어(예: Pub/Sub의 처리되지 않은 메시지/바이트)를 조정해야 합니다.
- 핫 키(예: 인기 있는 ID)는 처리 지연 작업(straggler)을 유발합니다. 키 샤딩(key#N), 부분 사전 집계 후 재키 지정(re-keying), 또는 스케치 기반 근사치를 통해 완화하세요.
- 특이 레코드(대용량 페이로드)나 버스트성 게시자로 인한 스큐(skew)는 게시자별 파티션, 배치 처리 또는 압축이 필요할 수 있습니다.
Pub/Sub 통합:
- 수집에는 Pub/Sub 주제를 사용하고, 메타데이터(예: deviceId, 이벤트 타임스탬프)를 위해 메시지 속성을 활성화하세요.
- PubSubIO로 수집하고, 속성이나 페이로드에서 이벤트 타임스탬프를 추출하세요. 그렇지 않으면 게시 시간으로 대체됩니다.
- 순서 지정 키(Ordering key)는 키별 순서를 제공하지만, Dataflow는 최소 한 번 전송(at-least-once delivery) 보장으로 인해 다운스트림에서 멱등성 동작이 필요합니다.
BigQuery로 스트리밍하는 패턴:
- 높은 처리량, 낮은 지연 시간, 그리고 스트림 오프셋과 자동 재시도를 통한 스트림 내 ‘정확히 한 번’ 시맨틱스를 위해 Storage Write API와 함께 BigQueryIO를 사용하는 것을 선호합니다.
- 처리율이 낮은 간단한 파이프라인의 경우 스트리밍 삽입도 가능합니다. 클라이언트 재시도 시 중복을 제거하기 위해 insertId를 설정하세요.
- 스트리밍 버퍼에 대한 쿼리는 최종적 일관성을 가집니다. 시간이 중요한 분석의 경우, 버퍼 지연 시간(예: 관찰된 가용성 지연 시간의 약 2배)을 기다린 후 쿼리하거나, 마이크로 배치 윈도우와 Storage Write API 커밋 모드를 통해 구체화하세요.
정확히 한 번(Exactly-once) 효과, 멱등성, 재처리, 싱크:
- Beam은 최소 한 번 처리(at-least-once processing)를 보장합니다. ‘정확히 한 번(exactly-once)‘은 싱크에서 멱등성 쓰기, 트랜잭션 또는 중복 제거 키를 사용하여 달성해야 합니다.
- BigQuery: 스트림 내에서 정확히 한 번을 보장하려면 Storage Write API의 기본 스트림 또는 커밋된 스트림을 사용하세요. 스트리밍 삽입을 사용할 경우, 안정적인 insertId를 설정하세요.
- 파일: 고유한 이름으로 임시 파일을 쓰고, 윈도우 완료 시 확정하며, 원자적 이름 변경을 보장하세요. 부분적인 중복을 방지하기 위해 덮어쓰기를 피하세요.
- 외부 데이터베이스: 안정적인 ID를 키로 사용하는 업서트(upsert)를 사용하거나 중복 제거 윈도우를 구현하세요.
- 재처리를 고려한 설계: 결정론적 변환을 유지하고, 싱크가 재시도 시 중복을 제거하도록 보장하세요.
데드 레터 처리, 오류 라우팅, 관측 가능성:
- 위험한 파싱/보강 로직은 ParDo 내에서 try/catch로 감싸고, 실패 시 TupleTag를 통해 데드 레터 PCollection으로 내보내세요. 페이로드, 오류 코드, 컨텍스트를 포함해야 합니다.
- 분석을 위해 DLQ(데드 레터 큐)를 BigQuery나 Cloud Storage로 라우팅하세요. 재처리를 위해 별도의 Pub/Sub 주제를 고려할 수 있습니다.
- 관측 가능성: Dataflow 작업 측정항목(워터마크 지연, 시스템 지연, 처리량), 커스텀 카운터, 분포 측정항목, 그리고 Cloud Logging의 단계별 로그를 사용하세요. Cloud Monitoring에서 지연 및 오류율에 대한 알림을 생성하세요. 예외를 집계하려면 Error Reporting을 사용하세요.
성능 튜닝 패턴:
- 효율적인 읽기: BigQuery 소스의 경우, 필요한 필드와 필터만 선택하는 Storage Read API 또는 쿼리 기반 읽기를 선호합니다.
- 결합 리프팅: GroupByKey 전에 컴바이너(combiner)를 사용하여 셔플 볼륨을 줄이세요.
- 사이드 입력: 작은 참조 데이터는 메모리에 캐시하세요. 팬아웃(fanout)과 업데이트 주기를 주시하세요.
- 직렬화: 압축된 스키마(Avro/Proto)를 사용하고, 핫 경로(hot path)에서 과도한 JSON 파싱을 피하세요.
배포, 템플릿, 업그레이드 전략
Flex 템플릿:
- 재현 가능한 배포를 위해 파이프라인을 컨테이너화되고 매개변수화된 템플릿으로 패키징합니다. Flex 템플릿은 커스텀 종속성, GPU 이미지, 환경 격리를 지원합니다.
- 런타임 매개변수(예: 입력 구독, 출력 테이블, 데드 레터 싱크, maxWorkers)를 외부화하여 환경별 배포를 지원합니다.
파이프라인 업데이트 및 호환성:
- Dataflow는 변환 이름, 상태 사양, 출력 유형이 호환성을 유지하는 경우 많은 스트리밍 파이프라인에 대해 인플레이스(in-place) 업데이트를 지원합니다. 안정적인 PTransform 이름을 사용하세요.
- 호환되지 않는 그래프 또는 상태 변경의 경우, 제어된 전환을 수행합니다. 새 작업을 시작한 다음, 이전 작업을 드레이닝(drain)하여 진행 중인 작업을 완료하고 새 요소 읽기를 중지합니다.
드레이닝(Draining) 및 스냅샷:
- 드레이닝은 처리를 정상적으로 완료하고, 남은 출력을 쓰며, 종료됩니다. 데이터 공백을 피하기 위해 Pub/Sub 보관 또는 스냅샷과 연계해야 합니다.
- 연속성을 보장하기 위해 Pub/Sub 스냅샷을 생성하고, 스냅샷 또는 적절한 타임스탬프로 이동하여 새 파이프라인을 시작하고, 출력을 확인한 다음, 이전 작업을 드레이닝할 수 있습니다.
구성 예시:
- 조기/지연 트리거 및 누적을 사용한 윈도우 예시:
undefined
- Storage Write API를 사용한 BigQueryIO 예시:
undefined
- 일반적인 함정:
- 스트리밍에서 윈도우 쓰기 없이 파일 기반 싱크에 쓰면 최종화가 지연될 수 있습니다. 윈도우 쓰기와 트리거를 활성화하세요.
- 무한한 증가: 상태나 허용된 지연 시간에 경계를 설정하는 것을 잊으면 메모리 누수 및 확장 실패를 유발할 수 있습니다.
- 타임스탬프 누락: 이벤트 타임스탬프를 할당하지 않으면 파이프라인이 기본적으로 처리 시간을 사용하게 되어 가변적인 지연 상황에서 정확성을 잃게 됩니다.
실용적인 문제 시나리오
NovaTrack Inc.는 50,000개의 온도 센서에서 글로벌 IoT 텔레메트리를 수집하고, 분 단위 집계 데이터를 제공하며, 원시 데이터를 영구 저장하고, 실시간 대시보드를 제공해야 합니다. 간헐적으로 형식이 잘못된 메시지와 순서가 맞지 않는 전송이 발생할 것으로 예상됩니다. 솔루션은 자동 확장되고, 검사를 위해 잘못된 레코드를 노출하며, 무중단 업그레이드를 지원해야 합니다.
접근 방식:
수집 및 시간 시맨틱
- 지역별 Pub/Sub 주제와 deviceId 및 eventTs(RFC3339) 속성을 가진 리전별 게시자를 생성합니다. 가능한 경우 deviceId를 기준으로 순서 지정 키를 활성화합니다.
- 근거: Pub/Sub은 최소 한 번 전송(at-least-once delivery)을 보장하는 내구성 있고 탄력적인 인그레스(ingress)를 제공합니다. 엣지에서 이벤트 타임스탬프를 첨부하면 실제 이벤트 시간을 보존할 수 있으며, 기기별 순서 지정은 중앙 병목 현상 없이 기기 내 순서 재정렬을 줄여줍니다.
이벤트 시간 윈도우를 사용한 Dataflow 스트리밍 파이프라인
- PubSubIO를 통해 전용 구독에서 읽고, eventTs를 Beam 타임스탬프로 추출하며, 누락된 경우 publishTime으로 대체합니다.
- 1분 단위의 FixedWindows를 적용하고, 30초에 조기 트리거를, 각 지연 요소에 대해 지연 실행을 설정합니다. 허용된 지연 시간은 10분으로 설정하고 창(pane)을 누적합니다.
- 근거: 이벤트 시간 윈도우는 정확한 분 단위 집계를 보장합니다. 조기 실행은 1분 미만의 최신 데이터로 대시보드를 채우고, 지연 실행은 지연된 데이터가 도착할 때 집계를 수정합니다. 지연 시간 경계는 상태 크기와 비용을 제한합니다.
유효성 검사, 보강 및 데드 레터 라우팅
- 작업 시작 시 BigQuery에서 로드된 사이드 입력(side input)을 통해 JSON을 파싱하고, 스키마와 범위를 검증하며, 작은 정적 참조 데이터로 보강하는 ParDo를 구현합니다.
- TupleTags를 사용하여 유효한 레코드는 주 출력으로, 실패한 레코드는 페이로드, 오류, deviceId, 파싱 타임스탬프를 포함하는 데드 레터 PCollection으로 내보냅니다. DLQ는 파티션된 BigQuery 테이블에 씁니다.
- 근거: 사이드 입력은 참조 데이터를 메모리에 유지하여 낮은 지연 시간을 보장합니다. 데드 레터 캡처는 주 흐름을 차단하지 않고 잘못된 행을 검사하고 대상 재처리를 가능하게 합니다.
집계 및 핫 키(hot key) 완화
- deviceId로 키를 지정하고 CombineFns를 사용하여 분당 평균/최소/최대값을 계산합니다. 상위 N개 지역 메트릭의 경우, 핫 키를 피하기 위해 region#N으로 샤딩한 다음 다시 집계합니다.
- 근거: Combiner는 셔플 볼륨과 비용을 최소화합니다. 키 샤딩은 지역별 팬인(fan-in) 중 단일 키 병목 현상을 방지합니다.
싱크 및 정확히 한 번(exactly-once) 효과
- Storage Write API를 사용하는 BigQueryIO를 통해 검증된 원시 이벤트와 분 단위 집계를 BigQuery에 씁니다. 커스텀 재시도 시 멱등성을 위해 deviceId + eventTs를 기반으로 안정적인 삽입 ID를 설정합니다.
- 근거: Storage Write API는 스트림 내에서 정확히 한 번(exactly-once) 시맨틱을 갖춘 높은 처리량과 낮은 지연 시간의 수집을 제공합니다. 안정적인 ID는 재실행이 발생할 경우 다운스트림에서 중복 제거를 보장합니다.
대시보드 일관성 전략
- 대시보드는 워터마크 대비 2분의 조회 기간(lookback) 또는 스트리밍 데이터에 대해 관찰된 가용성 지연 시간의 2배에 해당하는 고정 지연 시간을 두고 파티션된 집계 테이블을 쿼리합니다.
- 근거: BigQuery 스트리밍 가시성은 최종적 일관성을 가집니다. 읽기를 약간 지연시키면 거의 실시간 동작을 유지하면서 진행 중인 행이 누락되는 것을 방지할 수 있습니다.
운영: 자동 확장 및 스트리밍 엔진
- Streaming Engine을 활성화합니다. 예상 피크(예: 평균의 3배)를 기반으로 maxWorkers를 설정하고, CPU 집약적인 파싱 및 암호화에 맞는 머신 유형을 선택하며, 일시적인 셔플을 수용하기 위해 부팅 디스크를 늘립니다.
- 워터마크 지연, 백로그 시간(초), CPU, 단계별 처리량을 모니터링합니다. 지속적인 지연 및 DLQ 비율 급증 시 알림을 설정합니다.
- 근거: Streaming Engine은 상태/셔플을 외부화하여 탄력성과 더 간단한 업그레이드를 제공합니다. 적절한 크기 조정 및 모니터링은 조용한 SLO 위반을 방지합니다.
Flex 템플릿을 사용한 배포 및 업그레이드
- 파이프라인을 입력 구독, 출력 테이블, DLQ 테이블, maxWorkers, 리전 등의 매개변수를 가진 Flex 템플릿으로 패키징합니다. 호환되지 않는 변경의 경우, 동일한 주제를 대상으로 하는 새 구독으로 새 파이프라인을 시작하고, 출력을 확인한 다음, 이전 작업을 드레이닝합니다. 선택적으로 Pub/Sub 스냅샷을 생성하고 새 구독을 스냅샷으로 이동시켜 데이터 공백이 없도록 보장할 수 있습니다.
- 근거: Flex 템플릿은 반복 가능하고 매개변수화된 배포를 가능하게 합니다. 드레이닝을 통한 검증된 블루/그린(blue/green) 전환은 데이터 손실 제로와 최소한의 다운타임을 달성합니다.
재처리 및 배치 백필
- 사이드 출력을 통해 원시 이벤트의 압축된 Avro 파일을 Cloud Storage에 저장합니다. 모델이나 스키마가 변경될 때 배치 Dataflow 파이프라인을 실행하여 BigQuery로 백필(backfill)하거나 재처리합니다.
- 근거: 내구성 있는 원시 아카이브는 핫 패스(hot path)에 영향을 주지 않으면서 재현성과 스키마 진화를 지원합니다.
이 설계는 전 세계적인 규모에서 순서가 맞지 않는 데이터와 지연된 데이터를 처리하면서, 제한된 비용으로 정확하고 지연 시간이 짧은 집계를 생성하고, 명확한 오류 격리, 강력한 관측 가능성, 안전한 업그레이드 경로를 제공합니다.
← BigQuery 분석 및 웨어하우스 엔지니어링 · 모든 도메인 · 메시징 →
이 문제 연습하기 → · ExamRoll.io에서 시간 제한 연습 →
Pass the whole exam — not just this question
You found this answer. Get every verified question and explanation in one place, and save hours of prep. Free to start.
시험 합격하기 →