Google PDE: Spark, Dataproc i rozproszone przetwarzanie danych — 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

Apache Spark w Google Cloud Dataproc dostarcza zarządzaną, elastyczną platformę do rozproszonego przetwarzania danych. Można wybierać między długo działającymi lub efemerycznymi klastrami Dataproc a Dataproc Serverless for Spark, w zależności od potrzeb w zakresie kontroli, zmienności czasu wykonania i narzutu związanego z zarządzaniem. Spark oferuje odporne abstrakcje (RDD), relacyjne API (DataFrames i Spark SQL) oraz odporny na błędy silnik wykonawczy DAG, zoptymalizowany pod kątem iteracyjnego i wsadowego ETL na dużą skalę. W Google Cloud, Cloud Storage zastępuje HDFS jako trwały i tani magazyn danych; konektor BigQuery umożliwia bezpośrednie odciążenie analityczne; a Dataproc Metastore centralizuje zarządzanie schematami. Efektywne rozwiązania dopasowują cykle życia pamięci masowej i zasobów obliczeniowych, dostrajają Sparka do obciążenia, wdrażają obserwowalność oraz stosują zabezpieczenia zgodnie z zasadą najmniejszych uprawnień i izolacji sieciowej.

Architektura Dataproc: Klastry, Serverless, Pamięć masowa i Magazyn metadanych

    --conf mapreduce.fileoutputcommitter.algorithm.version=2
    ```

  - Używaj formatów Parquet/ORC z przycinaniem kolumn (column pruning) i wpychaniem predykatów (predicate pushdown). Zarządzaj małymi plikami poprzez kompakcję, dążąc do rozmiaru 128–512 MiB na plik w celu zapewnienia wydajnego skanowania.
- Magazyn metadanych Hive (Hive metastore)
  - Centralizuj schematy i metadane tabel w Dataproc Metastore (zarządzany Apache Hive Metastore) lub w magazynie metadanych opartym na Cloud SQL, aby współdzielić katalogi między klastrami.
  - Używaj tabel zewnętrznych (external tables) wskazujących na GCS w celu zapewnienia trwałości; partycjonuj według daty/godziny, aby ograniczyć koszt skanowania.
- Zadania, inicjalizacja i przepływy pracy
  - Uruchamiaj zadania typu spark, pyspark, spark-sql lub hadoop. Akcje inicjalizacyjne (initialization actions) instalują dodatkowe biblioteki lub agentów podczas tworzenia klastra (np. konektory, biblioteki Pythona).
  - Szablony przepływów pracy (workflow templates) parametryzują wieloetapowe potoki; mogą one tworzyć efemeryczne klastry dla każdego przepływu, a następnie je usuwać. Poprawia to izolację i redukuje koszty bezczynności.
  - Klastry efemeryczne są zalecane dla wsadowego ETL; dane i magazyn metadanych znajdują się poza klastrem (w GCS, Dataproc Metastore, BigQuery).
- Integracja z BigQuery
  - Konektor Spark BigQuery odczytuje/zapisuje dane bezpośrednio w BigQuery; rozważ użycie BigQuery Storage Read API dla wysokiej przepustowości oraz Write API dla wstawiania strumieniowego z niższym opóźnieniem i semantyką „dokładnie raz” (exactly-once).
  - W celu utrzymania tabel, wykonuj operacje MERGE lub nadpisywanie partycji w BigQuery, aby atomowo finalizować procesy ładowania danych.
### Model Spark, dostrajanie wydajności i niezawodność

- API i wykonanie
  - RDD: niskopoziomowe, niezmienne, bezpieczne typowo w Scala/Java; kontrolujesz partycjonowanie i persystencję.
  - DataFrame/Dataset: relacyjne, zoptymalizowane przez Catalyst; preferowane dla ETL ze względu na optymalizację zapytań i generowanie kodu.
  - Transformacje są leniwe (map, filter, join); akcje wyzwalają wykonanie (count, collect, save). Spark buduje DAG etapów (stages) rozdzielonych przez operacje shuffle; zadania (tasks) uruchamiane są per partycja.
- Partycjonowanie i shuffle
  - Partycjonowanie wejściowe: wystarczająca liczba partycji, aby wykorzystać wszystkie rdzenie; zacznij od 2–4x całkowitej liczby rdzeni egzekutorów. Kontroluj za pomocą `spark.default.parallelism` (dla RDD) i opcji czytnika (dla DataFrame).
  - Partycje shuffle: domyślna wartość 200 często jest zbyt mała lub zbyt duża. Dostosuj:
--conf spark.sql.shuffle.partitions= {total_executor_cores * 2 to 3}
```
      --conf spark.sql.autoBroadcastJoinThreshold=64m
      ```

    - Dodawaj sól (salt) do kluczy dla gorących partycji (hot partitions); stosuj preagregację po stronie mapowania (map-side pre-aggregation); filtruj dane jak najwcześniej.
    - Włącz Adaptive Query Execution (AQE), aby łączyć partycje po operacji shuffle (post-shuffle partitions) i obsługiwać nierównomierne złączenia (skewed joins):
  --conf spark.sql.adaptive.enabled=true
  ```
      --conf spark.dynamicAllocation.enabled=true
      --conf spark.shuffle.service.enabled=true
      --conf spark.dynamicAllocation.minExecutors=0
      --conf spark.dynamicAllocation.maxExecutors=200
      ```

- Wzorce odporności na błędy dla wsadowego ETL
  - Idempotentne zapisy: zapisuj do tymczasowej ścieżki/ścieżki przejściowej (staging), a następnie atomowo przenieś za pomocą zatwierdzenia na poziomie katalogu; dla BigQuery, zapisz do tabeli przejściowej i użyj `MERGE`:
MERGE target t USING staging s
ON t.id = s.id
WHEN MATCHED THEN UPDATE SET ...
WHEN NOT MATCHED THEN INSERT (...)
```

