Google PDE: Streamverwerking met Dataflow en Apache Beam — 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

Streamverwerking op Google Cloud draait om het uniforme programmeermodel van Apache Beam, uitgevoerd door de Dataflow runner. Beam biedt een logische abstractie—pipelines van transforms over PCollections—die uw code loskoppelt van uitvoeringsdetails zoals parallellisme, autoscaling en fouttolerantie. Bij streaming hangt de correctheid af van tijdsemantiek (event time vs. processing time), windowing (fixed, sliding, session, global), watermarks, triggers en de verwerking van late data. Operationele uitmuntendheid op Dataflow vereist de juiste worker-grootte, autoscaling-beleid, streaming engine, shuffle-keuzes, een idempotent sink-ontwerp, dead-letter-verwerking en robuuste observability.

Apache Beam-model en tijdsemantiek

Faalscenario’s en afwegingen:

Dataflow beheren voor streaming workloads

Implementatie, templates en upgradestrategieën

undefined

undefined

Praktisch probleemscenario

NovaTrack Inc. verwerkt wereldwijde IoT-telemetrie van 50.000 temperatuursensoren en moet aggregaties op minuutniveau leveren, ruwe data persisteren en een real-time dashboard voorzien. Af en toe worden misvormde berichten en ‘out-of-order’ levering verwacht. De oplossing moet automatisch schalen, foute records tonen voor inspectie en upgrades zonder downtime ondersteunen.

Aanpak:

  1. Ingestie en tijdsemantiek

    • Maak een regionaal Pub/Sub-topic en per-regio publishers aan met de attributen deviceId en eventTs (RFC3339). Schakel waar mogelijk ‘ordering keys’ in op basis van deviceId.
    • Rationale: Pub/Sub biedt duurzame, elastische ingress met ‘at-least-once’ levering. Het toevoegen van event-timestamps aan de ’edge’ behoudt de ware event-tijd; ordening per apparaat vermindert herschikking binnen een apparaat zonder centrale knelpunten.
  2. Dataflow streaming-pipeline met event-time windows

    • Lees vanuit een toegewijde subscription via PubSubIO, extraheer eventTs als de Beam-timestamp en val terug op publishTime indien deze ontbreekt.
    • Pas FixedWindows van 1 minuut toe met een vroege trigger na 30 seconden en ’late firings’ voor elk te laat element; stel de toegestane vertraging (‘allowed lateness’) in op 10 minuten en accumuleer ‘panes’.
    • Rationale: Event-time windows zorgen voor nauwkeurige aggregaties per minuut; vroege ‘firings’ voeden het dashboard met een actualiteit van minder dan een minuut; late ‘firings’ corrigeren aggregaties naarmate vertraagde data binnenkomt. De begrenzing op vertraging beperkt de omvang van de state en de kosten.
  3. Validatie, verrijking en dead-letter-routering

    • Implementeer een ParDo die JSON parset, schema en bereiken valideert, en verrijkt met kleine statische referentiedata via een ‘side input’ die bij het starten van de job uit BigQuery wordt geladen.
    • Gebruik TupleTags om valide records naar de hoofdoutput te sturen en fouten naar een dead-letter PCollection die de payload, foutmelding, deviceId en parse-timestamp bevat; schrijf de DLQ naar een gepartitioneerde BigQuery-tabel.
    • Rationale: ‘Side inputs’ houden referentiedata in het geheugen voor lage latentie. Het vastleggen in een dead-letter-queue maakt inspectie en gerichte herverwerking van foute rijen mogelijk zonder de hoofdstroom te blokkeren.
  4. Aggregatie en hot-key-mitigatie

    • Gebruik deviceId als sleutel en bereken per minuut avg/min/max met CombineFns. Voor top-N regionale statistieken, shard op region#N om ‘hot keys’ te vermijden en aggregeer daarna opnieuw.
    • Rationale: Combiners minimaliseren het shuffle-volume en de kosten; ‘key sharding’ voorkomt knelpunten bij één enkele sleutel tijdens de regionale ‘fan-in’.
  5. Sinks en ’exactly-once’-effecten

    • Schrijf ruwe gevalideerde events en aggregaties per minuut naar BigQuery met BigQueryIO en de Storage Write API. Stel een stabiele ‘insert id’ in op basis van deviceId + eventTs voor idempotentie bij eventuele aangepaste ‘retries’.
    • Rationale: De Storage Write API biedt high-throughput, low-latency ingestie met ’exactly-once’-semantiek binnen een stream. Stabiele id’s zorgen voor downstream ontdubbeling als ‘replays’ optreden.
  6. Strategie voor dashboardconsistentie

    • Het dashboard bevraagt gepartitioneerde aggregatietabellen met een terugkijkperiode van 2 minuten ten opzichte van de watermark of een vaste vertraging van 2x de waargenomen beschikbaarheidslatentie voor streaming data.
    • Rationale: De zichtbaarheid van streaming data in BigQuery is ’eventually consistent’; het iets uitstellen van leesoperaties voorkomt het missen van rijen die onderweg zijn, terwijl het near-real-time gedrag behouden blijft.
  7. Operations: autoscaling en streaming engine

    • Schakel Streaming Engine in; stel maxWorkers in op basis van de verwachte piek (bijv. 3x het gemiddelde), selecteer een machinetype dat is gedimensioneerd voor CPU-gebonden parsing en encryptie, en vergroot de ‘boot disk’ om ruimte te bieden voor tijdelijke shuffle.
    • Monitor de watermark-vertraging, backlog in seconden, CPU en doorvoer per stap; stel alerts in voor aanhoudende vertraging en pieken in de DLQ-rate.
    • Rationale: Streaming Engine externaliseert state/shuffle voor elasticiteit en eenvoudigere upgrades; juiste dimensionering en monitoring voorkomen stille schendingen van de SLO.
  8. Implementatie en upgrades met Flex Templates

    • Verpak de pipeline als een Flex Template met parameters: input-subscription, output-tabellen, DLQ-tabel, maxWorkers en regio. Start voor een incompatibele wijziging de nieuwe pipeline die gericht is op hetzelfde topic met een nieuwe subscription, verifieer de outputs en ‘drain’ vervolgens de oude job. Maak optioneel een Pub/Sub-snapshot aan en laat de nieuwe subscription vanaf het snapshot beginnen om te garanderen dat er geen hiaten zijn.
    • Rationale: Flex Templates maken herhaalbare, geparametriseerde implementaties mogelijk. Een geverifieerde blue/green-overschakeling met ‘drain’ zorgt voor nul dataverlies en minimale downtime.
  9. Herverwerking en batch-backfills

    • Sla gecomprimeerde Avro-bestanden van ruwe events op in Cloud Storage via een ‘side output’; voer een batch Dataflow-pipeline uit om te ‘backfillen’ of opnieuw te verwerken naar BigQuery wanneer modellen of schema’s veranderen.
    • Rationale: Duurzame ruwe archieven ondersteunen reproduceerbaarheid en schema-evolutie zonder het ‘hot path’ te beïnvloeden.

