Google PDE: Workflow-orkestratie en Pipeline-automatisering — Studiegids
Onderdeel van de Google Professional Data Engineer — Studiegids. Oefen met geverifieerde antwoorden in het Google-examencentrum, of doe getimede oefentests op ExamRoll.io.
Betrouwbaarheid, Foutafhandeling en Idempotentie
Retries, timeouts en backoff
- Gebruik begrensde exponentiële backoff voor tijdelijke fouten en beperk de totale retry-vensters tot de SLA van de taak. Een frontend of taak die bijvoorbeeld elke 15 minuten een database pollt, moet met exponentiële backoff tot 15 minuten proberen te herhalen en daarna een gecontroleerde fout tonen.
- Configureer
execution_timeoutper taak en globale DAG SLA’s in Airflow; stel in Workflows timeouts per stap en retry-beleid in metmax_doublingsenmax_retry_duration. Voor Cloud Run-jobs, stel het aantal retries en de backoff in.
Backfills, catchup en foutafhandeling
- Activeer catchup voor historische herberekeningen wanneer taken idempotent zijn en bronnen op datum zijn gepartitioneerd. Voor niet-deterministische outputs of externe neveneffecten, overweeg DAG’s die alleen voor backfill zijn of write-audit tabellen om bij te houden wat er is geproduceerd.
- Gebruik dead-letter topics/tabellen voor fouten op recordniveau in streaming/batch-transformaties. Voor batch Dataflow, vang foute rijen op met fouttags en aggregeer foutstatistieken; voor streaming, gebruik Pub/Sub DLQ’s.
Idempotent taakontwerp en heruitvoeringen
- BigQuery: geef de voorkeur aan MERGE of INSERT met ontdubbelingssleutels; gebruik
insertIdom streaming inserts te ontdubbelen. Voor batch, schrijf naar een staging-tabel en voer vervolgens een MERGE uit naar de doeltabel binnen een transactioneel veilige stap om volledige heruitvoeringen mogelijk te maken. - Cloud Storage: gebruik generatie-voorwaarden en deterministische objectnamen (bijv. prefix/datum/hash) zodat heruitvoeringen alleen veilig overschrijven wanneer dit verwacht wordt.
- Pub/Sub en Dataflow: ontwerp voor at-least-once levering. Voeg bericht-ID’s toe (bijv. Package ID, logische event-timestamp) zodat downstream systemen kunnen ontdubbelen en redeneren over vertraging. Als de bedrijfsregels “eerste verwerkte event wint”-semantiek accepteren, documenteer die afweging en monitor op scheeftrekking; anders, bepaal de winnaars op basis van event-tijd met tie-breakers.
- Herstel van gedeeltelijke fouten: partitioneer outputs op
run_idof datum, schrijf voltooiingsmarkeringen en maak downstream taken afhankelijk van deze markeringen. Verwerk alleen de partities opnieuw die als onvolledig zijn gemarkeerd.
Probleemoplossing en schaalbaarheid
- Wanneer een streaming dashboard events mist maar Pub/Sub aangeeft dat ze aanwezig zijn, voer dan een bekende, vaste dataset door de Dataflow-pipeline om transformatiedefecten te isoleren. Valideer de windowing, triggers en
allowed lateness. - Veelvoorkomende foutmodus: het creëren van een streaming pipeline zonder geschikte windowing/triggers voor onbegrensde bronnen of het incorrect gebruiken van een gesharde window kan het aanmaken van de pipeline laten mislukken of state-explosies veroorzaken.
- Schaal Dataflow via
max workersen het autoscaling-algoritme; voor pieken (bijv. 50.000 installaties), verhoog het maximale aantal workers om horizontale schaalvergroting tijdens piekmomenten mogelijk te maken.
Beveiliging, Parametrisering, Omgevingen en CI/CD
Parametrisering en configuratiebeheer
- Externaliseer configuratie per omgeving. Gebruik in Composer Variables, Connections en omgevingsvariabelen; maak DAG-parameters template-gebaseerd op uitvoeringsdatum of partitie. Gebruik in Workflows runtime-argumenten en aparte workflows per omgeving, of lees de configuratie uit Secret Manager.
- Gebruik metadata-gestuurde orkestratie door een controletabel te lezen (bijv. een BigQuery config-dataset) die clients, bronnen of partities opsomt. Genereer taken dynamisch zodat codewijzigingen losgekoppeld zijn van datagestuurde wijzigingen.
Secrets, service accounts en least privilege
- Sla credentials op in Secret Manager en refereer ernaar tijdens runtime. Vermijd het inbedden van secrets in code of Airflow Variables.
- Wijs een aparte service account toe per pipeline met de minimaal benodigde IAM-rollen. Voor gereguleerde BigQuery-toegang, isoleer klantdata in aparte datasets, ken dataset-specifieke rollen alleen toe aan goedgekeurde gebruikers en beperk de toegang tot de BigQuery API tot goedgekeurde principals. Voor multitenancy, maak een dataset per klant aan en koppel alleen de juiste rollen.
CI/CD en infrastructure as code
- Beheer infrastructuur (Composer-omgevingen, Workflows, Scheduler-jobs, Pub/Sub-topics, log sinks) met Terraform. Gebruik modules om projecten/omgevingen, secrets en service accounts te standaardiseren.
- Bouw en test pipeline-code met Cloud Build of GitHub Actions. Automatiseer unit tests, SQL linting, Dataform dry-runs en Airflow DAG-validatie. Promoot artefacten via tags; voor Composer, verpak DAG’s als deployable bundels; voor Dataform, gebruik release branches die promoten nadat assertions slagen.
- Deployment-promotie: dev → test → prod via aparte projecten en geparametriseerde configuraties. Gebruik continuous delivery met handmatige goedkeuringspoorten en change windows voor promoties met een hoog risico.
Observeerbaarheid, Alarmering en Draaiboeken
Telemetrie en alarmering
- Stuur alle orkestratielogs naar Cloud Logging met gestructureerde velden (pipeline, dag_id, run_id, task_id, partition). Exporteer foutenlogs naar Monitoring via op logs gebaseerde metrics. Alarmeer bij:
- Gemiste schedules of SLA-overschrijdingen
- Opeenvolgende taakfouten
- Groei van de backlog (bijv. Pub/Sub unacked-berichten, Dataflow-systeemvertraging)
- Fouten bij data-kwaliteitsasserties
- Cloud Composer: monitor de duur van DAG’s/taken, slagingspercentage, wachtrijdiepte en de gezondheid van de scheduler. Configureer on_failure_callback voor paging en hersteldraaiboeken.
- Cloud Workflows: inspecteer Execution-logs en stap-latencies; voeg expliciete retries en error handlers toe; stuur aangepaste logs uit met correlatie-ID’s.
- BigQuery-tabelwijzigingsnotificaties: maak een Logging-sink op projectniveau met een geavanceerd filter voor insert-jobs gericht op een specifieke tabel en exporteer naar Pub/Sub; uw monitoringtool abonneert zich op het topic voor directe meldingen zonder ruis van andere tabellen.
Ontwerp van draaiboeken
- Documenteer voor elke pipeline de triggers, afhankelijkheden, SLA’s, rollback/retry-procedures en veilige backfill-stappen. Neem “fixed dataset replay” voor Dataflow op, hoe een streaming-job te ‘drainen’, hoe mislukte partities opnieuw te verwerken en hoe DLQ-berichten te herstellen.
- Leg veelvoorkomende foutsymptomen vast (bijv. permission denied, quota exceeded, schema mismatch) met beslisbomen en escalatiepaden.
Praktisch Probleemscenario
Acme Retail Analytics moet dagelijkse CSV-leveringen van partners verwerken die af en toe corrupte rijen bevatten, de geldige data transformeren en laden naar BigQuery, en de foute rijen beschikbaar maken voor onderzoek. Ze willen ook event-driven verrijking voor bijna-realtime prijsupdates en een veilige promotie van dev naar prod.
Aanpak:
Opslag en event-triggers
- Maak een dedicated Cloud Storage-bucket met object-versiebeheer en uniforme toegang op bucket-niveau. Schakel ‘object finalize’-notificaties naar Pub/Sub in via Eventarc.
- Rationale: ‘Object finalization’ is een betrouwbare event om downstream-ingest te triggeren; versiebeheer ondersteunt herhalingen en audits.
Batch-ingest met dead-letter-afhandeling
- Gebruik Cloud Composer om dagelijks om 02:00 een Airflow DAG te schedulen met ‘catchup’ ingeschakeld. De DAG start een Dataflow-batchjob die CSV’s parset, het schema valideert en geldige records naar BigQuery schrijft met behulp van deterministische staging-tabellen, en vervolgens MERGE’t naar gepartitioneerde doeltabellen. Routeer corrupte/mislukte records naar een BigQuery dead-letter-tabel.
- Rationale: Dataflow schaalt het parsen/valideren; MERGE zorgt voor idempotent gedrag; het vastleggen in een dead-letter-tabel maakt inspectie mogelijk zonder de pipeline te blokkeren, wat overeenkomt met het aanbevolen patroon voor corrupte rijen.
Event-driven verrijking
- Implementeer een Cloud Run-job voor lichtgewicht verrijking voor incrementele prijsupdates. Trigger deze via Cloud Workflows die luistert naar Pub/Sub-berichten van Eventarc wanneer gedurende de dag kleine update-bestanden binnenkomen.
- Rationale: Serverless containers met Workflows bieden low-latency, low-ops orkestratie voor kleine events, terwijl zware transformaties in batch blijven.
Betrouwbaarheidsmaatregelen
- Configureer retries met exponential backoff voor tijdelijke fouten in Dataflow- en Cloud Run-jobs, waarbij de totale retry-tijd wordt beperkt tot de SLA van de DAG. Stel per-taak execution timeouts en on_failure callbacks in in Airflow; stel in Workflows max_doublings en max_retry_duration in.
- Rationale: Begrensde backoff beschermt SLA’s en voorkomt ongecontroleerde retries.
Beveiliging en ’least privilege’
- Voer elk component uit onder een dedicated service account: Composer orchestrator SA, Dataflow worker SA, Cloud Run job SA. Wijs alleen de benodigde rollen toe: GCS-leesrechten op de ingest-bucket voor Dataflow, BigQuery dataEditor op de doel-datasets en Viewer op de logs. Sla secrets op in Secret Manager en verwijs ernaar tijdens runtime.
- Rationale: Dwingt ’least privilege’ af en isoleert de ‘blast radius’.
Metadata-gedreven orkestratie
- Onderhoud een BigQuery-controletabel met partnerbronnen, bestandspatronen en doel-datasets. Tijdens de runtime van de DAG bevraagt Airflow deze tabel en gebruikt ‘dynamic task mapping’ om per partner taken te genereren.
- Rationale: Het toevoegen van een partner wordt een datawijziging in plaats van een codewijziging, wat het implementatierisico verkleint.
Observeerbaarheid en alarmering
- Stuur gestructureerde logs uit met run_id en partner_id. Maak alarmeringsbeleid voor SLA-overschrijdingen van de DAG, Dataflow-systeemvertraging en niet-lege dead-letter-tellingen. Configureer voor BigQuery-inserts in de doeltabel een Cloud Logging-sink met een geavanceerd filter voor die tabel naar een Pub/Sub-topic dat wordt geconsumeerd door de monitoringtool van Acme.
- Rationale: Fijnmazige alarmering maakt snelle triage mogelijk zonder ruis.
CI/CD en promotie
- Beheer de infrastructuur (buckets, Pub/Sub, Eventarc, Composer, Workflows, BigQuery-datasets) in Terraform. Gebruik Cloud Build om de syntaxis van Airflow DAG’s te valideren, unit tests uit te voeren en te implementeren in een dev Composer-omgeving. Promoot naar test en prod met geparametriseerde configuraties en handmatige goedkeuringsstappen nadat Dataform-asserties en integratietests zijn geslaagd.
- Rationale: Declaratieve, herhaalbare implementaties en veilige promotie tussen omgevingen.
Draaiboek en herstel
- Documenteer de stappen om een specifieke datum opnieuw af te spelen: herstel de CSV vanuit object-versiebeheer, voer de Dataflow-job voor die partitie opnieuw uit, MERGE de resultaten en controleer de DLQ-records. Neem een “fixed dataset replay”-procedure op om transformatiebugs te isoleren als er discrepanties optreden.
- Rationale: Een idempotent ontwerp en gedocumenteerd herstel stroomlijnen het oplossen van gedeeltelijke fouten.
← Data-ingestie · Alle domeinen · Machine Learning →
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 →