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 모델과 시간 시맨틱

실패 모드와 트레이드오프:

스트리밍 워크로드를 위한 Dataflow 운영

배포, 템플릿, 업그레이드 전략

undefined

undefined

실용적인 문제 시나리오

NovaTrack Inc.는 50,000개의 온도 센서에서 글로벌 IoT 텔레메트리를 수집하고, 분 단위 집계 데이터를 제공하며, 원시 데이터를 영구 저장하고, 실시간 대시보드를 제공해야 합니다. 간헐적으로 형식이 잘못된 메시지와 순서가 맞지 않는 전송이 발생할 것으로 예상됩니다. 솔루션은 자동 확장되고, 검사를 위해 잘못된 레코드를 노출하며, 무중단 업그레이드를 지원해야 합니다.

접근 방식:

  1. 수집 및 시간 시맨틱

    • 지역별 Pub/Sub 주제와 deviceId 및 eventTs(RFC3339) 속성을 가진 리전별 게시자를 생성합니다. 가능한 경우 deviceId를 기준으로 순서 지정 키를 활성화합니다.
    • 근거: Pub/Sub은 최소 한 번 전송(at-least-once delivery)을 보장하는 내구성 있고 탄력적인 인그레스(ingress)를 제공합니다. 엣지에서 이벤트 타임스탬프를 첨부하면 실제 이벤트 시간을 보존할 수 있으며, 기기별 순서 지정은 중앙 병목 현상 없이 기기 내 순서 재정렬을 줄여줍니다.
  2. 이벤트 시간 윈도우를 사용한 Dataflow 스트리밍 파이프라인

    • PubSubIO를 통해 전용 구독에서 읽고, eventTs를 Beam 타임스탬프로 추출하며, 누락된 경우 publishTime으로 대체합니다.
    • 1분 단위의 FixedWindows를 적용하고, 30초에 조기 트리거를, 각 지연 요소에 대해 지연 실행을 설정합니다. 허용된 지연 시간은 10분으로 설정하고 창(pane)을 누적합니다.
    • 근거: 이벤트 시간 윈도우는 정확한 분 단위 집계를 보장합니다. 조기 실행은 1분 미만의 최신 데이터로 대시보드를 채우고, 지연 실행은 지연된 데이터가 도착할 때 집계를 수정합니다. 지연 시간 경계는 상태 크기와 비용을 제한합니다.
  3. 유효성 검사, 보강 및 데드 레터 라우팅

    • 작업 시작 시 BigQuery에서 로드된 사이드 입력(side input)을 통해 JSON을 파싱하고, 스키마와 범위를 검증하며, 작은 정적 참조 데이터로 보강하는 ParDo를 구현합니다.
    • TupleTags를 사용하여 유효한 레코드는 주 출력으로, 실패한 레코드는 페이로드, 오류, deviceId, 파싱 타임스탬프를 포함하는 데드 레터 PCollection으로 내보냅니다. DLQ는 파티션된 BigQuery 테이블에 씁니다.
    • 근거: 사이드 입력은 참조 데이터를 메모리에 유지하여 낮은 지연 시간을 보장합니다. 데드 레터 캡처는 주 흐름을 차단하지 않고 잘못된 행을 검사하고 대상 재처리를 가능하게 합니다.
  4. 집계 및 핫 키(hot key) 완화

    • deviceId로 키를 지정하고 CombineFns를 사용하여 분당 평균/최소/최대값을 계산합니다. 상위 N개 지역 메트릭의 경우, 핫 키를 피하기 위해 region#N으로 샤딩한 다음 다시 집계합니다.
    • 근거: Combiner는 셔플 볼륨과 비용을 최소화합니다. 키 샤딩은 지역별 팬인(fan-in) 중 단일 키 병목 현상을 방지합니다.
  5. 싱크 및 정확히 한 번(exactly-once) 효과

    • Storage Write API를 사용하는 BigQueryIO를 통해 검증된 원시 이벤트와 분 단위 집계를 BigQuery에 씁니다. 커스텀 재시도 시 멱등성을 위해 deviceId + eventTs를 기반으로 안정적인 삽입 ID를 설정합니다.
    • 근거: Storage Write API는 스트림 내에서 정확히 한 번(exactly-once) 시맨틱을 갖춘 높은 처리량과 낮은 지연 시간의 수집을 제공합니다. 안정적인 ID는 재실행이 발생할 경우 다운스트림에서 중복 제거를 보장합니다.
  6. 대시보드 일관성 전략

    • 대시보드는 워터마크 대비 2분의 조회 기간(lookback) 또는 스트리밍 데이터에 대해 관찰된 가용성 지연 시간의 2배에 해당하는 고정 지연 시간을 두고 파티션된 집계 테이블을 쿼리합니다.
    • 근거: BigQuery 스트리밍 가시성은 최종적 일관성을 가집니다. 읽기를 약간 지연시키면 거의 실시간 동작을 유지하면서 진행 중인 행이 누락되는 것을 방지할 수 있습니다.
  7. 운영: 자동 확장 및 스트리밍 엔진

    • Streaming Engine을 활성화합니다. 예상 피크(예: 평균의 3배)를 기반으로 maxWorkers를 설정하고, CPU 집약적인 파싱 및 암호화에 맞는 머신 유형을 선택하며, 일시적인 셔플을 수용하기 위해 부팅 디스크를 늘립니다.
    • 워터마크 지연, 백로그 시간(초), CPU, 단계별 처리량을 모니터링합니다. 지속적인 지연 및 DLQ 비율 급증 시 알림을 설정합니다.
    • 근거: Streaming Engine은 상태/셔플을 외부화하여 탄력성과 더 간단한 업그레이드를 제공합니다. 적절한 크기 조정 및 모니터링은 조용한 SLO 위반을 방지합니다.
  8. Flex 템플릿을 사용한 배포 및 업그레이드

    • 파이프라인을 입력 구독, 출력 테이블, DLQ 테이블, maxWorkers, 리전 등의 매개변수를 가진 Flex 템플릿으로 패키징합니다. 호환되지 않는 변경의 경우, 동일한 주제를 대상으로 하는 새 구독으로 새 파이프라인을 시작하고, 출력을 확인한 다음, 이전 작업을 드레이닝합니다. 선택적으로 Pub/Sub 스냅샷을 생성하고 새 구독을 스냅샷으로 이동시켜 데이터 공백이 없도록 보장할 수 있습니다.
    • 근거: Flex 템플릿은 반복 가능하고 매개변수화된 배포를 가능하게 합니다. 드레이닝을 통한 검증된 블루/그린(blue/green) 전환은 데이터 손실 제로와 최소한의 다운타임을 달성합니다.
  9. 재처리 및 배치 백필

    • 사이드 출력을 통해 원시 이벤트의 압축된 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.

시험 합격하기 →

Google 찾아보기 →

Related guides

올인원 액세스

하나의 구독. 모든 시험.

모든 플랜은 무제한 답변 검색, 모의고사, AI 해설, 전체 자료 라이브러리를 20개 이상의 언어로 잠금 해제합니다.

월간
24.87
Just €0.83/day
모든 포함:
  • 무제한 답변 검색
  • 무제한 모의고사
  • AI 기반 해설
  • 전체 자료 라이브러리
  • 20개 이상의 언어
  • 주간 콘텐츠 업데이트
  • 보상 및 추천
  • 우선 지원
무료 체험 시작

신용카드 필요 없음*

최고의 가치
12개월
179.87
Just €0.49/daySave 40%
모든 포함:
  • 무제한 답변 검색
  • 무제한 모의고사
  • AI 기반 해설
  • 전체 자료 라이브러리
  • 20개 이상의 언어
  • 주간 콘텐츠 업데이트
  • 보상 및 추천
  • 우선 지원
무료 체험 시작

신용카드 필요 없음*

✓ 무료 플랜 포함 · ✓ 언제든지 취소 가능 · ✓ 모든 플랜은 전체 제품을 잠금 해제합니다