Odporne rozproszone zbiory danych w PySpark

Wprowadzenie do PySpark

Benjamin Schmidt

Data Engineer

Czym jest zrównoleglanie w PySpark?

  • Automatyczne zrównoleglanie danych i obliczeń na wielu węzłach klastra
  • Rozproszone przetwarzanie dużych zbiorów danych na wielu węzłach
  • Węzły robocze przetwarzają dane równolegle, łącząc wyniki na końcu zadania
  • Szybsze przetwarzanie dużych zbiorów (gigabajty lub terabajty)

Zrównoleglanie

Wprowadzenie do PySpark

Zrozumienie RDD

RDD czyli Resilient Distributed Datasets:

  • Rozproszone kolekcje danych w klastrze z automatycznym odtwarzaniem po awarii węzła
  • Przydatne przy dużych zbiorach danych
  • Niemutowalne; można je przekształcać operacjami jak map() lub filter(), a akcjami jak collect() czy paralelize() pobierać wyniki lub tworzyć nowe RDD
Wprowadzenie do PySpark

Tworzenie 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()
Wprowadzenie do PySpark

Prezentacja 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)
```    
Wprowadzenie do PySpark

RDD a DataFrames

DataFrames

  • Wysokopoziomowe: zoptymalizowane pod kątem łatwości użycia
  • Operacje SQL: zapytania podobne do SQL, złożone operacje przy mniejszej ilości kodu
  • Schemat: kolumny i typy danych jak w tabeli SQL

RDD

  • Niskopoziomowe: bardziej elastyczne, ale wymagają więcej kodu
  • Bezpieczeństwo typów: zachowują typy danych, bez optymalizacji DataFrames
  • Brak schematu: trudniejsza praca z danymi strukturalnymi jak SQL
  • Duża skalowalność
  • Bardzo rozbudowane w porównaniu do DataFrames i słabe przy analizie danych
Wprowadzenie do PySpark

Przydatne funkcje i metody

  • map(): metoda stosująca funkcje (w tym wyrażenia lambda) do zbioru danych: rdd.map(map_function)
  • collect(): pobiera dane z klastra: rdd.collect()
Wprowadzenie do PySpark

Czas na ćwiczenia!

Wprowadzenie do PySpark

Preparing Video For Download...