Google PDE: Orkiestracja przepływów pracy i automatyzacja potoków — Przewodnik do nauki
Część Google Professional Data Engineer — Przewodnik do nauki. Ćwicz ze zweryfikowanymi odpowiedziami w centrum egzaminów Google, albo rozwiąż testy na czas na ExamRoll.io.
Przegląd
Orkiestracja przepływów pracy i automatyzacja potoków koordynują zadania przetwarzania danych między usługami, tak aby pozyskiwanie, transformacja, kontrole jakości i publikowanie odbywały się w sposób niezawodny, bezpieczny i opłacalny. W Google Cloud orkiestracja musi być zgodna z modelem wykonania każdego obciążenia: zaplanowane zadania wsadowe (batch), strumienie sterowane zdarzeniami, zadania ad-hoc lub zadania długotrwałe. Cele projektowe to powtarzalność, idempotentność, obserwowalność, zasada najmniejszych uprawnień i bezpieczne promowanie między środowiskami.
Kluczowe wybory:
- Orkiestracja wsadowa zorientowana na kod za pomocą Cloud Composer (Apache Airflow) dla DAG-ów, zależności między zadaniami i zaawansowanego harmonogramowania.
- Bezserwerowa choreografia API za pomocą Cloud Workflows dla lekkich, sterowanych zdarzeniami sekwencji obejmujących wiele usług.
- Punkty końcowe wykonania, takie jak zadania Cloud Run lub Dataproc, wyzwalane przez Cloud Scheduler dla zadań typu cron lub przez Eventarc dla zdarzeń.
- Orkiestracja natywna dla SQL za pomocą Dataform dla transformacji w BigQuery, asercji i zarządzania wydaniami.
Model operacyjny kładzie nacisk na ponowienia prób z ograniczonym wykładniczym czasem oczekiwania (exponential backoff), limity czasu (timeouts), umowy SLA, nadrabianie zaległości (catchup) i uzupełnianie danych (backfills), projektowanie zadań w sposób idempotentny dla bezpiecznego ponownego uruchamiania oraz solidną obsługę błędów z przechwytywaniem w kolejce niedoręczonych wiadomości (dead-letter). Bezpieczeństwo jest egzekwowane poprzez konta serwisowe dedykowane dla każdego potoku, izolację sekretów, parametryzację i zasadę najmniejszych uprawnień w IAM. CI/CD, infrastruktura jako kod i kompleksowa telemetria dopełniają podejście gotowe do wdrożenia produkcyjnego.
Orkiestracja w Google Cloud: Narzędzia i Wzorce
Cloud Composer (Airflow)
- DAG-i definiują skierowane grafy acykliczne (DAG) wykonania z jawnymi zależnościami. Użyj TaskFlow API lub operatorów (np. BigQuery, Dataflow, Dataproc, Cloud Run) do wyrażania zadań. Sensory i operatory odraczalne (deferrable operators) zmniejszają obciążenie harmonogramu w przypadku warunków oczekiwania (np. finalizacja obiektu w Cloud Storage lub pojawienie się partycji w BigQuery).
- Harmonogramowanie: wyrażenia cron,
start_date,end_dateicatchupkontrolują historyczne uruchomienia. Użyjcatchupdo uzupełniania danych (backfills); wyłącz dla celów powiązanych ze strumieniowaniem lub nieidempotentnych. Ogranicz współbieżność za pomocąmax_active_runsi pul (pools), aby chronić systemy podrzędne. - Zależności:
set_upstream/set_downstreamlub zależności w taskflow. W przypadku orkiestracji sterowanej metadanymi, dynamicznie generuj zadania na podstawie tabeli kontrolnej w BigQuery (np. lista klientów/partycji) używając dynamicznego mapowania zadań (dynamic task mapping), utrzymując stabilny czas parsowania DAG-a i sprawiając, że zadania są sterowane danymi. - Przykład (skrócony) fragmentu DAG-a:
undefined
undefined
undefined
undefined
undefined
undefined
undefined
undefined
Cloud Workflows, Cloud Scheduler, zadania Cloud Run i wykonanie sterowane zdarzeniami
- Cloud Workflows orkiestruje interfejsy API Google i punkty końcowe HTTP z wbudowanymi ponowieniami prób, pętlami, gałęziami równoległymi i logiką kompensacyjną. Jest idealny do lekkiego przepływu sterowania między usługami takimi jak BigQuery, Dataflow, Batch i zadania Cloud Run.
- Cloud Scheduler wyzwala Workflows, tematy Pub/Sub lub usługi HTTP w celu automatyzacji w stylu cron. Dla codziennego zadania wsadowego o 02:00, zaplanuj Workflow, który uruchamia zadanie Dataflow lub Dataproc.
- Zadania Cloud Run wykonują skonteneryzowane kroki wsadowe z automatycznym ponawianiem prób i minimalnym nakładem operacyjnym. Dobrze komponują się z Workflows w wieloetapowych zadaniach przetwarzania danych lub w pre/post-processingu wokół Dataflow lub BigQuery.
- Sterowanie zdarzeniami: użyj Eventarc, aby kierować zdarzenia finalizacji obiektu Cloud Storage, wiadomości Pub/Sub lub dzienniki audytu (Audit Logs) do Cloud Run lub Workflows. Dla powiadomień o zadaniach wstawiania do pojedynczej tabeli BigQuery, utwórz ujście (sink) Cloud Logging z zaawansowanym filtrem do Pub/Sub, a następnie wyzwalaj swojego konsumenta z tego tematu.
Dataform: przepływy pracy SQL dla BigQuery
- Modeluj grafy zależności za pomocą
ref(), definiuj tabele/widoki/przyrosty i orkiestruj procesy budowania według tagów lub harmonogramów. Dataform kompiluje SQLX do uporządkowanych planów wykonania, umożliwiając orkiestrację sterowaną metadanymi na podstawie deklaratywnych definicji. - Asercje zapewniają jakość danych. Asercja to zapytanie, które musi zwrócić zero wierszy, aby zostało zaliczone. Przykład asercji: – definitions/assert_non_negative_prices.sqlx
undefined
undefined
- Wydania i kontrola repozytorium: przechowuj kod w repozytorium, używaj gałęzi i przeglądów kodu (reviews) oraz promuj otagowane wydania do środowisk (np. dev, test, prod) ze zmiennymi specyficznymi dla danego środowiska. Warunkuj wdrożenia poprzez sprawdzenia CI/CD i wyniki asercji.
Dataproc, Dataflow i wzorce przechowywania danych
- Aby ponownie wykorzystać kod Hadoop/Spark przy minimalnym nakładzie operacyjnym, użyj Dataproc z konektorem GCS, aby przechowywać dane dłużej niż cykl życia klastra i zminimalizować koszty dysków trwałych. Twórz efemeryczne klastry dla każdego zadania w celu izolacji i kontroli kosztów; orkiestruj za pomocą Composer lub Workflows.
- W przypadku pozyskiwania wsadowego z nieprawidłowo sformatowanymi wierszami, uruchom Dataflow, aby zapisać prawidłowe rekordy do BigQuery, a błędy parsowania/walidacji skierować do tabeli BigQuery typu dead-letter w celu ich analizy.
Niezawodność, obsługa błędów i idempotencja
Ponowienia, limity czasu i wycofywanie wykładnicze (backoff)
- Używaj ograniczonego wycofywania wykładniczego (bounded exponential backoff) dla błędów przejściowych i ograniczaj całkowite okna ponowień do SLA zadania. Na przykład, frontend lub zadanie, które odpytuje bazę danych co 15 minut, powinno ponawiać próby z wycofywaniem wykładniczym do 15 minut, a następnie zgłosić kontrolowany błąd.
- W Airflow konfiguruj
execution_timeoutdla poszczególnych zadań i globalne SLA dla DAG-ów; w Workflows ustawiaj limity czasu dla poszczególnych kroków i polityki ponowień zmax_doublingsimax_retry_duration. Dla zadań Cloud Run ustaw liczbę ponowień i wycofywanie (backoff).
Uzupełnianie danych historycznych (backfill), nadrabianie (catchup) i obsługa błędów
- Włącz nadrabianie (catchup) w celu ponownego przeliczania danych historycznych, gdy zadania są idempotentne, a źródła partycjonowane według daty. W przypadku niedeterministycznych wyników lub zewnętrznych efektów ubocznych, rozważ użycie DAG-ów przeznaczonych tylko do uzupełniania danych (backfill) lub tabel audytowych zapisu, aby śledzić, co zostało wyprodukowane.
- Używaj tematów/tabel martwych listów (dead-letter) dla błędów na poziomie rekordu w transformacjach strumieniowych/wsadowych. W przypadku wsadowego Dataflow, przechwytuj nieprawidłowe wiersze za pomocą tagów błędów i agreguj metryki błędów; w przypadku strumieniowania, używaj DLQ w Pub/Sub.
Projektowanie zadań idempotentnych i ponowne uruchomienia
- BigQuery: preferuj
MERGElubINSERTz kluczami do deduplikacji; używajinsertIddo deduplikacji wstawień strumieniowych. W przypadku przetwarzania wsadowego, zapisuj dane do tabeli przejściowej (staging), a następnie wykonajMERGEdo tabeli docelowej w ramach kroku bezpiecznego transakcyjnie, aby umożliwić pełne ponowne uruchomienia. - Cloud Storage: używaj warunków wstępnych generacji (generation preconditions) i deterministycznych nazw obiektów (np. prefix/data/hash), aby ponowne uruchomienia bezpiecznie nadpisywały dane tylko wtedy, gdy jest to oczekiwane.
- Pub/Sub i Dataflow: projektuj z myślą o dostarczaniu co najmniej raz (at-least-once). Dołączaj identyfikatory wiadomości (np. ID paczki, logiczny znacznik czasu zdarzenia), aby systemy docelowe mogły deduplikować dane i analizować opóźnienia. Jeśli reguły biznesowe akceptują semantykę „pierwsze przetworzone zdarzenie wygrywa”, udokumentuj ten kompromis i monitoruj rozbieżności; w przeciwnym razie rozstrzygaj zwycięzców na podstawie czasu zdarzenia z użyciem reguł rozstrzygania remisów.
- Odzyskiwanie po częściowej awarii: partycjonuj wyniki według
run_idlub daty, zapisuj znaczniki ukończenia i uzależniaj zadania w dalszych etapach od tych znaczników. Przetwarzaj ponownie tylko te partycje, które są oznaczone jako nieukończone.
Rozwiązywanie problemów i skalowalność
- Gdy w panelu strumieniowym brakuje zdarzeń, ale Pub/Sub pokazuje, że są one obecne, uruchom znany, stały zbiór danych przez potok Dataflow, aby wyizolować błędy w transformacji. Zweryfikuj okienkowanie (windowing), wyzwalacze (triggers) i dozwolone opóźnienie (allowed lateness).
- Częsty tryb awarii: tworzenie potoku strumieniowego bez odpowiedniego okienkowania/wyzwalaczy dla źródeł nieograniczonych lub nieprawidłowe użycie okna podzielonego na fragmenty (sharded window) może spowodować niepowodzenie tworzenia potoku lub eksplozję stanu.
- Skaluj Dataflow za pomocą
max workersi algorytmu autoskalowania; w przypadku skoków obciążenia (np. 50 000 instalacji), podnieś maksymalną liczbę workerów, aby umożliwić skalowanie horyzontalne w okresach szczytowych.
Bezpieczeństwo, parametryzacja, środowiska i CI/CD
Parametryzacja i zarządzanie konfiguracją
- Eksternalizuj konfigurację według środowiska. W Composer używaj zmiennych (Variables), połączeń (Connections) i zmiennych środowiskowych; twórz szablony parametrów DAG według daty wykonania lub partycji. W Workflows używaj argumentów czasu wykonania i oddzielnych przepływów pracy dla każdego środowiska lub odczytuj konfigurację z Secret Manager.
- Stosuj orkiestrację sterowaną metadanymi, odczytując tabelę kontrolną (np. zbiór danych konfiguracyjnych w BigQuery), która zawiera listę klientów, źródeł lub partycji. Generuj zadania dynamicznie, aby zmiany w kodzie były oddzielone od zmian sterowanych danymi.
Sekrety, konta usług i zasada najmniejszych uprawnień
- Przechowuj poświadczenia w Secret Manager i odwołuj się do nich w czasie wykonania. Unikaj osadzania sekretów w kodzie lub w zmiennych (Variables) Airflow.
- Przypisz odrębne konto usługi do każdego potoku z minimalnymi wymaganymi rolami IAM. W przypadku regulowanego dostępu do BigQuery, izoluj dane klientów w oddzielnych zbiorach danych, przyznawaj role specyficzne dla zbioru danych tylko zatwierdzonym użytkownikom i ograniczaj dostęp do BigQuery API do zatwierdzonych jednostek (principals). W przypadku wielodostępności (multitenancy), utwórz zbiór danych dla każdego klienta i przypisuj tylko odpowiednie role.
CI/CD i infrastruktura jako kod
- Zarządzaj infrastrukturą (środowiska Composer, Workflows, zadania Scheduler, tematy Pub/Sub, ujścia logów) za pomocą Terraform. Używaj modułów do standaryzacji projektów/środowisk, sekretów i kont usług.
- Buduj i testuj kod potoku za pomocą Cloud Build lub GitHub Actions. Automatyzuj testy jednostkowe, linting SQL, uruchomienia na sucho (dry-runs) w Dataform oraz walidację DAG-ów Airflow. Promuj artefakty za pomocą tagów; dla Composer, pakuj DAG-i jako pakiety wdrożeniowe; dla Dataform, używaj gałęzi wydań (release branches), które promują zmiany po pomyślnym przejściu asercji.
- Promocja wdrożeń: dev → test → prod poprzez oddzielne projekty i sparametryzowane konfiguracje. Stosuj ciągłe dostarczanie (continuous delivery) z bramkami ręcznego zatwierdzania i oknami zmian dla promocji wysokiego ryzyka.
Obserwowalność, alerty i runbooki
Telemetria i alerty
- Przekieruj wszystkie logi orkiestracji do Cloud Logging z ustrukturyzowanymi polami (pipeline, dag_id, run_id, task_id, partition). Eksportuj logi błędów do Monitoring za pomocą metryk opartych na logach. Ustaw alerty dla:
- Pominiętych harmonogramów lub naruszeń SLA
- Kolejnych niepowodzeń zadań
- Wzrostu zaległości (np. wiadomości bez potwierdzenia w Pub/Sub, opóźnienie systemowe w Dataflow)
- Niepowodzeń asercji jakości danych
- Cloud Composer: monitoruj czas trwania DAG-a/zadania, wskaźnik powodzenia, głębokość kolejki i kondycję harmonogramu (scheduler). Skonfiguruj
on_failure_callbackdo powiadamiania osób dyżurujących (paging) i uruchamiania runbooków naprawczych. - Cloud Workflows: sprawdzaj logi wykonania (Execution logs) i opóźnienia kroków; dodawaj jawne ponowienia i procedury obsługi błędów; emituj niestandardowe logi z identyfikatorami korelacji.
- Powiadomienia o zmianach w tabeli BigQuery: utwórz ujście (sink) na poziomie projektu w Cloud Logging z zaawansowanym filtrem dla zadań wstawiania (insert jobs) do określonej tabeli i eksportuj do Pub/Sub; Twoje narzędzie monitorujące subskrybuje ten temat, aby otrzymywać natychmiastowe alerty bez szumu z innych tabel.
Projektowanie runbooków
- Dla każdego potoku udokumentuj wyzwalacze, zależności, umowy SLA, procedury wycofywania/ponawiania oraz bezpieczne kroki uzupełniania danych (backfill). Uwzględnij „odtwarzanie na stałym zbiorze danych” (fixed dataset replay) dla Dataflow, sposób opróżniania zadania strumieniowego (drain), ponownego przetwarzania nieudanych partycji oraz naprawiania wiadomości z kolejki DLQ.
- Zarejestruj typowe sygnatury błędów (np. odmowa dostępu, przekroczenie limitu, niezgodność schematu) wraz z drzewami decyzyjnymi i ścieżkami eskalacji.
Praktyczny scenariusz problemu
Firma Acme Retail Analytics musi codziennie pozyskiwać pliki CSV od partnerów, które czasami zawierają nieprawidłowo sformatowane wiersze, transformować i ładować prawidłowe dane do BigQuery oraz udostępniać błędne wiersze do analizy. Chce również wprowadzić wzbogacanie sterowane zdarzeniami w celu aktualizacji cen w czasie zbliżonym do rzeczywistego oraz zapewnić bezpieczne promowanie zmian ze środowiska deweloperskiego (dev) na produkcyjne (prod).
Podejście:
Przechowywanie i wyzwalacze zdarzeń
- Utwórz dedykowany bucket Cloud Storage z włączonym wersjonowaniem obiektów i jednolitym dostępem na poziomie bucketa. Włącz powiadomienia o finalizacji obiektu (object finalize) do Pub/Sub za pośrednictwem Eventarc.
- Uzasadnienie: Finalizacja obiektu jest niezawodnym zdarzeniem do wyzwalania dalszego pozyskiwania danych; wersjonowanie wspiera ponowne uruchomienia i audyty.
Przetwarzanie wsadowe z obsługą kolejki niedostarczonych wiadomości (dead-letter)
- Użyj Cloud Composer, aby uruchamiać codzienny DAG Airflow o godzinie 02:00 z włączoną opcją
catchup. DAG uruchamia zadanie wsadowe Dataflow, które parsuje pliki CSV, waliduje schemat i zapisuje prawidłowe rekordy do BigQuery, używając deterministycznych tabel przejściowych (staging), a następnie wykonuje operacjęMERGEdo docelowych tabel partycjonowanych. Przekierowuj nieprawidłowo sformatowane/nieudane rekordy do tabeli dead-letter w BigQuery. - Uzasadnienie: Dataflow skaluje parsowanie/walidację;
MERGEzapewnia idempotentność; przechwytywanie do kolejki dead-letter umożliwia inspekcję bez blokowania potoku, co jest zgodne z zalecanym wzorcem dla nieprawidłowo sformatowanych wierszy.
- Użyj Cloud Composer, aby uruchamiać codzienny DAG Airflow o godzinie 02:00 z włączoną opcją
Wzbogacanie sterowane zdarzeniami
- Wdróż zadanie Cloud Run do wykonywania lekkiego wzbogacania dla przyrostowych aktualizacji cen. Wyzwalaj je za pomocą Cloud Workflows, które nasłuchują na wiadomości Pub/Sub z Eventarc, gdy w ciągu dnia pojawiają się małe pliki z aktualizacjami.
- Uzasadnienie: Kontenery bezserwerowe z Workflows zapewniają orkiestrację o niskim opóźnieniu i niskich wymaganiach operacyjnych dla małych zdarzeń, pozostawiając ciężkie transformacje w trybie wsadowym.
Mechanizmy niezawodności
- Skonfiguruj ponowienia z wykładniczym czasem oczekiwania (exponential backoff) dla błędów przejściowych w zadaniach Dataflow i Cloud Run, ograniczając całkowity czas ponowień do SLA dla DAG-a. Ustaw limity czasu wykonania dla poszczególnych zadań (
execution_timeouts) i wywołania zwrotneon_failurew Airflow; w Workflows ustawmax_doublingsimax_retry_duration. - Uzasadnienie: Ograniczony backoff chroni umowy SLA i zapobiega niekontrolowanym ponowieniom.
- Skonfiguruj ponowienia z wykładniczym czasem oczekiwania (exponential backoff) dla błędów przejściowych w zadaniach Dataflow i Cloud Run, ograniczając całkowity czas ponowień do SLA dla DAG-a. Ustaw limity czasu wykonania dla poszczególnych zadań (
Bezpieczeństwo i zasada najmniejszych uprawnień
- Uruchamiaj każdy komponent z dedykowanym kontem serwisowym (SA): SA orkiestratora Composer, SA workera Dataflow, SA zadania Cloud Run. Nadawaj tylko wymagane role: odczyt GCS w buckecie z danymi wejściowymi dla Dataflow,
dataEditorw BigQuery dla docelowych zbiorów danych orazViewerdla logów. Przechowuj sekrety w Secret Manager i odwołuj się do nich w czasie działania. - Uzasadnienie: Wymusza zasadę najmniejszych uprawnień i izoluje promień rażenia (blast radius).
- Uruchamiaj każdy komponent z dedykowanym kontem serwisowym (SA): SA orkiestratora Composer, SA workera Dataflow, SA zadania Cloud Run. Nadawaj tylko wymagane role: odczyt GCS w buckecie z danymi wejściowymi dla Dataflow,
Orkiestracja sterowana metadanymi
- Utrzymuj tabelę kontrolną w BigQuery, zawierającą listę źródeł partnerskich, wzorce plików i docelowe zbiory danych. W czasie działania DAG-a, Airflow odpytuje tę tabelę i używa dynamicznego mapowania zadań (dynamic task mapping) do tworzenia zadań dla każdego partnera.
- Uzasadnienie: Dodanie partnera staje się zmianą danych, a nie zmianą w kodzie, co zmniejsza ryzyko wdrożenia.
Obserwowalność i alerty
- Emituj ustrukturyzowane logi z
run_idipartner_id. Utwórz polityki alertów dla naruszeń SLA DAG-a, opóźnień systemowych Dataflow i niezerowej liczby wiadomości w kolejce dead-letter. Dla operacji wstawiania do tabeli docelowej w BigQuery, skonfiguruj ujście (sink) Cloud Logging z zaawansowanym filtrem dla tej tabeli do tematu Pub/Sub, który jest konsumowany przez narzędzie monitorujące Acme. - Uzasadnienie: Szczegółowe alerty umożliwiają szybką klasyfikację problemów (triage) bez zbędnego szumu.
- Emituj ustrukturyzowane logi z
CI/CD i promowanie zmian
- Zarządzaj infrastrukturą (buckety, Pub/Sub, Eventarc, Composer, Workflows, zbiory danych BigQuery) w Terraform. Użyj Cloud Build do walidacji składni DAG-ów Airflow, uruchamiania testów jednostkowych i wdrażania na środowisko deweloperskie Composer. Promuj na środowisko testowe i produkcyjne za pomocą sparametryzowanych konfiguracji i ręcznych bramek zatwierdzających, po pomyślnym przejściu asercji Dataform i testów integracyjnych.
- Uzasadnienie: Deklaratywne, powtarzalne wdrożenia i bezpieczne promowanie zmian między środowiskami.
Runbook i odzyskiwanie po awarii
- Udokumentuj kroki w celu ponownego przetworzenia danych z określonej daty: przywróć plik CSV z wersjonowania obiektów, ponownie uruchom zadanie Dataflow dla tej partycji, scal (
MERGE) wyniki i przejrzyj rekordy z DLQ. Uwzględnij procedurę „odtwarzania na stałym zbiorze danych” (fixed dataset replay), aby izolować błędy transformacji, jeśli pojawią się rozbieżności. - Uzasadnienie: Idempotentny projekt i udokumentowane procedury odzyskiwania usprawniają naprawę częściowych awarii.
- Udokumentuj kroki w celu ponownego przetworzenia danych z określonej daty: przywróć plik CSV z wersjonowania obiektów, ponownie uruchom zadanie Dataflow dla tej partycji, scal (
← Pozyskiwanie · Wszystkie domeny · Uczenie maszynowe →
Przećwicz te pytania → · Testy na czas na 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.
Zdaj egzamin →