Jeux de données distribués résilients dans PySpark

Introduction à PySpark

Benjamin Schmidt

Data Engineer

Qu'est-ce que la parallélisation dans PySpark ?

  • Parallélisation automatique des données et calculs sur plusieurs nœuds d'un grappe
  • Traitement distribué de grands ensembles de données sur plusieurs nœuds
  • Les nœuds travailleurs traitent en parallèle, puis combinent à la fin de la tâche
  • Traitement plus rapide à grande échelle (gigaoctets voire téraoctets)

Parallélisation

Introduction à PySpark

Comprendre les RDD

RDD ou « Resilient Distributed Datasets » :

  • Collections de données distribuées sur une grappe avec reprise automatique en cas de défaillance d'un nœud
  • Adaptés aux données à grande échelle
  • Immuables ; transformations avec map() ou filter() et actions comme collect() ou paralelize() pour obtenir des résultats ou créer des RDD
Introduction à PySpark

Créer un 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 à PySpark

Afficher 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 à PySpark

RDD vs DataFrames

DataFrames

  • Haut niveau : optimisés pour la simplicité
  • Opérations de type SQL : requêtes SQL et opérations complexes avec moins de code
  • Schéma : colonnes et types comme une table SQL

RDD

  • Bas niveau : plus souples, mais plus de code pour les opérations complexes
  • Typage préservé : conservent les types, sans les optimisations des DataFrames
  • Pas de schéma : moins adaptés aux données structurées (SQL, relationnel)
  • Excellente mise à l'échelle
  • Très verbeux comparés aux DataFrames et peu adaptés à l'analytique
Introduction à PySpark

Fonctions et méthodes utiles

  • map() : applique une fonction (y compris une lambda) à chaque élément : rdd.map(map_function)
  • collect() : récupère les données à travers la grappe : rdd.collect()
Introduction à PySpark

Passons à la pratique !

Introduction à PySpark

Preparing Video For Download...