Google PDE: Messaging, Event Ingestion en Real-Time Services — Studiegids
Onderdeel van de Google Professional Data Engineer — Studiegids. Oefen met geverifieerde antwoorden in het Google-examencentrum, of doe getimede oefentests op ExamRoll.io.
Overzicht
Messaging, event-ingestie en real-time services op Google Cloud zijn gecentreerd rond Cloud Pub/Sub en Eventarc voor ontkoppeld, duurzaam transport; Dataflow voor stateful streamverwerking; en sinks zoals BigQuery, Cloud Storage en operationele databases. Ontwerpen voor at-least-once delivery, idempotente consumptie en observeerbaarheid zorgt voor veerkrachtige systemen die elastisch schalen en tegelijkertijd correctheid behouden bij storingen, backpressure en schema-evolutie.
Core Messaging met Pub/Sub
- Topics en abonnementen
- Publishers sturen berichten naar een topic; subscribers koppelen via abonnementen (meerdere subscribers kunnen onafhankelijk van elkaar dezelfde berichten consumeren).
- Abonnementstypes:
- Pull: clients halen expliciet berichten op; gebruik streaming pull voor de hoogste doorvoer en minder round trips.
- Push: Pub/Sub levert via HTTPS; uw endpoint moet een 2xx-statuscode retourneren ter bevestiging (acknowledge).
- Export naar BigQuery: een BigQuery-abonnement levert berichten aan een BigQuery-tabel zonder code; het beste wanneer payloads overeenkomen met het gedeclareerde schema en low-latency ingestie in analytics vereist is.
- Ordering keys
- Activeer berichtvolgorde op het topic en het abonnement om levering in de juiste volgorde per ordering key te ontvangen. De doorvoer per key wordt geserialiseerd: één in-flight bericht per key kan de volgende blokkeren; gebruik veel keys (bijvoorbeeld hash(device_id)) om te schalen.
- Fan-out en replay
- Maak afzonderlijke abonnementen voor verschillende consumers om workloads en retentie te isoleren.
- Gebruik seek of een snapshot om te replayen vanaf een timestamp of snapshot voor herstel en backfills.
Afwegingen:
- Volgordebepaling vermindert parallellisme en doorvoer per key; schakel dit uit tenzij het strikt noodzakelijk is.
- Push vereenvoudigt de clientcode, maar introduceert aandachtspunten voor schaalbaarheid van het HTTP-endpoint, beveiliging en backoff; pull geeft meer controle en stabiliteit bij hoge doorvoer.
Delivery Semantics, Acknowledgment, Retentie en Dead Lettering
- Bevestiging (acknowledgment) en deadlines
- At-least-once delivery: duplicaten kunnen voorkomen.
- Elke levering heeft een ack-deadline (standaard 10 seconden). Verleng deze (ModifyAckDeadline) tijdens het verwerken van langdurige taken; het niet bevestigen (acken) vóór de deadline is de meest voorkomende oorzaak van dubbele push-leveringen.
- Een nack of het verstrijken van de deadline maakt het bericht opnieuw beschikbaar voor herlevering.
- Retentie
- Niet-bevestigde (unacknowledged) berichten worden bewaard gedurende de ack-deadline van het abonnement en opnieuw geprobeerd; bevestigde (acknowledged) berichten kunnen worden bewaard tot de retentieduur van het topic voor replay.
- Configureer de retentieperiode om uw maximale uitval plus hersteltijd te dekken.
- Herpogingen (retries)
- Pull: herlevering vindt plaats nadat de ack-deadline is verstreken; beheer concurrency met flow control-limieten.
- Push: exponentiële backoff; alleen HTTP 2xx is een succes. 3xx/4xx/5xx triggeren herpogingen. Implementeer idempotente handlers om herhalingen te tolereren.
- Dead-letter topics (DLT’s)
- Configureer een DL-topic en een maximaal aantal leveringspogingen per abonnement om ‘poison messages’ in quarantaine te plaatsen.
- Monitor het DLQ-volume; creëer triage-workflows en publiceer na correctie opnieuw naar het hoofdtopic.
Voorbeeld:
undefined
Samenvatting delivery semantics:
- Pub/Sub: at-least-once, best-effort-volgorde binnen een ordering key indien ingeschakeld.
- Sinks: BigQuery insert-API’s bieden duplicaatmitigatie (insertId of Storage Write API stream-offsets), maar ontwerp consumers en writers desondanks idempotent.
Schema’s, Compatibiliteit en Validatie
- Pub/Sub-schema’s
- Ingebouwde ondersteuning voor Avro en Protocol Buffers met centraal opgeslagen schema’s.
- Schema-instellingen op topic-niveau: codering (Avro of Protobuf) en handhaving (geen, alleen valideren, of vereisen).
- De producer publiceert gecodeerde payloads; Pub/Sub valideert deze tegen het huidige schema wanneer handhaving is ingeschakeld.
- Evolutie en compatibiliteit
- Gebruik backward-compatible wijzigingen (voeg optionele velden toe, voeg velden met standaardwaarden toe in Avro, hergebruik nooit tags in Protobuf, vermijd het verwijderen of hernoemen van velden).
- Versioneer schema’s expliciet. Voor breaking changes, gebruik dual-publish naar v1- en v2-topics, of voeg een versienummer-veld toe en routeer op basis daarvan.
- Producer-consumer-contracten
- Consumers moeten onbekende velden negeren en ontbrekende velden voorzien van een standaardwaarde.
- Test de schemacompatibiliteit met alle consumers vóór promotie; valideer in staging-abonnementen met dezelfde schemahandhaving als in productie.
Kort Avro-voorbeeld (fragment):
undefined
Event-Driven Integratie, Eventarc en Kafka-interoperabiliteit
- Eventarc en CloudEvents
- Eventarc routeert events van Google Cloud-services (en aangepaste bronnen via Pub/Sub) naar Cloud Run, GKE of Workflows met behulp van de CloudEvents-specificatie. Attributen zoals type, source en subject maken fijnmazige filtering en auditeerbaarheid mogelijk.
- Gebruik attribuutfilters om fan-out te minimaliseren en de downstream-belasting te verminderen.
- Levering is at-least-once; maak handlers waar mogelijk idempotent en stateless.
- Voorbeeld van een Eventarc-trigger:
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-interoperabiliteit en beheerde migratie
- Dataflow-templates verbinden Kafka <-> Pub/Sub voor een gefaseerde migratie. Spiegel topics met behoud van keys; zet eerst consumers over, dan producers, of gebruik dual-write tijdens de overgang.
- Pub/Sub Lite biedt gepartitioneerde, op capaciteit ingerichte streaming met key-gebaseerde routing en lagere kosten; het is regionaal/zonaal en geschikt voor Kafka-achtige workloads waarbij voorspelbare capaciteit en volgorde per partitie de belangrijkste aandachtspunten zijn.
- Migratieoverwegingen:
- Volgorde: map Kafka-keys naar Pub/Sub ordering keys of Lite-partities.
- Offsets: neem offsets mee als berichtattributen voor diagnostiek; consumers kunnen na de migratie niet vertrouwen op Kafka-offsets.
- Levering: accepteer at-least-once; dwing idempotency af in downstream-systemen.
- Schema’s: migreer Confluent Schema Registry-definities naar Pub/Sub-schema’s of standaardiseer op Protobuf/Avro met compatibele evolutieregels.
Patronen voor Streaming Ingestie, Doorvoer, Schaalbaarheid, Beveiliging en Operations
- Real-time ingestiepatronen
- Pub/Sub -> Dataflow -> BigQuery: gebruik de BigQuery Storage Write API-sink voor hoge doorvoer en idempotentie met stream offsets; stuur mislukte berichten door naar een dead-letter-tabel voor inspectie.
- Pub/Sub -> Dataflow -> Cloud Storage: archiveer onbewerkte events voor herverwerking; gebruik windowed, gecomprimeerde writes om een balans te vinden tussen kosten en latency.
- Pub/Sub -> operationele datastores: schrijf naar Bigtable voor lookups met lage latency, naar Spanner voor sterk consistente transacties, of naar Cloud SQL/Firestore afhankelijk van de workload-eisen. Zorg voor idempotente upserts met een unieke event-ID als sleutel.
- At-least-once, preventie van duplicaten en idempotentie
- Neem een unieke event_id en event_time op in elk bericht; dwing UUID’s af aan de kant van de producer.
- BigQuery streaming ontdubbeling: stel insertId in of gebruik de Storage Write API met geordende streams; beveilig query’s desondanks met ontdubbelingslogica.
- Voorbeeld van ontdubbeling tijdens query-uitvoering: 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;
- Voor push-eindpunten, retourneer alleen een 2xx-status na succesvolle verwerking; verwacht anders een nieuwe afleverpoging.
- Berichtendoorvoer, quota en schaalbaarheid
- Publishers: batch berichten en hergebruik verbindingen; paralleliseer over meerdere clients. Gebruik veel ordering keys om geordende workloads te schalen.
- Subscribers: geef de voorkeur aan streaming pull met flow control (maximaal aantal openstaande bytes/berichten). Stem ack-deadlines af op de verwerkingstijd en verleng ze indien nodig.
- Monitor en vraag quotumverhogingen aan voor publish- en subscribe-doorvoer naarmate de volumes groeien; ontwerp met extra capaciteit (bijvoorbeeld 2x de verwachte piek) om pieken op te vangen.
- Consistentie en beschikbaarheid
- BigQuery streaming is eventually consistent voor zichtbaarheid in query’s; voor interactieve query’s die gestreamde rijen moeten bevatten, wacht op basis van de waargenomen latency (bijvoorbeeld 2x de P50-beschikbaarheidsvertraging) of ontwerp met op watermarks afgestemde aggregaties in Dataflow en bevraag de gematerialiseerde resultaten.
- Beveiliging
- IAM: ken rollen toe volgens het least-privilege-principe (pubsub.publisher aan producers op het topic; pubsub.subscriber aan consumers op de subscription). Gebruik toegewijde serviceaccounts voor elke workload.
- Push-authenticatie: configureer push-subscriptions om OIDC-tokens van een serviceaccount mee te sturen; dwing audience-validatie af op het eindpunt. Geef de voorkeur aan private Cloud Run-eindpunten voor ingebouwde authenticatie en TLS.
- Encryptie: Pub/Sub versleutelt data in transit en at rest; gebruik CMEK op topics voor door de klant beheerde sleutels. Pas VPC Service Controls toe om het risico op data-exfiltratie te verminderen. Gebruik client-side encryptie voor gevoelige payload-velden indien nodig.
- Operationele diagnose van lag, herlevering en falende subscribers
- Monitoren met Cloud Monitoring:
- subscription/num_undelivered_messages en oldest_unacked_message_age voor de backlog.
- expired_ack_deadline_count om gemiste acks te detecteren die duplicaten veroorzaken.
- publish_request_count en pull_request_count voor de doorvoer.
- Onderzoek ontbrekende dashboard-events door een bekende dataset opnieuw door de pipeline af te spelen en de output van elke fase te vergelijken om de foutieve transformatie of sink te isoleren.
- Voor Dataflow streaming:
- Gebruik autoscaling met een geschikte maxWorkers-waarde om de belasting van vele bronnen op te vangen.
- Drain pipelines bij incompatibele updates om in-flight werk te laten voltooien en dataverlies te voorkomen.
- Voor BigQuery insert-notificaties, routeer Cloud Logging audit-entries via een sink die filtert op specifieke tabellen naar een Pub/Sub-topic voor alarmering.
- Monitoren met Cloud Monitoring:
Praktijkscenario
Contoso Freight heeft een wereldwijd, real-time eventing-platform nodig om 10.000 IoT-telemetrieberichten per minuut van vrachtwagens te ingesteren, events te verrijken, interactieve analyses mogelijk te maken en workflows te triggeren bij het ontvangen van bestanden van externe partners. Sommige CSV’s van partners bevatten corrupte rijen, en het analytics-team moet fouten kunnen inspecteren zonder de stream te blokkeren.
- Creëer de kernlaag voor messaging en schema’s
- Actie: Definieer een Avro-schema voor telemetrie en koppel dit aan een Pub/Sub-topic
telemetrymet schema-validatie ingesteld op ‘vereist’. Activeer berichtvolgorde en publiceer metordering_key = hash(device_id). - Rationale: Schema-validatie op topic-niveau wijst corrupte events vroegtijdig af. Volgorde per apparaat ondersteunt geordende verwerking waar nodig, terwijl hashen de keys spreidt om de doorvoer te behouden.
- Provisioneer subscriptions met isolatie en dead-lettering
- Actie: Maak een pull-subscription
telemetry-stream-subvoor Dataflow met een dead-letter-topictelemetry-dltenmax_delivery_attempts=10. Voeg een BigQuery-subscriptiontelemetry-raw-bqtoe om onbewerkte events op te slaan in een op tijd gepartitioneerde tabel voor lineage en herhaling. - Rationale: Een DLQ isoleert ‘poison messages’ voor onderzoek. Een aparte BigQuery-subscription biedt een low-ops exportpad voor het bewaren van onbewerkte events, onafhankelijk van de verwerkingspipeline.
- Bouw een Dataflow streaming-pipeline voor verrijking en sinks
- Actie: Ingesteer vanuit
telemetry-stream-submet behulp van streaming pull met flow control. Valideer tegen het schema, verrijk met referentiedata en bereken windowed aggregaties. Schrijf naar BigQuery met de Storage Write API met een benoemde stream eninsertId = event_id; schrijf elk uur onbewerkte back-ups naar Cloud Storage; stuur ongeldige/mislukte records door naar een dead-letter BigQuery-tabel. - Rationale: De Storage Write API levert writes met hoge doorvoer en lage latency, met idempotentie via
insertId/stream offsets. Een dead-letter-tabel ondersteunt inspectie zonder de stream te blokkeren, en Cloud Storage-archieven maken herhaling mogelijk.
- Behandel duplicaten en eventual consistency in analytics
- Actie: Voor interactieve query’s die duplicaten moeten uitsluiten, publiceer
event_idenevent_timein elk record en gebruik een ontdubbelings-view: 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; Introduceer een korte vertraging in de query op basis van de waargenomen beschikbaarheid van BigQuery streaming (bijvoorbeeld tweemaal de mediane latency). - Rationale: At-least-once-levering vereist idempotente writes en ontdubbeling tijdens de query. Wachten vermindert het missen van in-flight rijen, gegeven de zichtbaarheidslatency van streaming.
- Integreer het ontvangen van partnerbestanden met Eventarc
- Actie: Configureer Eventarc om Cloud Storage
object.finalized-events voor de bucketpartner-dropste routeren naar een Cloud Run-service die een batch Dataflow-job start om CSV’s in BigQuery te laden, waarbij parseerfouten naar een dead-letter-tabel worden gestuurd. - Rationale: Eventarc biedt event-driven orkestratie met CloudEvents-filtering op bucket en object-prefix. Een batch Dataflow-job scheidt corrupte rijen voor analyse, terwijl de correcte data direct wordt geladen.
- Beveilig het platform
- Actie: Gebruik afzonderlijke serviceaccounts: producers krijgen
pubsub.publisheroptelemetry; de Dataflow worker SA krijgtpubsub.subscriberoptelemetry-stream-suben schrijftoegang tot de doel-datasets in BigQuery en Cloud Storage; de Eventarc-trigger gebruikt een toegewijde SA met deinvoker-rol op Cloud Run. Activeer CMEK op hettelemetry-topic en de BigQuery-datasets. Configureer eventuele push-eindpunten met OIDC en audience-controles. - Rationale: Least-privilege IAM en CMEK voldoen aan beveiligings- en compliance-eisen; geauthenticeerde aflevering voorkomt spoofing.
- Betrouwbaar opereren en schalen
- Actie: Stel Dataflow autoscaling in met een ruime
maxWorkers-waarde om pieken op te vangen. Monitorsubscription/oldest_unacked_message_ageenexpired_ack_deadline_count; alarmeer wanneer drempelwaarden worden overschreden. Voor pipeline-wijzigingen die de compatibiliteit verbreken, deploy met de ‘drain’-optie om berichtverlies te voorkomen. Als de lag toeneemt, verhoog dan het parallellisme van de subscriber en verleng de ack-deadlines proportioneel aan de verwerkingstijd. - Rationale: Proactieve monitoring detecteert lag en herleveringen vroegtijdig. Autoscaling en afgestemde ack-deadlines voorkomen ‘duplicate storms’. Draining behoudt in-flight berichten tijdens upgrades.
Dit ontwerp levert een veerkrachtige, veilige en observeerbare real-time ingestie met event-driven batch-integratie, ondersteunt tolerantie voor duplicaten en schema-evolutie, en biedt snelle analyses terwijl incorrecte data wordt geïsoleerd voor gerichte correctie.
← Streamverwerking met Dataflow en Apache Beam · Alle domeinen · Spark →
Oefen deze vragen → · Getimede oefening op 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.
Slaag voor je examen →