Google PDE: 워크플로 오케스트레이션 및 파이프라인 자동화 — 학습 가이드
다음의 일부입니다: Google Professional Data Engineer — 학습 가이드. 검증된 답안으로 연습하기: Google 시험 허브, 또는 다음에서 시간 제한 모의고사 풀기: ExamRoll.io.
개요
워크플로 오케스트레이션과 파이프라인 자동화는 여러 서비스에 걸친 데이터 작업을 조정하여 수집, 변환, 품질 검사, 게시가 안정적이고 안전하며 비용 효율적으로 이루어지도록 합니다. Google Cloud에서 오케스트레이션은 각 워크로드의 실행 모델(예약된 배치, 이벤트 기반 스트림, 애드혹 또는 장기 실행 작업)에 맞춰야 합니다. 설계 목표는 반복성, 멱등성, 관측 가능성, 최소 권한, 그리고 환경 간의 안전한 프로모션입니다.
주요 선택 사항:
- DAG, 태스크 종속성, 고급 스케줄링을 위해 Cloud Composer(Apache Airflow)를 사용한 코드 중심의 배치 오케스트레이션.
- 경량의 이벤트 기반 교차 서비스 시퀀스를 위한 Cloud Workflows를 사용한 서버리스 API 코레오그래피.
- cron을 위한 Cloud Scheduler 또는 이벤트를 위한 Eventarc에 의해 트리거되는 Cloud Run jobs 또는 Dataproc jobs와 같은 실행 엔드포인트.
- BigQuery 변환, 어설션, 릴리스 관리를 위한 Dataform을 사용한 SQL 네이티브 오케스트레이션.
운영 모델은 제한된 지수 백오프를 사용한 재시도, 타임아웃, SLA, 따라잡기(catchup) 및 백필, 안전한 재실행을 위한 멱등성 태스크 설계, 데드 레터 캡처를 통한 강력한 장애 처리를 강조합니다. 보안은 파이프라인별 서비스 계정, 보안 비밀 격리, 매개변수화, 최소 권한 IAM을 통해 강화됩니다. CI/CD, 코드형 인프라, 포괄적인 원격 측정은 프로덕션에 즉시 사용 가능한 접근 방식을 완성합니다.
Google Cloud에서의 오케스트레이션: 도구 및 패턴
Cloud Composer (Airflow)
- DAG는 명시적인 종속성을 가진 방향성 비순환 실행 그래프를 정의합니다. TaskFlow API 또는 오퍼레이터(예: BigQuery, Dataflow, Dataproc, Cloud Run)를 사용하여 태스크를 표현합니다. 센서와 지연 가능한(deferrable) 오퍼레이터는 대기 조건(예: Cloud Storage의 객체 완료 또는 BigQuery의 파티션 생성)에 대한 스케줄러 부하를 줄입니다.
- 스케줄링: cron 표현식, start_date, end_date, catchup은 과거 실행을 제어합니다. 백필에는 catchup을 사용하고, 스트리밍에 인접하거나 멱등성이 없는 대상에는 비활성화합니다. 다운스트림 시스템을 보호하기 위해 max_active_runs와 풀(pool)로 동시성을 제한합니다.
- 종속성: set_upstream/set_downstream 또는 태스크플로우 종속성. 메타데이터 기반 오케스트레이션을 위해, 동적 태스크 매핑을 사용하여 BigQuery 제어 테이블(예: 클라이언트/파티션 목록)에서 동적으로 태스크를 생성하여 DAG 파싱 시간을 안정적으로 유지하고 태스크를 데이터 기반으로 만듭니다.
- 예시 (간결한) DAG 프래그먼트:
undefined
undefined
undefined
undefined
undefined
undefined
undefined
undefined
undefined
Cloud Workflows, Cloud Scheduler, Cloud Run jobs 및 이벤트 기반 실행
- Cloud Workflows는 내장된 재시도, 루프, 병렬 분기, 보상 로직을 사용하여 Google API 및 HTTP 엔드포인트를 오케스트레이션합니다. BigQuery, Dataflow, Batch, Cloud Run jobs와 같은 서비스 전반에 걸친 가벼운 제어 흐름에 이상적입니다.
- Cloud Scheduler는 cron 스타일 자동화를 위해 Workflows, Pub/Sub 주제 또는 HTTP 서비스를 트리거합니다. 매일 02:00에 배치 작업을 실행하려면 Dataflow 작업이나 Dataproc 작업을 시작하는 Workflow를 예약합니다.
- Cloud Run jobs는 자동 재시도 기능과 최소한의 운영으로 컨테이너화된 배치 단계를 실행합니다. 다단계 데이터 작업이나 Dataflow 또는 BigQuery 전후 처리 작업을 위해 Workflows와 잘 연동됩니다.
- 이벤트 기반: Eventarc를 사용하여 Cloud Storage 객체 완료, Pub/Sub 메시지 또는 감사 로그를 Cloud Run 또는 Workflows로 라우팅합니다. 단일 테이블에 대한 BigQuery 삽입 작업 알림을 받으려면, 고급 필터가 있는 Cloud Logging 싱크를 Pub/Sub으로 생성한 다음 해당 주제에서 컨슈머를 트리거합니다.
Dataform: BigQuery를 위한 SQL 워크플로
- ref()로 종속성 그래프를 모델링하고, 테이블/뷰/증분(incremental)을 정의하며, 태그나 일정에 따라 빌드를 오케스트레이션합니다. Dataform은 SQLX를 순서가 지정된 실행 계획으로 컴파일하여 선언적 정의로부터 메타데이터 기반 오케스트레이션을 가능하게 합니다.
- 어설션은 데이터 품질을 보장합니다. 어설션은 통과하기 위해 0개의 행을 반환해야 하는 쿼리입니다. 예시 어설션: – definitions/assert_non_negative_prices.sqlx
undefined
undefined
- 릴리스 및 리포지토리 제어: 코드를 리포지토리에 저장하고, 브랜치와 리뷰를 사용하며, 환경별 변수를 사용하여 태그가 지정된 릴리스를 환경(예: dev, test, prod)으로 프로모션합니다. CI/CD 검사 및 어설션 결과를 통해 배포를 제어합니다.
Dataproc, Dataflow 및 스토리지 패턴
- 최소한의 운영으로 Hadoop/Spark를 재사용하려면, GCS 커넥터와 함께 Dataproc을 사용하여 클러스터 수명을 넘어 데이터를 유지하고 영구 디스크 비용을 최소화합니다. 격리 및 비용 제어를 위해 작업별로 임시(ephemeral) 클러스터를 생성하고, Composer 또는 Workflows로 오케스트레이션합니다.
- 형식이 잘못된 행이 있는 배치 수집의 경우, Dataflow를 실행하여 유효한 레코드를 BigQuery에 쓰고, 파싱/유효성 검사 오류는 검사를 위해 데드 레터 BigQuery 테이블로 라우팅합니다.
관측성, 알림, 런북
텔레메트리 및 알림
- 모든 오케스트레이션 로그를 구조화된 필드(pipeline, dag_id, run_id, task_id, partition)와 함께 Cloud Logging으로 라우팅합니다. 로그 기반 측정항목을 통해 오류 로그를 Monitoring으로 내보냅니다. 다음에 대해 알림을 설정합니다:
- 누락된 스케줄 또는 SLA 위반
- 연속적인 태스크 실패
- 백로그 증가 (예: Pub/Sub 미확인 메시지, Dataflow 시스템 지연)
- 데이터 품질 어설션 실패
- Cloud Composer: DAG/태스크 기간, 성공률, 큐 깊이, 스케줄러 상태를 모니터링합니다. 페이징 및 복구 런북을 위해 on_failure_callback을 구성합니다.
- Cloud Workflows: 실행 로그와 단계별 지연 시간을 검사하고, 명시적인 재시도 및 오류 핸들러를 추가하며, 상관관계 ID를 포함한 커스텀 로그를 내보냅니다.
- BigQuery 테이블 변경 알림: 특정 테이블을 대상으로 하는 삽입 작업에 대한 고급 필터를 사용하여 프로젝트 수준의 Logging 싱크를 생성하고 Pub/Sub으로 내보냅니다. 모니터링 도구는 이 토픽을 구독하여 다른 테이블의 노이즈 없이 즉각적인 알림을 받습니다.
런북 설계
- 각 파이프라인에 대해 트리거, 종속성, SLA, 롤백/재시도 절차, 안전한 백필 단계를 문서화합니다. Dataflow를 위한 “고정 데이터셋 재실행”, 스트리밍 작업 드레이닝 방법, 실패한 파티션 재처리 방법, DLQ 메시지 해결 방법을 포함합니다.
- 일반적인 실패 시그니처(예: 권한 거부, 할당량 초과, 스키마 불일치)를 의사 결정 트리 및 에스컬레이션 경로와 함께 캡처합니다.
실제 문제 시나리오
Acme Retail Analytics는 때때로 형식이 잘못된 행을 포함하는 파트너의 일일 CSV 파일을 수집하고, 유효한 데이터를 변환하여 BigQuery에 로드하며, 잘못된 행은 조사를 위해 표면화해야 합니다. 또한 거의 실시간에 가까운 가격 업데이트를 위한 이벤트 기반 보강과 개발 환경에서 프로덕션 환경으로의 안전한 승격을 원합니다.
접근 방식:
스토리지 및 이벤트 트리거
- 객체 버전 관리 및 균일한 버킷 수준 액세스가 설정된 전용 Cloud Storage 버킷을 생성합니다. Eventarc를 통해 객체 완료 알림을 Pub/Sub으로 보내도록 활성화합니다.
- 근거: 객체 완료는 다운스트림 수집을 트리거하는 신뢰할 수 있는 이벤트이며, 버전 관리는 재실행 및 감사를 지원합니다.
데드 레터 처리를 포함한 배치 수집
- Cloud Composer를 사용하여 catchup이 활성화된 일일 Airflow DAG를 02:00에 스케줄링합니다. 이 DAG는 CSV를 파싱하고 스키마를 검증하며, 결정론적 스테이징 테이블을 사용하여 유효한 레코드를 BigQuery에 쓴 다음, 파티션된 대상 테이블에 MERGE하는 Dataflow 배치 작업을 시작합니다. 형식이 잘못되거나 실패한 레코드는 BigQuery 데드 레터 테이블로 라우팅합니다.
- 근거: Dataflow는 파싱/검증을 확장하고, MERGE는 멱등성을 보장하며, 데드 레터 캡처는 파이프라인을 차단하지 않고 검사를 지원하여 형식이 잘못된 행에 대한 권장 패턴과 일치합니다.
이벤트 기반 보강
- 증분 가격 업데이트를 위한 경량 보강을 수행하는 Cloud Run 작업을 배포합니다. 낮 동안 작은 업데이트 파일이 도착하면 Eventarc의 Pub/Sub 메시지를 수신하는 Cloud Workflows를 통해 트리거합니다.
- 근거: Workflows와 함께 사용하는 서버리스 컨테이너는 작은 이벤트에 대해 낮은 지연 시간과 적은 운영 부담의 오케스트레이션을 제공하면서, 무거운 변환 작업은 배치로 유지합니다.
안정성 제어
- Dataflow 및 Cloud Run 작업의 일시적인 실패에 대해 지수 백오프를 사용한 재시도를 구성하고, 총 재시도 시간을 DAG SLA에 맞게 제한합니다. Airflow에서는 태스크별 실행 시간 초과와 on_failure 콜백을 설정하고, Workflows에서는 max_doublings와 max_retry_duration을 설정합니다.
- 근거: 제한된 백오프는 SLA를 보존하고 통제 불가능한 재시도를 방지합니다.
보안 및 최소 권한 원칙
- 각 구성 요소를 전용 서비스 계정(Composer 오케스트레이터 SA, Dataflow 워커 SA, Cloud Run 작업 SA)으로 실행합니다. Dataflow에는 수집 버킷에 대한 GCS 읽기 권한, 대상 데이터셋에 대한 BigQuery dataEditor 권한, 로그에 대한 Viewer 권한 등 필요한 역할만 부여합니다. 보안 비밀은 Secret Manager에 저장하고 런타임에 참조합니다.
- 근거: 최소 권한 원칙을 적용하고 장애 영향을 격리합니다.
메타데이터 기반 오케스트레이션
- 파트너 소스, 파일 패턴, 대상 데이터셋을 나열하는 BigQuery 제어 테이블을 유지 관리합니다. DAG 런타임에 Airflow는 이 테이블을 쿼리하고 동적 태스크 매핑을 사용하여 파트너별 태스크를 생성합니다.
- 근거: 파트너 추가가 코드 변경이 아닌 데이터 변경이 되어 배포 위험을 줄입니다.
관측성 및 알림
- run_id와 partner_id를 포함한 구조화된 로그를 내보냅니다. DAG SLA 위반, Dataflow 시스템 지연, 비어 있지 않은 데드 레터 수에 대한 알림 정책을 생성합니다. 대상 테이블에 대한 BigQuery 삽입의 경우, 해당 테이블에 대한 고급 필터를 사용하여 Cloud Logging 싱크를 구성하고, Acme의 모니터링 도구가 소비하는 Pub/Sub 토픽으로 보냅니다.
- 근거: 세분화된 알림을 통해 노이즈 없이 신속한 문제 분류가 가능합니다.
CI/CD 및 승격
- 인프라(버킷, Pub/Sub, Eventarc, Composer, Workflows, BigQuery 데이터셋)를 Terraform으로 관리합니다. Cloud Build를 사용하여 Airflow DAG 구문을 검증하고, 단위 테스트를 실행하며, 개발 Composer 환경에 배포합니다. Dataform 어설션 및 통합 테스트가 통과된 후, 파라미터화된 구성과 수동 승인 게이트를 사용하여 테스트 및 프로덕션 환경으로 승격합니다.
- 근거: 선언적이고 반복 가능한 배포와 환경 간의 안전한 승격을 보장합니다.
런북 및 복구
- 특정 날짜를 재실행하는 단계를 문서화합니다: 객체 버전 관리에서 CSV를 복원하고, 해당 파티션에 대해 Dataflow 작업을 재실행하며, 결과를 MERGE하고, DLQ 레코드를 검토합니다. 불일치가 발생할 경우 변환 버그를 격리하기 위한 “고정 데이터셋 재실행” 절차를 포함합니다.
- 근거: 멱등성 설계와 문서화된 복구 절차는 부분적인 장애 복구를 간소화합니다.
이 문제 연습하기 → · 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.
시험 합격하기 →