Resilientní distribuované datové sady v PySparku

Introduction to PySpark

Benjamin Schmidt

Data Engineer

Co je paralelizace v PySparku?

  • Automatická paralelizace dat a výpočtů napříč uzly clusteru
  • Distribuované zpracování velkých datových sad
  • Pracovní uzly zpracovávají data paralelně a výsledky se spojí na konci
  • Rychlejší zpracování ve velkém měřítku (gigabajty až terabajty)

Paralelizace

Introduction to PySpark

Porozumění RDD

RDD neboli Resilient Distributed Datasets:

  • Distribuované kolekce dat v clusteru s automatickým obnovením po selhání uzlu
  • Vhodné pro velké objemy dat
  • Neměnné; transformují se pomocí operací jako map() nebo filter() a akcí jako collect() či paralelize()
Introduction to PySpark

Vytvoření RDD

# Initialize a Spark session
from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("RDDExample").getOrCreate()

# Create a DataFrame from a csv census_df = spark.read.csv("/census.csv")
# Convert DataFrame to RDD census_rdd = census_df.rdd
# Show the RDD's contents using collect() census_rdd.collect()
Introduction to PySpark

Ukázka metody collect

# Collect the entire DataFrame into a local Python list of Row objects
data_collected = df.collect()

# Print the collected data
for row in data_collected:
    print(row)
```    
Introduction to PySpark

RDD vs. DataFrames

DataFrames

  • Vysoká úroveň abstrakce: optimalizované pro snadné použití
  • SQL operace: práce s SQL dotazy a složité operace s méně kódem
  • Schéma: obsahují sloupce a datové typy jako tabulka v SQL

RDD

  • Nízká úroveň abstrakce: flexibilnější, ale vyžadují více kódu
  • Typová bezpečnost: zachovávají datové typy, bez optimalizací DataFrames
  • Bez schématu: obtížnější práce se strukturovanými daty
  • Vysoká škálovatelnost
  • Velmi rozsáhlý kód a nevhodné pro analytiku
Introduction to PySpark

Užitečné funkce a metody

  • map(): metoda aplikuje funkce (včetně vlastních, např. lambda) na datovou sadu: rdd.map(map_function)
  • collect(): shromažďuje data z celého clusteru: rdd.collect()
Introduction to PySpark

Pojďme si procvičit!

Introduction to PySpark

Preparing Video For Download...