Dit ontwerp levert correcte aggregaties met lage latentie en begrensde kosten, duidelijke foutisolatie, sterke observeerbaarheid en veilige upgradepaden, terwijl het omgaat met out-of-order en late data op wereldwijde schaal.


BigQuery Analytics en Warehouse Engineering · Alle domeinen · Messaging

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 →

Blader door Google →

Related guides

Alles-in-één toegang

Eén abonnement. Elk examen.

Elk plan ontgrendelt onbeperkt zoeken naar antwoorden, oefentests, AI-uitleg en de volledige bronnenbibliotheek — in meer dan 20 talen.

Maandelijks
24.87
Just €0.83/day
Alles inbegrepen:
  • Onbeperkt zoeken naar antwoorden
  • Onbeperkte oefentests
  • AI-gestuurde uitleg
  • Volledige bronnenbibliotheek
  • 20+ talen
  • Wekelijkse contentupdates
  • Beloningen & verwijzingen
  • Prioriteitsondersteuning
Start gratis proefperiode

Geen creditcard vereist*

Beste waarde
12 maanden
179.87
Just €0.49/daySave 40%
Alles inbegrepen:
  • Onbeperkt zoeken naar antwoorden
  • Onbeperkte oefentests
  • AI-gestuurde uitleg
  • Volledige bronnenbibliotheek
  • 20+ talen
  • Wekelijkse contentupdates
  • Beloningen & verwijzingen
  • Prioriteitsondersteuning
Start gratis proefperiode

Geen creditcard vereist*

✓ Gratis plan inbegrepen · ✓ Annuleer op elk moment · ✓ Alle plannen ontgrendelen het volledige product