Google PDE: Spark, Dataproc e Elaborazione Distribuita dei Dati — Guida allo studio

Fa parte della Google Professional Data Engineer — Guida allo studio. Esercitati con risposte verificate nel centro esami Google, oppure fai test cronometrati su ExamRoll.io.

Panoramica

Apache Spark su Google Cloud Dataproc fornisce una piattaforma gestita ed elastica per l’elaborazione distribuita dei dati. È possibile scegliere tra cluster Dataproc a lunga esecuzione o effimeri e Dataproc Serverless per Spark, a seconda delle esigenze di controllo, della variabilità del runtime e dell’overhead di gestione. Spark offre astrazioni resilienti (RDD), API relazionali (DataFrame e Spark SQL) e un motore di esecuzione DAG fault-tolerant ottimizzato per ETL iterativi e batch su larga scala. Su Google Cloud, Cloud Storage sostituisce HDFS per uno storage durevole e a basso costo; il connettore BigQuery consente un offload analitico diretto; e Dataproc Metastore centralizza la gestione degli schemi. Le soluzioni efficaci allineano i cicli di vita dello storage e del calcolo, ottimizzano Spark per il carico di lavoro, strumentano l’osservabilità e applicano la sicurezza con il principio del minimo privilegio e l’isolamento di rete.

Architettura di Dataproc: Cluster, Serverless, Storage e Metastore

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

  - Utilizzare Parquet/ORC con column pruning e predicate pushdown. Gestire i file di piccole dimensioni tramite compattazione per raggiungere una dimensione target di 128–512 MiB per file per una scansione efficiente.
- Hive metastore
  - Centralizzare schemi e metadati delle tabelle in Dataproc Metastore (un Apache Hive Metastore gestito) o in un metastore basato su Cloud SQL per condividere i cataloghi tra più cluster.
  - Utilizzare tabelle esterne che puntano a GCS per la durabilità; partizionare per data/ora per limitare i costi di scansione.
- Job, inizializzazione e workflow
  - Inviare job spark, pyspark, spark-sql o hadoop. Le azioni di inizializzazione installano librerie o agenti aggiuntivi alla creazione del cluster (ad esempio, connettori, librerie Python).
  - I modelli di workflow parametrizzano pipeline multi-step; possono creare cluster effimeri per ogni workflow e poi eliminarli. Ciò migliora l'isolamento e riduce i costi di inattività.
  - I cluster effimeri sono consigliati per ETL batch; i dati e il metastore risiedono all'esterno del cluster (GCS, Dataproc Metastore, BigQuery).
- Integrazione con BigQuery
  - Il connettore Spark BigQuery legge/scrive direttamente su BigQuery; considerare la BigQuery Storage Read API per il throughput e la Write API per inserimenti in streaming a bassa latenza e con semantica "exactly-once".
  - Per la manutenzione delle tabelle, eseguire operazioni MERGE o di sovrascrittura delle partizioni a valle in BigQuery per finalizzare i caricamenti in modo atomico.
### Modello Spark, Ottimizzazione delle Prestazioni e Affidabilità

- API ed esecuzione
  - RDD: di basso livello, immutabili, type-safe in Scala/Java; si controllano il partizionamento e la persistenza.
  - DataFrame/Dataset: relazionali, ottimizzati con Catalyst; preferire questi per l'ETL grazie all'ottimizzazione delle query e alla generazione di codice.
  - Le trasformazioni sono lazy (map, filter, join); le azioni ne attivano l'esecuzione (count, collect, save). Spark costruisce un DAG di stage separati da shuffle; i task vengono eseguiti per ogni partizione.
- Partizionamento e shuffle
  - Partizionamento in input: un numero di partizioni sufficiente a utilizzare tutti i core; iniziare con un valore pari a 2-4 volte il numero totale di core degli executor. Controllare tramite `spark.default.parallelism` (per gli RDD) e le opzioni del reader (per i DataFrame).
  - Partizioni di shuffle: il valore predefinito di 200 è spesso insufficiente o eccessivo. Ottimizzare:
--conf spark.sql.shuffle.partitions= {total_executor_cores * 2 to 3}
```
      --conf spark.sql.autoBroadcastJoinThreshold=64m
      ```

    - Aggiungere un "salt" alle chiavi per le partizioni "hot"; applicare una pre-aggregazione lato map; filtrare i dati il prima possibile.
    - Abilitare l'Adaptive Query Execution (AQE) per unire le partizioni post-shuffle e gestire gli skewed join:
  --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
      ```

- Pattern di tolleranza ai guasti per ETL batch
  - Scritture idempotenti: scrivere su un percorso temporaneo/di staging, quindi promuovere atomicamente con un commit a livello di directory; per BigQuery, scrivere su una tabella di staging ed eseguire un `MERGE`:
MERGE target t USING staging s
ON t.id = s.id
WHEN MATCHED THEN UPDATE SET ...
WHEN NOT MATCHED THEN INSERT (...)
```

