Google PDE: Przesyłanie komunikatów, pozyskiwanie zdarzeń i usługi czasu rzeczywistego — 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
Usługi przesyłania komunikatów, pozyskiwania zdarzeń i przetwarzania w czasie rzeczywistym w Google Cloud opierają się na Cloud Pub/Sub i Eventarc w celu zapewnienia rozproszonego i trwałego transportu, Dataflow do stanowego przetwarzania strumieniowego oraz ujściach (sinks) takich jak BigQuery, Cloud Storage i operacyjne bazy danych. Projektowanie z myślą o dostarczaniu co najmniej raz, idempotentnej konsumpcji i obserwowalności zapewnia odporne systemy, które skalują się elastycznie, zachowując poprawność działania w warunkach awarii, przeciwciśnienia (backpressure) i ewolucji schematu.
Podstawy przesyłania komunikatów z Pub/Sub
- Tematy i subskrypcje
- Wydawcy (publishers) wysyłają komunikaty do tematu; subskrybenci dołączają za pomocą subskrypcji (wielu subskrybentów może niezależnie konsumować te same komunikaty).
- Typy subskrypcji:
- Pull: klienci jawnie pobierają komunikaty; użyj
streaming pulldla najwyższej przepustowości i mniejszej liczby zapytań i odpowiedzi (round trips). - Push: Pub/Sub dostarcza komunikaty przez HTTPS; twój punkt końcowy musi zwrócić kod 2xx, aby potwierdzić odbiór.
- Eksport do BigQuery: subskrypcja BigQuery dostarcza komunikaty do tabeli BigQuery bez konieczności pisania kodu; najlepsze rozwiązanie, gdy ładunki (payloads) pasują do zadeklarowanego schematu i wymagane jest pozyskiwanie danych do analityki z niskim opóźnieniem.
- Pull: klienci jawnie pobierają komunikaty; użyj
- Klucze porządkujące
- Włącz porządkowanie komunikatów w temacie i subskrypcji, aby otrzymywać je w kolejności dla danego klucza porządkującego. Przepustowość na klucz jest serializowana: jeden przetwarzany komunikat na klucz może blokować kolejne; używaj wielu kluczy (na przykład
hash(device_id)), aby skalować system.
- Włącz porządkowanie komunikatów w temacie i subskrypcji, aby otrzymywać je w kolejności dla danego klucza porządkującego. Przepustowość na klucz jest serializowana: jeden przetwarzany komunikat na klucz może blokować kolejne; używaj wielu kluczy (na przykład
- Fan-out i ponowne odtwarzanie (replay)
- Twórz oddzielne subskrypcje dla różnych konsumentów, aby izolować obciążenia i retencję.
- Użyj operacji
seeklubsnapshot, aby odtworzyć komunikaty od określonego znacznika czasu lub z migawki w celu odzyskiwania danych i uzupełniania braków (backfills).
Kompromisy:
- Porządkowanie zmniejsza równoległość i przepustowość na klucz; wyłączaj porządkowanie, chyba że jest to absolutnie konieczne.
- Tryb Push upraszcza kod klienta, ale wprowadza kwestie związane ze skalowaniem punktu końcowego HTTP, bezpieczeństwem i mechanizmami backoff; tryb Pull daje większą kontrolę i stabilność przy wysokiej przepustowości.
Semantyka dostarczania, potwierdzanie, retencja i obsługa niedostarczalnych komunikatów
- Potwierdzanie i terminy (deadlines)
- Dostarczanie co najmniej raz: duplikaty mogą się zdarzyć.
- Każde dostarczenie ma termin potwierdzenia (
ack deadline, domyślnie 10 sekund). Przedłużaj go (ModifyAckDeadline) podczas przetwarzania długotrwałych zadań; brak potwierdzenia przed upływem terminu jest najczęstszą przyczyną zduplikowanych dostarczeń w trybie push. - Negatywne potwierdzenie (
nack) lub wygaśnięcie terminu sprawia, że komunikat kwalifikuje się do ponownego dostarczenia.
- Retencja
- Niepotwierdzone komunikaty są przechowywane przez czas określony w terminie potwierdzenia subskrypcji i ponawiane; potwierdzone komunikaty mogą być przechowywane przez okres retencji komunikatu tematu w celu ponownego odtworzenia. Skonfiguruj retencję tak, aby obejmowała maksymalny czas awarii plus czas potrzebny na odzyskanie systemu.
- Ponowienia
- Pull: ponowne dostarczenie następuje po wygaśnięciu terminu potwierdzenia; kontroluj współbieżność za pomocą limitów kontroli przepływu (flow control).
- Push: stosowany jest wykładniczy backoff (exponential backoff); tylko odpowiedź HTTP 2xx jest traktowana jako sukces. Kody 3xx/4xx/5xx powodują ponowienia. Implementuj idempotentne procedury obsługi, aby tolerować powtórzenia.
- Tematy wiadomości niedostarczalnych (DLT)
- Skonfiguruj temat DL i maksymalną liczbę prób dostarczenia dla każdej subskrypcji, aby przenosić do kwarantanny komunikaty powodujące błędy (poison messages).
- Monitoruj wolumen kolejki DLQ; twórz procesy analizy (triage) i ponownie publikuj do głównego tematu po poprawieniu błędu.
Przykład:
undefined
Podsumowanie semantyki dostarczania:
- Pub/Sub: dostarczanie co najmniej raz, porządkowanie typu „best-effort” (na miarę możliwości) w ramach klucza porządkującego, jeśli jest włączone.
- Ujścia (sinks): API wstawiania BigQuery zapewniają ograniczanie duplikatów (
insertIdlub przesunięcia strumienia w Storage Write API), ale i tak projektuj konsumentów i komponenty zapisujące tak, aby były idempotentne.
Schematy, kompatybilność i walidacja
- Schematy Pub/Sub
- Natywne wsparcie dla Avro i Protocol Buffers ze schematami przechowywanymi centralnie.
- Ustawienia schematu na poziomie tematu: kodowanie (Avro lub Protobuf) i egzekwowanie (brak, tylko walidacja lub wymagane).
- Producent publikuje zakodowane ładunki; Pub/Sub waliduje je względem bieżącego schematu, gdy egzekwowanie jest włączone.
- Ewolucja i kompatybilność
- Używaj zmian kompatybilnych wstecznie (dodawanie pól opcjonalnych, dodawanie pól z wartościami domyślnymi w Avro, nigdy nie używaj ponownie tagów w Protobuf, unikanie usuwania lub zmiany nazw pól).
- Jawnie wersjonuj schematy. W przypadku zmian niekompatybilnych (breaking changes), publikuj podwójnie do tematów v1 i v2 lub dodaj pole wersji i kieruj ruch odpowiednio.
- Kontrakty producent-konsument
- Konsumenci powinni ignorować nieznane pola i przypisywać wartości domyślne brakującym.
- Testuj kompatybilność schematu ze wszystkimi konsumentami przed wdrożeniem na produkcję; waliduj na subskrypcjach deweloperskich (staging) z takim samym egzekwowaniem schematu jak na produkcji.
Krótki przykład Avro (fragment):
undefined
Integracja sterowana zdarzeniami, Eventarc i interoperacyjność z Kafka
- Eventarc i CloudEvents
- Eventarc kieruje zdarzenia z usług Google Cloud (i niestandardowych źródeł przez Pub/Sub) do Cloud Run, GKE lub Workflows, używając specyfikacji CloudEvents. Atrybuty takie jak typ, źródło i temat (subject) umożliwiają szczegółowe filtrowanie i audytowalność.
- Używaj filtrów atrybutów, aby zminimalizować fan-out i zmniejszyć obciążenie systemów docelowych (downstream).
- Dostarczanie odbywa się w trybie co najmniej raz (at-least-once); w miarę możliwości implementuj procedury obsługi (handlery) jako idempotentne i bezstanowe.
- Przykład wyzwalacza Eventarc:
gcloud eventarc triggers create gcs-finalize-to-run
–destination-run-service=ingestor
–event-filters=“type=google.cloud.storage.object.v1.finalized”
–event-filters=“bucket=my-data-bucket”
–service-account=eventarc-sa@PROJECT_ID.iam.gserviceaccount.com - Interoperacyjność z Kafka i zarządzana migracja
- Szablony Dataflow łączą Kafka <-> Pub/Sub w celu przeprowadzania migracji etapowej. Odzwierciedlaj tematy (topics) z zachowaniem kluczy; przełącz najpierw konsumentów, potem producentów, lub stosuj podwójny zapis (dual-write) w okresie przejściowym.
- Pub/Sub Lite oferuje partycjonowany streaming z alokowaną pojemnością (capacity-provisioned), routing oparty na kluczach i niższy koszt; jest usługą regionalną/strefową i nadaje się do obciążeń podobnych do Kafki, gdzie przewidywalna przepustowość i kolejność w ramach partycji są kluczowe.
- Kwestie do rozważenia podczas migracji:
- Kolejność: mapuj klucze Kafka na klucze porządkujące (ordering keys) Pub/Sub lub partycje Lite.
- Przesunięcia (offsets): przenoś offsety jako atrybuty wiadomości do celów diagnostycznych; po migracji konsumenci nie mogą polegać na offsetach z Kafki.
- Dostarczanie: zaakceptuj tryb co najmniej raz (at-least-once); wymuś idempotencję w systemach docelowych (downstream).
- Schematy: migruj definicje z Confluent Schema Registry do schematów Pub/Sub lub standaryzuj na Protobuf/Avro z kompatybilnymi zasadami ewolucji.
Wzorce pozyskiwania strumieniowego, przepustowość, skalowanie, bezpieczeństwo i operacje
- Wzorce pozyskiwania w czasie rzeczywistym
- Pub/Sub -> Dataflow -> BigQuery: użyj ujścia (sink) BigQuery Storage Write API dla wysokiej przepustowości i idempotencji z offsetami strumienia; przekierowuj błędy do tabeli dead-letter w celu inspekcji.
- Pub/Sub -> Dataflow -> Cloud Storage: archiwizuj surowe zdarzenia w celu ponownego przetwarzania; używaj okienkowanych, skompresowanych zapisów, aby zrównoważyć koszty i opóźnienia.
- Pub/Sub -> operacyjne bazy danych: zapisuj do Bigtable dla odczytów o niskim opóźnieniu, do Spanner dla transakcji o silnej spójności lub do Cloud SQL/Firestore w zależności od potrzeb obciążenia. Zapewnij idempotentne operacje upsert z kluczem w postaci unikalnego ID zdarzenia.
- Co najmniej jednokrotne dostarczenie, zapobieganie duplikatom i idempotencja
- Przenoś unikalny
event_idievent_timew każdej wiadomości; wymuszaj stosowanie UUID po stronie producenta. - Deduplikacja w BigQuery streaming: ustaw
insertIdlub użyj Storage Write API z uporządkowanymi strumieniami; nadal zabezpieczaj zapytania logiką deduplikacji. - Przykład deduplikacji w czasie zapytania: WITH ranked AS ( SELECT t.*, ROW_NUMBER() OVER (PARTITION BY event_id ORDER BY event_time DESC) AS rn FROM dataset.events t ) SELECT * EXCEPT(rn) FROM ranked WHERE rn = 1;
- Dla punktów końcowych push, zwracaj kod 2xx tylko po pomyślnym przetworzeniu; w przeciwnym razie oczekuj ponownego dostarczenia.
- Przenoś unikalny
- Przepustowość wiadomości, limity (quotas) i skalowanie
- Producenci: grupuj wiadomości w paczki (batch) i ponownie używaj połączeń; zrównoleglaj pracę na wielu klientach. Używaj wielu kluczy porządkujących (
ordering keys), aby skalować uporządkowane obciążenia. - Subskrybenci: preferuj
streaming pullz kontrolą przepływu (flow control) (maks. liczba niepotwierdzonych bajtów/wiadomości). Dopasujack deadlinesdo czasu przetwarzania i wydłużaj je w razie potrzeby. - Monitoruj i wnioskuj o zwiększenie limitów (quotas) dla przepustowości publikowania i subskrybowania w miarę wzrostu wolumenu; projektuj z zapasem (np. 2x oczekiwanego szczytu), aby absorbować nagłe wzrosty obciążenia.
- Producenci: grupuj wiadomości w paczki (batch) i ponownie używaj połączeń; zrównoleglaj pracę na wielu klientach. Używaj wielu kluczy porządkujących (
- Spójność i dostępność
- Streaming do BigQuery jest ostatecznie spójny (eventually consistent) pod względem widoczności w zapytaniach; dla zapytań interaktywnych, które muszą uwzględniać wiersze ze strumienia, poczekaj w oparciu o obserwowane opóźnienie (np. 2x opóźnienie dostępności P50) lub projektuj z użyciem agregacji wyrównanych do znaków wodnych (watermarks) w Dataflow i odpytuj zmaterializowane wyniki.
- Bezpieczeństwo
- IAM: nadawaj role z najmniejszymi uprawnieniami (
least-privilege) (pubsub.publisherdla producentów na poziomie tematu;pubsub.subscriberdla konsumentów na poziomie subskrypcji). Używaj dedykowanych kont serwisowych (service accounts) dla każdego obciążenia. - Uwierzytelnianie push: skonfiguruj subskrypcje push, aby dołączały tokeny OIDC z konta serwisowego; wymuszaj walidację
audiencena punkcie końcowym. Preferuj prywatne punkty końcowe Cloud Run dla wbudowanego uwierzytelniania i TLS. - Szyfrowanie: Pub/Sub szyfruje dane w tranzycie i w spoczynku; używaj CMEK na tematach dla kluczy zarządzanych przez klienta (customer-managed keys). Zastosuj VPC Service Controls, aby zmniejszyć ryzyko eksfiltracji danych. W razie potrzeby użyj szyfrowania po stronie klienta dla wrażliwych pól w ładunku (payload).
- IAM: nadawaj role z najmniejszymi uprawnieniami (
- Diagnostyka operacyjna opóźnień, ponownych dostarczeń i awarii subskrybentów
- Monitoruj za pomocą Cloud Monitoring:
subscription/num_undelivered_messagesioldest_unacked_message_agedla zaległości.expired_ack_deadline_countdo wykrywania pominiętych potwierdzeń powodujących duplikaty.publish_request_countipull_request_countdla przepustowości.
- Badaj brakujące zdarzenia na dashboardzie, odtwarzając znany zbiór danych w potoku i porównując wyniki na poszczególnych etapach, aby wyizolować wadliwą transformację lub ujście (sink).
- Dla streamingu w Dataflow:
- Używaj autoskalowania z odpowiednią wartością
maxWorkers, aby absorbować obciążenie z wielu źródeł. - Używaj operacji
drainna potokach przy niekompatybilnych aktualizacjach, aby pozwolić na ukończenie przetwarzania w toku i zapobiec utracie danych.
- Używaj autoskalowania z odpowiednią wartością
- Dla powiadomień o wstawianiu do BigQuery, przekierowuj wpisy audytowe z Cloud Logging za pomocą ujścia (sink) przefiltrowanego do określonych tabel do tematu Pub/Sub w celu alertowania.
- Monitoruj za pomocą Cloud Monitoring:
Praktyczny scenariusz problemowy
Firma Contoso Freight potrzebuje globalnej platformy do obsługi zdarzeń w czasie rzeczywistym, która będzie pozyskiwać 10 000 komunikatów telemetrycznych IoT na minutę z ciężarówek, wzbogacać zdarzenia, zasilać interaktywne analizy i wyzwalać przepływy pracy po otrzymaniu plików od zewnętrznych partnerów. Niektóre pliki CSV od partnerów zawierają nieprawidłowo sformatowane wiersze, a zespół analityczny musi mieć możliwość inspekcji błędów bez blokowania strumienia.
- Stwórz podstawową warstwę wiadomości i schematu
- Działanie: Zdefiniuj schemat Avro dla danych telemetrycznych i dołącz go do tematu Pub/Sub o nazwie
telemetryz włączonym wymuszaniem schematu (schema enforcementustawionym narequire). Włącz porządkowanie wiadomości i publikuj zordering_key = hash(device_id). - Uzasadnienie: Wymuszanie schematu na poziomie tematu odrzuca nieprawidłowo sformatowane zdarzenia na wczesnym etapie. Porządkowanie per urządzenie wspiera uporządkowane przetwarzanie, gdy jest to potrzebne, a haszowanie rozprasza klucze, aby utrzymać przepustowość.
- Utwórz subskrypcje z izolacją i obsługą dead-lettering
- Działanie: Utwórz subskrypcję pull
telemetry-stream-subdla Dataflow z tematem dead-lettertelemetry-dltimax_delivery_attempts=10. Dodaj subskrypcję BigQuerytelemetry-raw-bq, aby zapisywać surowe zdarzenia w tabeli partycjonowanej czasowo w celu śledzenia pochodzenia danych (lineage) i ponownego odtwarzania. - Uzasadnienie: DLQ (Dead-Letter Queue) izoluje problematyczne wiadomości (
poison messages) w celu ich zbadania. Osobna subskrypcja BigQuery zapewnia ścieżkę eksportu o niskich wymaganiach operacyjnych do przechowywania surowych zdarzeń, niezależnie od potoku przetwarzania.
- Zbuduj potok strumieniowy Dataflow do wzbogacania danych i obsługi ujść (sinks)
- Działanie: Pobieraj dane z
telemetry-stream-subza pomocąstreaming pullz kontrolą przepływu. Waliduj dane względem schematu, wzbogacaj je danymi referencyjnymi i obliczaj agregacje okienkowe. Zapisuj do BigQuery za pomocą Storage Write API z nazwanym strumieniem iinsertId = event_id; zapisuj surowe kopie zapasowe do Cloud Storage co godzinę; przekierowuj błędne/nieprzetworzone rekordy do tabeli dead-letter w BigQuery. - Uzasadnienie: Storage Write API zapewnia zapisy o wysokiej przepustowości i niskim opóźnieniu z idempotencją dzięki
insertId/offsetom strumienia. Tabela dead-letter umożliwia inspekcję bez blokowania strumienia, a archiwa w Cloud Storage pozwalają na ponowne odtworzenie danych.
- Obsłuż duplikaty i ostateczną spójność w analityce
- Działanie: Dla zapytań interaktywnych, które muszą wykluczać duplikaty, publikuj
event_idievent_timew każdym rekordzie i użyj widoku deduplikującego: CREATE OR REPLACE VIEW analytics.latest_events AS SELECT * EXCEPT(rn) FROM ( SELECT e.*, ROW_NUMBER() OVER (PARTITION BY event_id ORDER BY event_time DESC) rn FROM analytics.events e ) WHERE rn = 1; Wprowadź krótkie opóźnienie w zapytaniach, oparte na obserwowanej dostępności streamingu BigQuery (np. dwukrotność mediany opóźnienia). - Uzasadnienie: Dostarczanie typu
at-least-oncewymaga idempotentnych zapisów i deduplikacji w czasie zapytania. Oczekiwanie zmniejsza ryzyko pominięcia danych w locie (in-flight) z powodu opóźnienia w widoczności danych strumieniowych.
- Zintegruj odbieranie plików od partnerów za pomocą Eventarc
- Działanie: Skonfiguruj Eventarc, aby przekierowywał zdarzenia
object.finalizedz Cloud Storage dla bucketapartner-dropsdo usługi Cloud Run, która uruchamia zadanie wsadowe (batch) Dataflow do ładowania plików CSV do BigQuery, wysyłając błędy parsowania do tabeli dead-letter. - Uzasadnienie: Eventarc zapewnia orkiestrację opartą na zdarzeniach z filtrowaniem CloudEvents na podstawie bucketa i prefiksu obiektu. Zadanie wsadowe Dataflow oddziela nieprawidłowo sformatowane wiersze do analizy, jednocześnie szybko ładując poprawne dane.
- Zabezpiecz platformę
- Działanie: Użyj oddzielnych kont serwisowych: producenci otrzymują rolę
pubsub.publisherna temacietelemetry; konto serwisowe workera Dataflow otrzymujepubsub.subscriberna subskrypcjitelemetry-stream-suboraz dostęp do zapisu w docelowych zbiorach danych BigQuery i Cloud Storage; wyzwalacz Eventarc używa dedykowanego SA z roląinvokerna usłudze Cloud Run. Włącz CMEK na temacietelemetryi zbiorach danych BigQuery. Skonfiguruj punkty końcowe push, jeśli istnieją, z OIDC i sprawdzaniemaudience. - Uzasadnienie: Zasada najmniejszych uprawnień (least-privilege) w IAM oraz CMEK spełniają wymagania bezpieczeństwa i zgodności; uwierzytelnione dostarczanie zapobiega podszywaniu się (spoofing).
- Zapewnij niezawodne działanie i skalowanie
- Działanie: Ustaw autoskalowanie Dataflow z dużą wartością
maxWorkers, aby absorbować szczytowe obciążenia
← Przetwarzanie strumieniowe z Dataflow i Apache Beam · Wszystkie domeny · Spark →
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 →