Miglioramenti delle prestazioni

Pulizia dei dati con PySpark

Mike Metzger

Data Engineering Consultant

Spiegare il piano di esecuzione di Spark

voter_df = df.select(df['VOTER NAME']).distinct()
voter_df.explain()
== Physical Plan ==
*(2) HashAggregate(keys=[VOTER NAME#15], functions=[])
+- Exchange hashpartitioning(VOTER NAME#15, 200)
   +- *(1) HashAggregate(keys=[VOTER NAME#15], functions=[])
      +- *(1) FileScan csv [VOTER NAME#15] Batched: false, Format: CSV, Location: 
      InMemoryFileIndex[file:/DallasCouncilVotes.csv.gz], 
      PartitionFilters: [], PushedFilters: [], 
      ReadSchema: struct<VOTER NAME:string>
Pulizia dei dati con PySpark

Cos’è lo shuffling?

Lo shuffling è lo spostamento dei dati tra i vari worker per completare un task

  • Nasconde la complessità all’utente
  • Può essere lento
  • Riduce il throughput complessivo
  • Spesso necessario, ma va minimizzato
Pulizia dei dati con PySpark

Come ridurre lo shuffling?

  • Limita l’uso di .repartition(num_partitions)
    • Preferisci .coalesce(num_partitions)
  • Attenzione quando usi .join()
  • Usa .broadcast()
  • Potresti non doverlo limitare
Pulizia dei dati con PySpark

Broadcasting

Broadcasting:

  • Fornisce una copia di un oggetto a ogni worker
  • Evita comunicazioni inutili tra nodi
  • Può velocizzare drasticamente le .join()

Usa il metodo .broadcast(<DataFrame>)

from pyspark.sql.functions import broadcast
combined_df = df_1.join(broadcast(df_2))
Pulizia dei dati con PySpark

Passiamo alla pratica!

Pulizia dei dati con PySpark

Preparing Video For Download...