Sicurezza, Osservabilità e Costi

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

  - Dataproc invia in streaming i log del driver e di YARN a Cloud Logging; esportali verso i sink per la conservazione/analisi forense.
  - Monitora con le metriche di Cloud Monitoring: container YARN in attesa, CPU, memoria, stato di HDFS (se utilizzato), throughput di GCS. Imposta avvisi (alert) per tentativi di stage prolungati, perdita di executor e picchi di esecuzione speculativa.
  - Analisi dei fallimenti: le cause comuni includono straggler indotti da skew, OOM (Out of Memory) degli executor durante lo shuffle, fallimenti nel commit sull'object-store e perdita di nodi preemptible/spot. Aumenta il numero di tentativi (retry) con giudizio; tentativi eccessivi possono amplificare costi e ritardi.
- Ottimizzazione dei costi
  - Utilizza cluster effimeri o Dataproc Serverless per evitare costi di inattività; mantieni i dati in GCS per ridurre al minimo l'uso di dischi permanenti.
  - Aggiungi worker secondari preemptible/spot per assorbire i picchi di domanda; progetta per la rielaborazione, poiché i task sui nodi persi vengono ritentati. Non posizionare i nodi master su nodi preemptible.
  - Dimensiona correttamente i tipi di macchina (right-sizing) e usa l'autoscaling per ridurre la capacità quando le code sono vuote. Prediligi Parquet/ORC con partition pruning per ridurre i costi di scansione e l'uso della CPU.
  - Evita file di piccole dimensioni compattando gli output; un numero minore di file più grandi riduce l'overhead dei metadati e il tempo di esecuzione del job.
  - Per job brevi e periodici (ad esempio, un ETL Spark settimanale di 30 minuti), i worker preemptible o la modalità serverless offrono spesso il miglior profilo di costo.

#### Scenario Pratico

Acme Retail sta migrando un cluster Hadoop on-premise da 30 nodi che esegue ETL notturni con Spark e Hive per alimentare sistemi di analytics a valle. L'obiettivo è riutilizzare i job esistenti con modifiche minime, evitare la gestione a tempo pieno dei cluster, persistere i dati oltre il ciclo di vita del cluster e ridurre i costi di storage.

Approccio:
1) Spostare dati e metadati su servizi gestiti
   - Archivia tutti i dati grezzi e curati in Cloud Storage utilizzando Parquet con partizionamento (ad esempio, dt=AAAA-MM-GG).
   - Motivazione: GCS è durevole, a basso costo e disaccoppia il calcolo (compute) dallo storage, consentendo a cluster effimeri e job serverless di funzionare senza dischi permanenti. Il formato Parquet partizionato abilita il predicate pushdown e scansioni efficienti.

2) Centralizzare il catalogo con Dataproc Metastore
   - Migra l'Hive metastore su Dataproc Metastore. Crea tabelle Hive esterne che fanno riferimento ai percorsi GCS e mantieni la logica di schema/partizionamento esistente.
   - Motivazione: Un metastore gestito consente a più cluster effimeri e job serverless di condividere le definizioni delle tabelle senza dover eseguire un'istanza HA di MySQL/PostgreSQL.

3) Utilizzare cluster Dataproc effimeri per l'ETL batch e workflow template per l'orchestrazione
   - Definisci un workflow template che crea un cluster con l'immagine richiesta (ad esempio, 2.1-debian11), esegue i job Spark (spark-sql e pyspark) ed elimina il cluster al termine. Aggiungi azioni di inizializzazione per installare eventuali librerie personalizzate.
   - Motivazione: I cluster effimeri eliminano i costi di inattività e isolano le dipendenze dei job. I workflow template offrono ripetibilità e parametrizzazione (date, percorsi di input).

4) Abilitare l'autoscaling e i worker preemptible
   - Associa una policy di autoscaling con un piccolo gruppo di worker principali (core) e un pool più ampio di worker secondari preemptible; ottimizza i periodi di cooldown per ridurre rapidamente la capacità dopo l'esecuzione.
   - Motivazione: I worker principali mantengono la stabilità del cluster; i worker preemptible assorbono gli shuffle e le trasformazioni ampie (wide transformations) a un costo inferiore. I tentativi (retry) di Spark/YARN gestiscono i task persi a causa della preemption.