Bezpieczeństwo, obserwowalność i koszty

    --conf spark.eventLog.enabled=true
    --conf spark.eventLog.dir=gs://bucket/spark-events/
    ```

  - Dataproc przesyła strumieniowo logi sterownika (driver) i YARN do Cloud Logging; eksportuj je do ujść (sinks) w celu przechowywania i analizy śledczej.
  - Monitoruj za pomocą metryk Cloud Monitoring: oczekujące kontenery YARN, użycie CPU, pamięć, stan HDFS (jeśli jest używany), przepustowość GCS. Ustawiaj alerty na przedłużające się ponowienia etapów (stage retries), utratę egzekutorów i nagłe wzrosty liczby wykonań spekulacyjnych.
  - Analiza awarii: typowe przyczyny to maruderzy (stragglers) spowodowani przez nierównomierne rozłożenie danych (skew), błędy OOM egzekutorów podczas tasowania (shuffle), niepowodzenia zatwierdzania (commit) w magazynie obiektów oraz utrata węzłów preemptible/spot. Zwiększaj liczbę ponowień rozważnie; nadmierne ponowienia mogą zwiększyć koszty i opóźnienia.
- Optymalizacja kosztów
  - Używaj klastrów efemerycznych lub Dataproc Serverless, aby uniknąć kosztów bezczynności; przechowuj dane w GCS, aby zminimalizować użycie dysków trwałych.
  - Dodawaj drugorzędne węzły robocze preemptible/spot, aby obsłużyć szczytowe zapotrzebowanie; projektuj z myślą o ponownym obliczaniu, ponieważ zadania na utraconych węzłach są ponawiane. Nie umieszczaj węzłów master na węzłach preemptible.
  - Dobieraj odpowiednie typy maszyn i używaj autoskalowania, aby zmniejszać zasoby, gdy kolejki są puste. Preferuj formaty Parquet/ORC z przycinaniem partycji (partition pruning), aby obniżyć koszty skanowania i zużycie CPU.
  - Unikaj małych plików poprzez kompaktowanie danych wyjściowych; mniejsza liczba większych plików zmniejsza narzut metadanych i czas wykonania zadania.
  - W przypadku krótkich, okresowych zadań (np. cotygodniowe 30-minutowe zadanie ETL w Spark), węzły robocze preemptible lub tryb serverless często zapewniają najlepszy profil kosztowy.

#### Scenariusz praktycznego problemu

Firma Acme Retail migruje 30-węzłowy klaster Hadoop on-premise, na którym działają nocne zadania ETL w Spark i Hive, zasilające dalsze systemy analityczne. Chcą ponownie wykorzystać istniejące zadania przy minimalnych zmianach, uniknąć zarządzania klastrami w pełnym wymiarze czasu, utrwalać dane poza cyklem życia klastra i zredukować koszty przechowywania.

Podejście:
1) Umieszczenie danych i metadanych w usługach zarządzanych
   - Przechowuj wszystkie surowe i przetworzone dane w Cloud Storage, używając formatu Parquet z partycjonowaniem (na przykład, dt=YYYY-MM-DD).
   - Uzasadnienie: GCS jest trwały, tani i oddziela warstwę obliczeniową od warstwy przechowywania, dzięki czemu klastry efemeryczne i zadania serverless mogą działać bez dysków trwałych. Partycjonowany Parquet umożliwia "predicate pushdown" i wydajne skanowanie.

2) Scentralizowanie katalogu za pomocą Dataproc Metastore
   - Zmigruj Hive metastore do Dataproc Metastore. Utwórz zewnętrzne tabele Hive odwołujące się do ścieżek w GCS i zachowaj istniejącą logikę schematów/partycji.
   - Uzasadnienie: Zarządzany metastore pozwala wielu klastrom efemerycznym i zadaniom serverless na współdzielenie definicji tabel bez konieczności utrzymywania instancji HA MySQL/PostgreSQL.

3) Użycie efemerycznych klastrów Dataproc do wsadowego ETL i szablonów przepływu pracy do orkiestracji
   - Zdefiniuj szablon przepływu pracy (workflow template), który tworzy klaster z wymaganym obrazem (np. 2.1-debian11), uruchamia zadania Spark (spark-sql i pyspark) i usuwa klaster po zakończeniu. Dodaj akcje inicjalizacyjne, aby zainstalować niestandardowe biblioteki.
   - Uzasadnienie: Klastry efemeryczne eliminują koszty bezczynności i izolują zależności zadań. Szablony przepływu pracy zapewniają powtarzalność i parametryzację (daty, ścieżki wejściowe).

4) Włączenie autoskalowania i węzłów roboczych preemptible
   - Dołącz politykę autoskalowania z małą grupą podstawowych węzłów roboczych (core workers) i większą pulą drugorzędnych węzłów roboczych preemptible; dostosuj okresy "cooldown", aby szybko skalować w dół po zakończeniu zadania.
   - Uzasadnienie: Podstawowe węzły robocze utrzymują stabilność klastra; węzły preemptible obsługują operacje "shuffle" i szerokie transformacje przy niższych kosztach. Mechanizmy ponawiania w Spark/YARN radzą sobie z utraconymi zadaniami w wyniku wywłaszczenia (preemption).

5) Integracja z BigQuery za pomocą konektora Spark BigQuery
   - W przypadku ładowania wymiarów/faktów, zapisuj wyniki Spark do tymczasowych tabel (staging) w BigQuery, a następnie wykonuj instrukcje MERGE, aby atomowo zaktualizować tabele docelowe. Tam, gdzie bezpośrednie nadpisanie jest bezpieczne, zapisuj do tabel partycjonowanych w trybie nadpisywania partycji (partition overwrite mode).
   - Uzasadnienie: BigQuery obsługuje analitykę i BI na dużą skalę; podejście staging+MERGE daje efekt podobny do transakcyjnych operacji "upsert" z wsadowego Sparka, redukując niespójności w systemach docelowych.

6) Dostrajanie Sparka pod kątem wydajności i niezawodności
   - Ustaw liczbę partycji "shuffle" proporcjonalnie do liczby rdzeni egzekutorów i włącz AQE:
 --conf spark.sql.shuffle.partitions=600
 --conf spark.sql.adaptive.enabled=true
 ```
  1. Wzmocnienie bezpieczeństwa i sieci

    • Uruchamiaj klastry na dedykowanych kontach serwisowych, nadając tylko role potrzebne do dostępu do ścieżek GCS, metastore i zbiorów danych BigQuery. Twórz klastry z prywatnymi adresami IP w ograniczonej podsieci z Private Google Access i ograniczaj dostęp do interfejsu użytkownika za pomocą reguł zapory sieciowej.
    • Uzasadnienie: Zasada najmniejszych uprawnień i izolacja sieciowa zmniejszają powierzchnię ataku; prywatny ruch wychodzący płaszczyzny sterowania (control-plane egress) pozwala uniknąć ekspozycji na sieć publiczną.
  2. Instrumentacja logowania, historii i alertów

    • Włącz logi zdarzeń Spark do GCS i wdróż History Server; kieruj logi sterownika (driver)/YARN do Cloud Logging z polityką retencji. Dodaj alerty w Monitoring dla długo oczekujących kontenerów, powtarzających się awarii zadań lub nadmiernego czasu trwania zadania.
    • Uzasadnienie: Scentralizowane logi wspierają analizę przyczyn źródłowych; proaktywne alerty pozwalają wcześnie wykryć nierównomierne rozłożenie danych (skew), błędy OOM lub pogorszoną wydajność I/O.
  3. Selektywna modernizacja z użyciem Dataproc Serverless dla zadań ad hoc i elastycznych skoków obciążenia

    • Przenieś sporadyczne lub eksploracyjne obciążenia Spark SQL do Dataproc Serverless; utrzymuj nocne potoki na klastrach efemerycznych, dopóki nie zostaną w pełni zweryfikowane w trybie serverless.
    • Uzasadnienie: Tryb Serverless eliminuje operacje związane z klastrem i skaluje się automatycznie, co jest idealne dla nieprzewidywalnych obciążeń; istniejące przepływy pracy działają dalej przy minimalnych zmianach w kodzie.
  4. Walidacja mechanizmów zatwierdzania (committers) w magazynie obiektów i zarządzanie małymi plikami

    • Ustaw algorytm FileOutputCommitter na v2 i kompaktuj pliki wyjściowe do rozmiaru 256–512 MiB na plik za pomocą repartition/coalesce przed zapisem.
    • Uzasadnienie: Magazyny obiektów nie mają atomowej operacji zmiany nazwy; zoptymalizowane mechanizmy zatwierdzania redukują narzut związany z kopiowaniem/zmianą nazwy. Kompaktowanie łagodzi problem małych plików, poprawiając wydajność i obniżając koszty.

Ten projekt pozwala na ponowne wykorzystanie istniejących zadań Spark i Hive przy minimalnym refaktoringu, zapewnia trwałość danych w GCS, centralizuje schematy, ogranicza promień rażenia w przypadku naruszenia bezpieczeństwa, zapewnia solidną obserwowalność i optymalizuje koszty poprzez klastry efemeryczne, autoskalowanie, zasoby preemptible oraz ukierunkowane wykorzystanie trybu serverless.


Przesyłanie komunikatów · Wszystkie domeny · Pozyskiwanie

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