Google PDE: 메시징, 이벤트 수집 및 실시간 서비스 — 학습 가이드
다음의 일부입니다: Google Professional Data Engineer — 학습 가이드. 검증된 답안으로 연습하기: Google 시험 허브, 또는 다음에서 시간 제한 모의고사 풀기: ExamRoll.io.
개요
Google Cloud의 메시징, 이벤트 수집, 실시간 서비스는 분리되고 내구성 있는 전송을 위한 Cloud Pub/Sub 및 Eventarc, 상태 저장 스트림 처리를 위한 Dataflow, 그리고 BigQuery, Cloud Storage, 운영 데이터베이스와 같은 싱크(sink)를 중심으로 합니다. 최소 한 번 이상 전달(at-least-once delivery), 멱등적 소비(idempotent consumption), 관측 가능성(observability)을 고려하여 설계하면 장애, 백프레셔(backpressure), 스키마 변경 상황에서도 정확성을 유지하며 탄력적으로 확장되는 복원력 있는 시스템을 보장할 수 있습니다.
Pub/Sub을 이용한 핵심 메시징
- Topic과 Subscription
- 게시자(Publisher)는 Topic으로 메시지를 보내고, 구독자(Subscriber)는 Subscription을 통해 연결됩니다 (여러 구독자가 동일한 메시지를 독립적으로 소비할 수 있습니다).
- Subscription 유형:
- Pull(가져오기): 클라이언트가 명시적으로 메시지를 가져옵니다. 가장 높은 처리량과 적은 왕복(round trip)을 위해 스트리밍 Pull을 사용합니다.
- Push(밀어넣기): Pub/Sub이 HTTPS를 통해 전달합니다. 엔드포인트는 확인(acknowledge)을 위해 2xx 상태 코드를 반환해야 합니다.
- BigQuery로 내보내기: BigQuery Subscription은 코드 없이 BigQuery 테이블로 메시지를 전달합니다. 페이로드가 선언된 스키마와 일치하고 분석 시스템으로의 낮은 지연 시간 수집이 필요할 때 가장 좋습니다.
- 순서 지정 키(Ordering key)
- Topic과 Subscription에서 메시지 순서 지정을 활성화하여 순서 지정 키별로 순서에 맞는 전달을 보장받을 수 있습니다. 키당 처리량은 직렬화됩니다. 즉, 키당 하나의 처리 중인(in-flight) 메시지가 후속 메시지를 차단할 수 있으므로, 확장성을 위해 많은 키(예: hash(device_id))를 사용해야 합니다.
- 팬아웃(Fan-out) 및 재재생(Replay)
- 워크로드와 보관 기간을 분리하기 위해 각기 다른 소비자를 위한 별도의 Subscription을 생성합니다.
- 복구 및 백필(backfill)을 위해 타임스탬프나 스냅샷(snapshot)으로부터 재재생하려면 탐색(seek) 또는 스냅샷을 사용합니다.
장단점:
- 순서 지정은 병렬성과 키당 처리량을 감소시키므로, 반드시 필요한 경우가 아니면 비활성화합니다.
- Push는 클라이언트 코드를 단순화하지만 HTTP 엔드포인트 확장, 보안, 백오프(backoff) 문제를 야기합니다. Pull은 높은 처리량에서 더 많은 제어권과 안정성을 제공합니다.
전달 시맨틱, 확인(Acknowledgment), 보관(Retention), 데드 레터링(Dead Lettering)
- 확인(Acknowledgment) 및 기한(Deadline)
- 최소 한 번 이상 전달(At-least-once delivery): 중복이 발생할 수 있습니다.
- 각 전달에는 ack 기한(ack deadline, 기본값 10초)이 있습니다. 장기 실행 작업을 처리하는 동안에는 기한을 연장(ModifyAckDeadline)해야 합니다. 기한 내에 ack를 보내지 못하는 것은 중복된 Push 전달의 가장 흔한 원인입니다.
- Nack(부정 확인) 또는 기한 만료 시 메시지는 재전달 대상이 됩니다.
- 보관(Retention)
- 확인되지 않은 메시지는 Subscription의 ack 기한 동안 보관되고 재시도됩니다. 확인된 메시지는 재재생을 위해 Topic의 메시지 보관 기간까지 보관될 수 있습니다. 최대 장애 시간과 복구 시간을 감당할 수 있도록 보관 기간을 구성하십시오.
- 재시도(Retry)
- Pull: ack 기한이 만료된 후 재전달이 발생합니다. 흐름 제어(flow control) 한도로 동시성을 제어합니다.
- Push: 지수 백오프(exponential backoff)를 사용하며, HTTP 2xx 응답만 성공으로 간주합니다. 3xx/4xx/5xx는 재시도를 유발합니다. 반복을 허용하기 위해 멱등성 핸들러(idempotent handler)를 구현하십시오.
- 데드 레터 토픽(Dead-letter topic, DLT)
- 포이즌 메시지(poison message)를 격리하기 위해 Subscription별로 DL 토픽과 최대 전달 시도 횟수를 구성합니다.
- DLQ(데드 레터 큐) 볼륨을 모니터링하고, 분류 워크플로를 생성하며, 수정 후 기본 Topic으로 다시 게시합니다.
예시:
undefined
전달 시맨틱 요약:
- Pub/Sub: 최소 한 번 이상 전달(at-least-once), 순서 지정 키가 활성화된 경우 키 내에서 최선 노력 순서 지정(best-effort ordering).
- 싱크(Sink): BigQuery 삽입 API는 중복 완화 기능(insertId 또는 Storage Write API 스트림 오프셋)을 제공하지만, 그럼에도 불구하고 소비자와 작성자는 멱등성을 갖도록 설계해야 합니다.
스키마, 호환성, 유효성 검사
- Pub/Sub 스키마
- 중앙에 저장된 스키마를 사용하여 Avro 및 Protocol Buffers를 기본적으로 지원합니다.
- Topic 수준 스키마 설정: 인코딩(Avro 또는 Protobuf) 및 적용 수준(없음, 유효성 검사만, 또는 필수).
- 생산자(Producer)는 인코딩된 페이로드를 게시합니다. 적용이 활성화된 경우 Pub/Sub은 현재 스키마에 대해 유효성을 검사합니다.
- 진화 및 호환성
- 하위 호환성을 유지하는 변경을 사용합니다 (선택적 필드 추가, Avro에서 기본값이 있는 필드 추가, Protobuf에서 태그 재사용 금지, 필드 제거 또는 이름 변경 방지).
- 스키마 버전을 명시적으로 관리합니다. 호환성이 손상되는 변경(breaking change)의 경우, v1 및 v2 Topic에 이중으로 게시하거나 버전 필드를 추가하여 그에 따라 라우팅합니다.
- 생산자-소비자 계약
- 소비자(Consumer)는 알 수 없는 필드를 무시하고 누락된 필드는 기본값으로 처리해야 합니다.
- 프로덕션으로 승격하기 전에 모든 소비자에 걸쳐 스키마 호환성을 테스트하고, 프로덕션과 동일한 스키마 적용 수준을 가진 스테이징 Subscription에서 유효성을 검사합니다.
간단한 Avro 예시 (발췌):
undefined
이벤트 기반 통합, Eventarc 및 Kafka 상호 운용성
- Eventarc와 CloudEvents
- Eventarc는 CloudEvents 사양을 사용하여 Google Cloud 서비스(및 Pub/Sub를 통한 커스텀 소스)의 이벤트를 Cloud Run, GKE 또는 Workflows로 라우팅합니다. type, source, subject과 같은 속성은 세분화된 필터링과 감사 가능성을 지원합니다.
- 속성 필터를 사용하여 팬아웃(fan-out)을 최소화하고 다운스트림 부하를 줄이세요.
- 전달은 최소 한 번(at-least-once) 보장되므로, 가능한 경우 핸들러를 멱등성(idempotent) 있고 상태 비저장(stateless)으로 만드세요.
- Eventarc 트리거 예시:
gcloud eventarc triggers create gcs-finalize-to-run
–destination-run-service=ingestor
–event-filters=“type=google.cloud.storage.object.v1.finalized”
–event-filters=“bucket=my-data-bucket”
–service-account=eventarc-sa@PROJECT_ID.iam.gserviceaccount.com - Kafka 상호 운용성 및 관리형 마이그레이션
- Dataflow 템플릿은 단계적 마이그레이션을 위해 Kafka와 Pub/Sub를 연결합니다. 키를 보존하여 토픽을 미러링하고, 소비자를 먼저 전환한 후 생산자를 전환하거나, 전환 기간 동안 이중 쓰기(dual-write)를 수행합니다.
- Pub/Sub Lite는 키 기반 라우팅과 저렴한 비용으로 파티션 분할 및 용량 프로비저닝 스트리밍을 제공합니다. 리전/영역별 서비스이며, 예측 가능한 용량과 파티션별 순서 보장이 주요 관심사인 Kafka와 유사한 워크로드에 적합합니다.
- 마이그레이션 고려 사항:
- 순서 지정: Kafka 키를 Pub/Sub 순서 지정 키 또는 Lite 파티션에 매핑합니다.
- 오프셋: 진단을 위해 오프셋을 메시지 속성으로 전달합니다. 마이그레이션 후 소비자는 Kafka 오프셋에 의존할 수 없습니다.
- 전달: 최소 한 번(at-least-once) 전달을 허용하고 다운스트림에서 멱등성을 강제합니다.
- 스키마: Confluent Schema Registry 정의를 Pub/Sub 스키마로 마이그레이션하거나, 호환 가능한 발전 규칙을 사용하여 Protobuf/Avro로 표준화합니다.
스트리밍 수집 패턴, 처리량, 확장, 보안 및 운영
- 실시간 수집 패턴
- Pub/Sub -> Dataflow -> BigQuery: 높은 처리량과 스트림 오프셋을 통한 멱등성을 위해 BigQuery Storage Write API 싱크를 사용합니다. 실패한 메시지는 검사를 위해 데드-레터 테이블로 라우팅합니다.
- Pub/Sub -> Dataflow -> Cloud Storage: 재처리를 위해 원시 이벤트를 보관합니다. 비용과 지연 시간의 균형을 맞추기 위해 윈도우 기반의 압축 쓰기를 사용합니다.
- Pub/Sub -> 운영 저장소: 짧은 지연 시간의 조회를 위해 Bigtable에 쓰거나, 강력한 일관성의 트랜잭션을 위해 Spanner에 쓰거나, 워크로드 요구에 따라 Cloud SQL/Firestore에 씁니다. 고유 이벤트 ID를 키로 사용하여 멱등성 있는 업서트(upsert)를 보장합니다.
- 최소 한 번 전송, 중복 방지, 멱등성
- 모든 메시지에 고유한
event_id와event_time을 포함하고, 생산자 측에서 UUID를 강제합니다. - BigQuery 스트리밍 중복 제거:
insertId를 설정하거나 순서가 지정된 스트림과 함께 Storage Write API를 사용합니다. 그래도 쿼리에 중복 제거 로직을 포함하여 방어적으로 설계합니다. - 쿼리 시점 중복 제거 예시: WITH ranked AS ( SELECT t.*, ROW_NUMBER() OVER (PARTITION BY event_id ORDER BY event_time DESC) AS rn FROM dataset.events t ) SELECT * EXCEPT(rn) FROM ranked WHERE rn = 1;
- 푸시 엔드포인트의 경우, 성공적으로 처리된 후에만 2xx를 반환합니다. 그렇지 않으면 재전송될 것으로 예상해야 합니다.
- 모든 메시지에 고유한
- 메시지 처리량, 할당량, 확장
- 게시자: 메시지를 일괄 처리하고 연결을 재사용합니다. 여러 클라이언트에 걸쳐 병렬화합니다. 순서가 지정된 워크로드를 확장하려면 많은 순서 지정 키를 사용합니다.
- 구독자: 흐름 제어(최대 미처리 바이트/메시지)가 있는 스트리밍 풀(pull)을 선호합니다. 처리 시간에 맞게 확인(ack) 기한을 정하고 필요할 때 연장합니다.
- 데이터 양이 증가함에 따라 게시 및 구독 처리량에 대한 할당량 증가를 모니터링하고 요청합니다. 급증(burst)을 흡수할 수 있도록 여유 공간(예: 예상 피크의 2배)을 두고 설계합니다.
- 일관성과 가용성
- BigQuery 스트리밍은 쿼리 가시성에 대해 최종적 일관성을 가집니다. 스트리밍된 행을 반드시 포함해야 하는 대화형 쿼리의 경우, 관찰된 지연 시간(예: P50 가용성 지연의 2배)을 기반으로 대기하거나, Dataflow에서 워터마크에 정렬된 집계를 사용하고 구체화된 결과를 쿼리하도록 설계합니다.
- 보안
- IAM: 최소 권한 역할을 부여합니다(토픽에 대한
pubsub.publisher는 생산자에게, 구독에 대한pubsub.subscriber는 소비자에게). 각 워크로드에 전용 서비스 계정을 사용합니다. - 푸시 인증: 서비스 계정에서 OIDC 토큰을 첨부하도록 푸시 구독을 구성합니다. 엔드포인트에서 대상(audience) 검증을 강제합니다. 내장된 인증 및 TLS를 위해 Cloud Run 비공개 엔드포인트를 선호합니다.
- 암호화: Pub/Sub는 전송 중 및 저장 시 데이터를 암호화합니다. 고객 관리형 키를 사용하려면 토픽에 CMEK를 사용합니다. 데이터 유출 위험을 줄이기 위해 VPC 서비스 제어를 적용합니다. 필요한 경우 민감한 페이로드 필드에 클라이언트 측 암호화를 사용합니다.
- IAM: 최소 권한 역할을 부여합니다(토픽에 대한
- 지연, 재전송, 구독자 실패의 운영 진단
- Cloud Monitoring으로 모니터링:
- 백로그 확인을 위한
subscription/num_undelivered_messages및oldest_unacked_message_age. - 중복을 유발하는 누락된 확인(ack)을 감지하기 위한
expired_ack_deadline_count. - 처리량 확인을 위한
publish_request_count및pull_request_count.
- 백로그 확인을 위한
- 누락된 대시보드 이벤트를 조사하려면, 알려진 데이터 세트를 파이프라인을 통해 재실행하고 단계별 출력을 비교하여 결함이 있는 변환 또는 싱크를 격리합니다.
- Dataflow 스트리밍의 경우:
- 여러 소스로부터의 부하를 흡수하기 위해 적절한
maxWorkers를 설정하여 자동 확장을 사용합니다. - 호환되지 않는 업데이트의 경우 파이프라인을 드레이닝(draining)하여 진행 중인 작업이 완료되도록 하고 데이터 손실을 방지합니다.
- 여러 소스로부터의 부하를 흡수하기 위해 적절한
- BigQuery 삽입 알림의 경우, 특정 테이블로 필터링된 싱크를 통해 Cloud Logging 감사 항목을 Pub/Sub 토픽으로 라우팅하여 알림을 보냅니다.
- Cloud Monitoring으로 모니터링:
실제 문제 시나리오
Contoso Freight는 트럭에서 분당 10,000개의 IoT 원격 측정 메시지를 수집하고, 이벤트를 보강하며, 대화형 분석을 지원하고, 외부 파트너의 파일 드롭 시 워크플로를 트리거하기 위한 글로벌 실시간 이벤트 플랫폼이 필요합니다. 일부 파트너 CSV에는 형식이 잘못된 행이 포함되어 있으며, 분석팀은 스트림을 차단하지 않고 오류를 검사해야 합니다.
- 핵심 메시징 및 스키마 레이어 생성
- 조치: 원격 측정을 위한 Avro 스키마를 정의하고, 스키마 적용이 필수로 설정된 Pub/Sub 토픽
telemetry에 연결합니다. 메시지 순서 지정을 활성화하고ordering_key = hash(device_id)로 게시합니다. - 근거: 토픽 수준의 스키마 적용은 형식이 잘못된 이벤트를 조기에 거부합니다. 기기별 순서 지정은 필요할 때 순서 있는 처리를 지원하며, 해싱은 키를 분산시켜 처리량을 유지합니다.
- 격리 및 데드-레터링 기능이 있는 구독 프로비저닝
- 조치: Dataflow용 풀(pull) 구독
telemetry-stream-sub을 생성하고, 데드-레터 토픽telemetry-dlt와max_delivery_attempts=10을 설정합니다. 계보 및 재실행을 위해 원시 이벤트를 시간 파티션 테이블에 저장하는 BigQuery 구독telemetry-raw-bq를 추가합니다. - 근거: DLQ는 조사를 위해 포이즌 메시지를 격리합니다. 별도의 BigQuery 구독은 처리 파이프라인과 독립적으로 원시 이벤트를 보존하기 위한 운영 부담이 적은 내보내기 경로를 제공합니다.
- 보강 및 싱크를 위한 Dataflow 스트리밍 파이프라인 구축
- 조치: 흐름 제어가 있는 스트리밍 풀을 사용하여
telemetry-stream-sub에서 수집합니다. 스키마에 대해 유효성을 검사하고, 참조 데이터로 보강하며, 윈도우 기반 집계를 계산합니다. 명명된 스트림과insertId = event_id를 사용하여 Storage Write API를 통해 BigQuery에 씁니다. 시간당 원시 백업을 Cloud Storage에 씁니다. 잘못되거나 실패한 레코드는 데드-레터 BigQuery 테이블로 리디렉션합니다. - 근거: Storage Write API는
insertId/스트림 오프셋을 통해 멱등성을 갖춘 높은 처리량과 짧은 지연 시간의 쓰기를 제공합니다. 데드-레터 테이블은 스트림을 차단하지 않고 검사를 지원하며, Cloud Storage 아카이브는 재실행을 가능하게 합니다.
- 분석에서 중복 및 최종적 일관성 처리
- 조치: 중복을 제외해야 하는 대화형 쿼리의 경우, 모든 레코드에
event_id와event_time을 게시하고 중복 제거 뷰를 사용합니다: CREATE OR REPLACE VIEW analytics.latest_events AS SELECT * EXCEPT(rn) FROM ( SELECT e.*, ROW_NUMBER() OVER (PARTITION BY event_id ORDER BY event_time DESC) rn FROM analytics.events e ) WHERE rn = 1; 관찰된 BigQuery 스트리밍 가용성(예: 중앙값 지연 시간의 두 배)을 기반으로 짧은 쿼리 지연을 도입합니다. - 근거: 최소 한 번 전송은 멱등성 있는 쓰기와 쿼리 시점 중복 제거를 요구합니다. 대기는 스트리밍 가시성 지연 시간을 고려할 때 진행 중인 행으로 인한 누락을 줄여줍니다.
- Eventarc와 파트너 파일 드롭 통합
- 조치:
partner-drops버킷에 대한 Cloud Storageobject.finalized이벤트를 Cloud Run 서비스로 라우팅하도록 Eventarc를 구성합니다. 이 서비스는 CSV를 BigQuery로 로드하는 배치 Dataflow 작업을 시작하고, 파싱 오류는 데드-레터 테이블로 보냅니다. - 근거: Eventarc는 버킷 및 객체 프리픽스에 대한 CloudEvents 필터링을 통해 이벤트 기반 오케스트레이션을 제공합니다. 배치 Dataflow 작업은 올바른 데이터를 신속하게 로드하면서 분석을 위해 형식이 잘못된 행을 분리합니다.
- 플랫폼 보안
- 조치: 별개의 서비스 계정을 사용합니다: 생산자는
telemetry에 대한pubsub.publisher권한을 얻고, Dataflow 워커 SA는telemetry-stream-sub에 대한pubsub.subscriber권한과 대상 BigQuery 데이터 세트 및 Cloud Storage에 대한 쓰기 액세스 권한을 얻습니다. Eventarc 트리거는 Cloud Run에 대한invoker권한이 있는 전용 SA를 사용합니다.telemetry토픽과 BigQuery 데이터 세트에 CMEK를 활성화합니다. 푸시 엔드포인트가 있는 경우 OIDC 및 대상(audience) 검사를 구성합니다. - 근거: 최소 권한 IAM 및 CMEK는 보안 및 규정 준수 요구사항을 충족합니다. 인증된 전송은 스푸핑을 방지합니다.
- 안정적인 운영 및 확장
- 조치: 피크를 흡수할 수 있도록 넉넉한
maxWorkers를 설정하여 Dataflow 자동 확장을 구성합니다.subscription/oldest_unacked_message_age및expired_ack_deadline_count를 모니터링하고, 임계값 초과 시 알림을 설정합니다. 호환성을 깨는 파이프라인 변경의 경우, 메시지 손실을 피하기 위해 드레이닝(drain)하여 배포합니다. 지연이 증가하면 구독자 병렬성을 높이고 처리 시간에 비례하여 확인(ack) 기한을 연장합니다. - 근거: 사전 모니터링은 지연 및 재전송을 조기에 감지합니다. 자동 확장 및 조정된 확인 기한은 중복 폭풍을 방지합니다. 드레이닝은 업그레이드 중 진행 중인 메시지를 보존합니다.
이 설계는 이벤트 기반 배치 통합을 통해 복원력 있고 안전하며 관찰 가능한 실시간 수집을 제공하고, 중복 허용 및 스키마 진화를 지원하며, 대상이 지정된 수정을 위해 잘못된 데이터를 격리하면서 빠른 분석을 제공합니다.
← Dataflow 및 Apache Beam을 사용한 스트림 처리 · 모든 도메인 · Spark →
이 문제 연습하기 → · 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.
시험 합격하기 →