5) Integrare con BigQuery tramite lo Spark BigQuery connector
   - Per i caricamenti di dimensioni/fatti, scrivi i risultati di Spark in tabelle di staging di BigQuery, quindi esegui istruzioni MERGE per aggiornare le tabelle di destinazione in modo atomico. Dove la sovrascrittura diretta è sicura, scrivi su tabelle partizionate utilizzando la modalità di sovrascrittura della partizione (partition overwrite mode).
   - Motivazione: BigQuery serve analytics e BI su larga scala; la combinazione staging+MERGE produce upsert di tipo transazionale da job Spark batch, riducendo l'incoerenza a valle.

6) Ottimizzare Spark per prestazioni e affidabilità
   - Imposta le partizioni di shuffle in relazione ai core degli executor e abilita AQE:
 --conf spark.sql.shuffle.partitions=600
 --conf spark.sql.adaptive.enabled=true
 ```
  1. Rafforzare la sicurezza e la rete

    • Esegui i cluster con service account dedicati, concedendo solo i ruoli necessari per i percorsi GCS, il metastore e i set di dati BigQuery. Crea cluster con IP privato in una subnet ristretta con Accesso Privato Google e limita l’accesso alla UI tramite regole di firewall.
    • Motivazione: Il principio del privilegio minimo e l’isolamento della rete riducono la superficie di attacco; l’egress privato del piano di controllo evita l’esposizione pubblica.
  2. Strumentare logging, cronologia e avvisi

    • Abilita i log degli eventi di Spark su GCS e distribuisci l’History Server; instrada i log del driver/YARN a Cloud Logging con una policy di conservazione. Aggiungi avvisi di Monitoring per container in attesa per lungo tempo, fallimenti ripetuti dei task o durata eccessiva del job.
    • Motivazione: I log centralizzati supportano l’analisi delle cause principali (root-cause analysis); gli avvisi proattivi rilevano precocemente skew, OOM o I/O degradato.
  3. Modernizzare selettivamente con Dataproc Serverless per carichi ad hoc e picchi elastici

    • Sposta i carichi di lavoro Spark SQL sporadici o esplorativi su Dataproc Serverless; mantieni le pipeline notturne su cluster effimeri finché non saranno completamente convalidate su serverless.
    • Motivazione: La modalità serverless elimina le operazioni di gestione del cluster e scala automaticamente, ideale per carichi imprevedibili; i workflow esistenti continuano a funzionare con modifiche minime al codice.
  4. Convalidare i committer per l’object-store e la gestione dei file di piccole dimensioni

    • Imposta l’algoritmo v2 di FileOutputCommitter e compatta gli output in file da 256–512 MiB tramite repartition/coalesce prima delle scritture.
    • Motivazione: Gli object store non dispongono di una ridenominazione atomica; i committer ottimizzati riducono l’overhead di copia/ridenominazione. La compattazione mitiga il problema dei file di piccole dimensioni, migliorando prestazioni e costi.

Questo design riutilizza i job Spark e Hive esistenti con un refactoring minimo, garantisce la durabilità dei dati in GCS, centralizza gli schemi, contiene il raggio d’azione di un’eventuale violazione della sicurezza (security blast radius), fornisce un’osservabilità robusta e ottimizza i costi attraverso cluster effimeri, autoscaling, capacità preemptible e l’uso mirato dell’esecuzione serverless.


Messaggistica · Tutti i domini · Ingestione

Esercitati su queste domande → · Pratica cronometrata su 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.

Supera l'esame →

Sfoglia Google →

Related guides

Accesso tutto incluso

Un abbonamento. Ogni esame.

Ogni piano sblocca la ricerca illimitata di risposte, test pratici, spiegazioni AI e la libreria completa di risorse — in oltre 20 lingue.

Mensile
24.87
Just €0.83/day
Tutto incluso:
  • Ricerca risposte illimitata
  • Test pratici illimitati
  • Spiegazioni basate su AI
  • Libreria completa di risorse
  • Oltre 20 lingue
  • Aggiornamenti settimanali dei contenuti
  • Premi e referral
  • Supporto prioritario
Inizia la prova gratuita

Nessuna carta di credito richiesta*

Miglior valore
12 mesi
179.87
Just €0.49/daySave 40%
Tutto incluso:
  • Ricerca risposte illimitata
  • Test pratici illimitati
  • Spiegazioni basate su AI
  • Libreria completa di risorse
  • Oltre 20 lingue
  • Aggiornamenti settimanali dei contenuti
  • Premi e referral
  • Supporto prioritario
Inizia la prova gratuita

Nessuna carta di credito richiesta*

✓ Piano gratuito incluso · ✓ Annulla in qualsiasi momento · ✓ Tutti i piani sbloccano il prodotto completo