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

Tryby awarii i kompromisy:

Obsługa Dataflow dla obciążeń strumieniowych

Wdrożenia, szablony i strategie aktualizacji

undefined

undefined

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:

  1. 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ł.
  2. 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.
  3. 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.
  4. 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).
  5. 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.
  6. 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.
  7. 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.
  8. 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.
  9. 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 →

Przeglądaj Google →

Related guides

Dostęp all-in-one

Jedna subskrypcja. Każdy egzamin.

Każdy plan odblokowuje nieograniczone wyszukiwanie odpowiedzi, testy praktyczne, wyjaśnienia AI i pełną bibliotekę zasobów — w ponad 20 językach.

Miesięczny
24.87
Just €0.83/day
Wszystko w cenie:
  • Nieograniczone wyszukiwanie odpowiedzi
  • Nieograniczone testy praktyczne
  • Wyjaśnienia wspomagane AI
  • Pełna biblioteka zasobów
  • Ponad 20 języków
  • Cotygodniowe aktualizacje treści
  • Nagrody i polecenia
  • Priorytetowe wsparcie
Rozpocznij bezpłatny okres próbny

Karta kredytowa nie jest wymagana*

Najlepsza wartość
12 miesięcy
179.87
Just €0.49/daySave 40%
Wszystko w cenie:
  • Nieograniczone wyszukiwanie odpowiedzi
  • Nieograniczone testy praktyczne
  • Wyjaśnienia wspomagane AI
  • Pełna biblioteka zasobów
  • Ponad 20 języków
  • Cotygodniowe aktualizacje treści
  • Nagrody i polecenia
  • Priorytetowe wsparcie
Rozpocznij bezpłatny okres próbny

Karta kredytowa nie jest wymagana*

✓ Plan darmowy w zestawie · ✓ Anuluj w dowolnym momencie · ✓ Wszystkie plany odblokowują pełny produkt