Google PDE: Przetwarzanie strumieniowe z Dataflow i Apache Beam — 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
Przetwarzanie strumieniowe w Google Cloud opiera się na zunifikowanym modelu programowania Apache Beam, wykonywanym przez runner Dataflow. Beam dostarcza logicznej abstrakcji — potoków transformacji działających na kolekcjach PCollection — która oddziela kod od szczegółów wykonania, takich jak równoległość, autoskalowanie i odporność na awarie. W przetwarzaniu strumieniowym poprawność zależy od semantyki czasu (czas zdarzenia a czas przetwarzania), okienkowania (stałe, przesuwne, sesyjne, globalne), znaków wodnych (watermarks), wyzwalaczy (triggers) i obsługi opóźnionych danych. Doskonałość operacyjna w Dataflow wymaga odpowiedniego doboru rozmiaru workerów, polityki autoskalowania, silnika strumieniowego, wyboru mechanizmu shuffle, projektu idempotentnego ujścia, obsługi martwych listów (dead-letter) i solidnej obserwowalności.
Model Apache Beam i semantyka czasu
Potoki, transformacje, PCollections, runnery:
- Potok Beam stosuje skierowany graf acykliczny transformacji PTransform do kolekcji PCollection (ograniczonych lub nieograniczonych).
- Runnery (Dataflow, Spark, Flink, Direct) wykonują potok; Dataflow zapewnia zarządzane autoskalowanie, checkpointing i wgląd operacyjny.
- Transformacje obejmują operacje na pojedynczych elementach (ParDo), grupowanie i łączenie (GroupByKey, Combine), złączenia (CoGroupByKey) oraz operacje wejścia/wyjścia (PubSubIO, BigQueryIO, FileIO).
Okna (Windows):
- Okna stałe (fixed windows): niepokrywające się przedziały czasowe (np. 1-minutowe okna stałe, tzw. tumbling windows) do okresowych agregacji.
- Okna przesuwne (sliding windows): nakładające się okna do płynnych metryk kroczących (np. 5-minutowe okna przesuwane co 1 minutę).
- Okna sesyjne (session windows): dynamiczne okna zamykane po przerwie w aktywności, idealne do śledzenia sesji użytkowników lub nagłych wzrostów aktywności urządzeń.
- Okno globalne (global window): domyślny widok całego nieograniczonego strumienia bez podziału na okna; często używane w parze z wyzwalaczami do okresowej materializacji.
Czas zdarzenia a czas przetwarzania:
- Czas zdarzenia (event time): moment, w którym zdarzenie wystąpiło w źródle; umożliwia logicznie spójne agregacje pomimo zmiennych opóźnień w transporcie.
- Czas przetwarzania (processing time): moment, w którym zdarzenie jest obserwowane przez potok; przydatny dla wyzwalaczy operacyjnych, ale nie dla poprawności semantycznej.
Znaki wodne (Watermarks):
- Znak wodny (watermark) szacuje kompletność czasu zdarzenia (przypuszczenie runnera, że zobaczył wszystkie zdarzenia do czasu T).
- Znaki wodne mogą przesuwać się nieregularnie lub zatrzymywać pod wpływem ciśnienia zwrotnego (backpressure) lub opóźnień w źródle; opóźnione dane to wszystko, co przybywa ze znacznikiem czasu mniejszym niż znak wodny.
Wyzwalacze (Triggers) i opóźnienia:
- Domyślnie: wyzwalacz AfterWatermark, który uruchamia się, gdy znak wodny przekroczy koniec okna; przy dozwolonym opóźnieniu (allowed lateness) = 0, opóźnione dane są odrzucane.
- Wczesne wyzwolenia (oparte na czasie przetwarzania lub liczbie elementów) dają wstępne wyniki o niskim opóźnieniu.
- Późne wyzwolenia pozwalają na korekty, gdy nadejdą opóźnione dane; tryb akumulacji (accumulation mode) decyduje, czy panele (panes) akumulują wyniki, czy odrzucają poprzednie.
- Dobierz dozwolone opóźnienie (allowed lateness) na podstawie tolerancji biznesowej i kompromisów między zasobami dyskowymi a obliczeniowymi; większe opóźnienie zwiększa przechowywanie stanu i koszty.
Przetwarzanie stanowe, timery, sesjonizacja, deduplikacja:
- Stanowe funkcje DoFn przechowują stan per-klucz (np. ostatnio widziane zdarzenie, agregacje bieżące) i ustawiają timery do emisji lub czyszczenia stanu.
- Sesjonizację naturalnie wyraża się za pomocą okien sesyjnych (SessionWindows); dla logiki niestandardowej użyj stanu kluczowanego (keyed state) oraz timerów opartych na czasie przetwarzania/zdarzenia.
- Deduplikacja: użyj stabilnego ID dla każdego zdarzenia i zastosuj Distinct/Combine w obrębie okna lub stan per-klucz (np. filtr Blooma lub zbiór z TTL). Należy znaleźć kompromis między zużyciem pamięci i wynikami fałszywie pozytywnymi a ścisłą dokładnością.
Tryby awarii i kompromisy:
- Używanie okien opartych na czasie przetwarzania dla metryk biznesowych powoduje dryf podczas skoków obciążenia lub ponownych prób; preferuj okna oparte na czasie zdarzenia.
- Zbyt małe okna z częstymi wczesnymi wyzwoleniami powodują nadmierną emisję paneli (panes) i wzmocnienie zapisu do ujścia.
- Nieograniczone dozwolone opóźnienie może prowadzić do rozrostu stanu; zawsze ograniczaj TTL stanu i ustawiaj timery do czyszczenia nieaktywnych kluczy.
Obsługa Dataflow dla obciążeń strumieniowych
Wymiarowanie i autoskalowanie maszyn roboczych:
- Autoskalowanie horyzontalne dodaje/usuwa maszyny robocze w oparciu o zaległości, opóźnienie znaku wodnego, użycie CPU i przepustowość; ustaw rozsądną wartość maxWorkers, aby absorbować skoki obciążenia.
- Dobieraj typy maszyn do wąskich gardeł: ograniczone przez CPU (więcej vCPU), ograniczone przez pamięć (typy z dużą ilością pamięci), ograniczone przez sieć (większe VM redukują narzut związany z operacją shuffle).
- Zwiększ dysk rozruchowy przy intensywnych operacjach shuffle lub ujściach opartych na plikach. Monitoruj opóźnienie systemowe (system lag) i zaległości w sekundach (backlog seconds).
Streaming Engine i shuffle:
- Streaming Engine przenosi stan i operacje shuffle do backendu usługi, poprawiając elastyczność, zmniejszając obciążenie pamięci maszyn roboczych i umożliwiając szybsze aktualizacje.
- W przypadku etapów z dużym obciążeniem wsadowym lub masowym grupowaniem kluczy, użyj Dataflow Shuffle, aby odciążyć I/O operacji shuffle z maszyn roboczych. Oba rozwiązania redukują awarie przeciążonych maszyn roboczych i nadmierne obciążenie dysku.
Backpressure, gorące klucze i skośność danych:
- Dataflow zarządza zjawiskiem backpressure poprzez dynamiczne równoważenie pracy; niemniej jednak, w stosownych przypadkach, dostrajaj kontrolę przepływu u źródła (np. w Pub/Sub liczbę nieprzetworzonych wiadomości/bajtów).
- Gorące klucze (np. popularne ID) tworzą procesy marudzące (stragglers). Łagodź ten problem przez sharding kluczy (key#N), częściową pre-agregację z ponownym kluczowaniem lub aproksymacje oparte na szkicach (sketch-based).
- Skośność danych wynikająca z rekordów odstających (ogromne ładunki) lub niestabilnych wydawców może wymagać partycjonowania per wydawca, batchingu lub kompresji.
Integracja z Pub/Sub:
- Używaj tematów Pub/Sub do pozyskiwania danych; włącz atrybuty wiadomości dla metadanych (np. deviceId, znacznik czasu zdarzenia).
- Pozyskuj dane za pomocą PubSubIO; wyodrębniaj znaczniki czasu zdarzeń z atrybutów lub z ładunku, w przeciwnym razie użyj czasu publikacji.
- Klucze porządkujące (ordering keys) zapewniają porządek w ramach danego klucza; Dataflow wciąż wymaga idempotentnego zachowania w systemach docelowych z powodu dostarczania co najmniej raz (at-least-once).
Wzorce strumieniowania do BigQuery:
- Preferuj BigQueryIO z Storage Write API dla wysokiej przepustowości i niskich opóźnień z semantyką „dokładnie raz” w ramach strumienia, dzięki przesunięciom strumienia (stream offsets) i automatycznym ponowieniom.
- Dla prostych potoków o niskiej przepustowości, wstawianie strumieniowe (streaming inserts) jest akceptowalne; ustawiaj insertId w celu deduplikacji ponowień po stronie klienta.
- Zapytania do buforów strumieniowych są ostatecznie spójne (eventually consistent); dla analityki krytycznej czasowo, wykonuj zapytania po opóźnieniu bufora (np. odczekaj ~2x zaobserwowane opóźnienie dostępności) lub materializuj dane za pomocą okien mikro-wsadowych i trybu committed w Storage Write API.
Efekty „dokładnie raz”, idempotencja, ponowne odtwarzanie i ujścia:
- Beam gwarantuje przetwarzanie co najmniej raz; efekt „dokładnie raz” musi być osiągnięty w ujściu (sink) za pomocą idempotentnych zapisów, transakcji lub kluczy deduplikacyjnych.
- BigQuery: używaj domyślnych strumieni (default streams) lub zatwierdzonych strumieni (committed streams) w Storage Write API dla semantyki „dokładnie raz” w ramach strumienia; przy wstawianiu strumieniowym ustawiaj stabilny insertId.
- Pliki: zapisuj do plików tymczasowych o unikalnych nazwach, finalizuj po zamknięciu okna i zapewnij atomowe zmiany nazw; unikaj nadpisywania, aby zapobiec częściowym duplikatom.
- Zewnętrzne bazy danych: używaj operacji upsert z kluczem opartym na stabilnym ID lub zaimplementuj okna deduplikacyjne.
- Projektuj z myślą o ponownym odtwarzaniu: utrzymuj deterministyczne transformacje; upewnij się, że ujścia deduplikują dane przy ponowieniu próby.
Obsługa niedostarczalnych wiadomości, routing błędów, obserwowalność:
- Ryzykowne operacje parsowania/wzbogacania umieszczaj w bloku try/catch wewnątrz ParDo i emituj błędy do PCollection na niedostarczalne wiadomości za pomocą TupleTag; dołączaj ładunek, kod błędu i kontekst.
- Kieruj kolejki DLQ do BigQuery lub Cloud Storage w celu analizy; rozważ użycie osobnego tematu Pub/Sub do ponownego przetwarzania.
- Obserwowalność: używaj metryk zadań Dataflow (opóźnienie znaku wodnego, opóźnienie systemowe, przepustowość), niestandardowych liczników, metryk dystrybucji i logów per krok w Cloud Logging. Twórz alerty w Cloud Monitoring dotyczące opóźnień i wskaźników błędów. Używaj Error Reporting do agregowania wyjątków.
Wzorce strojenia wydajności:
- Czytaj wydajnie: dla źródeł BigQuery preferuj Storage Read API lub odczyty oparte na zapytaniach, które wybierają tylko potrzebne pola i stosują filtry.
- Optymalizacja łączenia (Combine lifting): używaj funkcji łączących (combiners), aby zredukować wolumen danych w operacji shuffle przed GroupByKey.
- Wejścia boczne (side inputs): przechowuj małe dane referencyjne w pamięci podręcznej; zwracaj uwagę na fanout i kadencję aktualizacji.
- Serializacja: używaj kompaktowych schematów (Avro/Proto) i unikaj nadmiernego parsowania JSON na gorących ścieżkach (hot paths).
Wdrożenia, szablony i strategie aktualizacji
Flex Templates:
- Pakują potoki w skonteneryzowane, sparametryzowane szablony w celu zapewnienia powtarzalnych wdrożeń. Flex Templates obsługują niestandardowe zależności, obrazy GPU i izolację środowiska.
- Umożliwiają eksternalizację parametrów uruchomieniowych (np. subskrypcja wejściowa, tabela wyjściowa, ujście dla martwych listów, maxWorkers) w celu wdrożeń specyficznych dla danego środowiska.
Aktualizacje potoków i kompatybilność:
- Dataflow obsługuje aktualizacje w miejscu (in-place) dla wielu potoków strumieniowych, jeśli nazwy transformacji, specyfikacje stanu i typy wyjściowe pozostają kompatybilne. Należy używać stabilnych nazw PTransform.
- W przypadku niekompatybilnych zmian w grafie lub stanie, należy przeprowadzić kontrolowane przełączenie (cutover): uruchomić nowe zadanie, a następnie wygasić (drain) stare, aby zakończyć przetwarzanie bieżących elementów i przestać odczytywać nowe.
Wygaszanie (draining) i snapshoty:
- Wygaszanie w sposób kontrolowany kończy przetwarzanie, zapisuje pozostałe dane wyjściowe i kończy działanie; należy koordynować je z retencją Pub/Sub lub snapshotami, aby uniknąć przerw w danych.
- Aby zapewnić ciągłość, można utworzyć snapshot Pub/Sub, uruchomić nowy potok odczytujący od tego snapshotu lub odpowiedniego znacznika czasu, zweryfikować dane wyjściowe, a następnie wygasić stare zadanie.
Przykłady konfiguracji:
- Przykład okienkowania z wczesnymi/późnymi wyzwalaczami i akumulacją:
undefined
- Przykład BigQueryIO z użyciem Storage Write API:
undefined
- Częste pułapki:
- Zapisywanie do ujść plikowych w trybie strumieniowym bez zapisów okienkowych może zablokować finalizację; należy włączyć zapisy okienkowe i wyzwalacze.
- Nieograniczony wzrost: zapomnienie o ograniczeniu stanu lub dozwolonego opóźnienia (allowed lateness) może powodować wycieki pamięci i błędy skalowania.
- Brakujące znaczniki czasu: nieprzypisanie znaczników czasu zdarzeń powoduje, że potok domyślnie używa czasu przetwarzania i traci poprawność przy zmiennych opóźnieniach.
Praktyczny scenariusz problemowy
Firma NovaTrack Inc. przetwarza globalną telemetrię IoT z 50 000 czujników temperatury i musi dostarczać agregaty na poziomie minutowym, utrwalać surowe dane i prezentować pulpit nawigacyjny w czasie rzeczywistym. Oczekiwane są sporadyczne, źle sformatowane komunikaty i dostarczanie danych poza kolejnością. Rozwiązanie musi automatycznie się skalować, udostępniać błędne rekordy do inspekcji i wspierać aktualizacje bez przestojów (zero-downtime).
Podejście:
Pozyskiwanie danych i semantyka czasu
- Utworzenie regionalnego tematu Pub/Sub i publisherów w każdym regionie z atrybutami deviceId i eventTs (RFC3339). Włączenie kluczy porządkujących (ordering keys) według deviceId, gdy jest to możliwe.
- Uzasadnienie: Pub/Sub zapewnia trwałe, elastyczne wejście danych z gwarancją dostarczenia co najmniej raz (at-least-once). Dołączanie znaczników czasu zdarzeń na brzegu sieci (at the edge) zachowuje prawdziwy czas zdarzenia; porządkowanie według urządzenia zmniejsza problemy z kolejnością wewnątrz strumienia danych z jednego urządzenia bez tworzenia centralnych wąskich gardeł.
Potok strumieniowy Dataflow z oknami czasu zdarzenia
- Odczyt z dedykowanej subskrypcji za pomocą PubSubIO, wyodrębniając eventTs jako znacznik czasu Beam, z powrotem do publishTime w przypadku jego braku.
- Zastosowanie FixedWindows o długości 1 minuty z wczesnym wyzwalaczem po 30 sekundach i późnymi wyzwoleniami dla każdego opóźnionego elementu; ustawienie dozwolonego opóźnienia (allowed lateness) na 10 minut i akumulowanie paneli (accumulating panes).
- Uzasadnienie: Okna czasu zdarzenia zapewniają dokładne agregaty minutowe; wczesne wyzwolenia zasilają pulpit nawigacyjny ze świeżością poniżej minuty; późne wyzwolenia korygują agregaty w miarę napływania opóźnionych danych. Ograniczenie opóźnienia limituje rozmiar stanu i koszty.
Walidacja, wzbogacanie i routing do kolejki martwych listów (dead-letter)
- Implementacja ParDo, które parsuje JSON, waliduje schemat i zakresy oraz wzbogaca dane o małe, statyczne dane referencyjne za pomocą wejścia bocznego (side input) ładowanego z BigQuery przy starcie zadania.
- Użycie TupleTags do emitowania poprawnych rekordów do głównego wyjścia, a błędów do PCollection martwych listów (dead-letter), zawierającej payload, błąd, deviceId i znacznik czasu parsowania; zapis DLQ do partycjonowanej tabeli BigQuery.
- Uzasadnienie: Wejścia boczne (side inputs) przechowują dane referencyjne w pamięci, zapewniając niskie opóźnienia. Przechwytywanie martwych listów pozwala na inspekcję i ukierunkowane ponowne przetwarzanie błędnych wierszy bez blokowania głównego przepływu.
Agregacja i łagodzenie problemu gorących kluczy (hot-key)
- Kluczowanie według deviceId i obliczanie minutowych wartości avg/min/max za pomocą CombineFns. Dla metryk regionalnych typu top-N, shardowanie według region#N, aby uniknąć gorących kluczy, a następnie ponowna agregacja.
- Uzasadnienie: Funkcje Combine minimalizują wolumen i koszt operacji shuffle; shardowanie kluczy zapobiega wąskim gardłom związanym z pojedynczym kluczem podczas agregacji regionalnej (fan-in).
Ujścia i efekty “dokładnie raz” (exactly-once)
- Zapis surowych, zwalidowanych zdarzeń i agregatów minutowych do BigQuery przy użyciu BigQueryIO z Storage Write API. Ustawienie stabilnego identyfikatora wstawiania (insert id) opartego na deviceId + eventTs w celu zapewnienia idempotencji przy niestandardowych ponowieniach.
- Uzasadnienie: Storage Write API zapewnia wysokoprzepustowe pozyskiwanie danych o niskim opóźnieniu z semantyką “dokładnie raz” (exactly-once) w ramach jednego strumienia. Stabilne identyfikatory zapewniają deduplikację po stronie odbiorcy w przypadku ponownego odtwarzania danych.
Strategia spójności pulpitu nawigacyjnego
- Pulpit nawigacyjny odpytuje partycjonowane tabele agregatów z oknem czasowym 2 minut wstecz w stosunku do znaku wodnego (watermark) lub ze stałym opóźnieniem 2x obserwowanego opóźnienia dostępności dla danych strumieniowych.
- Uzasadnienie: Widoczność danych strumieniowych w BigQuery jest ostatecznie spójna (eventually consistent); lekkie opóźnienie odczytów zapobiega pomijaniu wierszy w trakcie przesyłania, zachowując jednocześnie zachowanie bliskie czasu rzeczywistego.
Operacje: autoskalowanie i silnik strumieniowy
- Włączenie Streaming Engine; ustawienie maxWorkers na podstawie oczekiwanego szczytu (np. 3x średniej), wybór typu maszyny zwymiarowanego pod kątem parsowania i szyfrowania intensywnie wykorzystujących CPU oraz zwiększenie dysku rozruchowego, aby pomieścić tymczasowe dane operacji shuffle.
- Monitorowanie opóźnienia znaku wodnego (watermark lag), zaległości w sekundach (backlog seconds), użycia CPU i przepustowości na krok; alertowanie przy trwałym opóźnieniu i skokach wskaźnika DLQ.
- Uzasadnienie: Streaming Engine eksternalizuje stan/shuffle, co zapewnia elastyczność i prostsze aktualizacje; odpowiednie zwymiarowanie i monitorowanie zapobiegają cichym naruszeniom SLO.
Wdrożenie i aktualizacje za pomocą Flex Templates
- Spakowanie potoku jako Flex Template z parametrami: subskrypcja wejściowa, tabele wyjściowe, tabela DLQ, maxWorkers i region. W przypadku niekompatybilnej zmiany, uruchomienie nowego potoku skierowanego na ten sam temat z nową subskrypcją, zweryfikowanie wyników, a następnie wygaszenie (drain) starego zadania. Opcjonalnie można utworzyć snapshot Pub/Sub i ustawić nową subskrypcję na odczyt od tego snapshotu, aby zagwarantować brak przerw w danych.
- Uzasadnienie: Flex Templates umożliwiają powtarzalne, sparametryzowane wdrożenia. Zweryfikowane przełączenie typu blue/green z wygaszaniem (drain) zapewnia zerową utratę danych i minimalny czas przestoju.
Ponowne przetwarzanie i uzupełnianie wsadowe
- Przechowywanie skompresowanych plików Avro z surowymi zdarzeniami w Cloud Storage za pomocą wyjścia bocznego (side output); uruchamianie wsadowego potoku Dataflow w celu uzupełnienia lub ponownego przetworzenia danych do BigQuery, gdy zmieniają się modele lub schematy.
- Uzasadnienie: Trwałe archiwa surowych danych wspierają powtarzalność i ewolucję schematów bez wpływu na ścieżkę gorącą (hot path).
Ten projekt zapewnia poprawne agregaty o niskim opóźnieniu i ograniczonym koszcie, z wyraźną izolacją błędów, silną obserwowalnością i bezpiecznymi ścieżkami aktualizacji, jednocześnie obsługując dane poza kolejnością i opóźnione w skali globalnej.
← Analityka BigQuery i inżynieria hurtowni danych · Wszystkie domeny · Przesyłanie komunikatów →
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 →