Travailler avec les Pair RDD en PySpark

Principes de Big Data avec PySpark

Upendra Devisetty

Science Analyst, CyVerse

Introduction aux Pair RDD en PySpark

  • Les ensembles réels sont souvent en paires clé/valeur

  • Chaque ligne est une clé associée à une ou plusieurs valeurs

  • Un Pair RDD est une structure conçue pour ces ensembles

  • Pair RDD : la clé est l'identifiant et la valeur est la donnée

Principes de Big Data avec PySpark

Créer des Pair RDD

  • Deux façons courantes de créer des Pair RDD

    • À partir d'une liste de tuples clé-valeur
    • À partir d'un RDD ordinaire
  • Mettre les données au format clé/valeur pour un Pair RDD

my_tuple = [('Sam', 23), ('Mary', 34), ('Peter', 25)]
pairRDD_tuple = sc.parallelize(my_tuple)
my_list = ['Sam 23', 'Mary 34', 'Peter 25']
regularRDD = sc.parallelize(my_list)
pairRDD_RDD = regularRDD.map(lambda s: (s.split(' ')[0], s.split(' ')[1]))
Principes de Big Data avec PySpark

Transformations sur les Pair RDD

  • Toutes les transformations régulières fonctionnent sur les Pair RDD

  • Il faut passer des fonctions qui opèrent sur des paires clé-valeur, et non des éléments isolés

  • Exemples de transformations sur Pair RDD

    • reduceByKey(func) : combine les valeurs ayant la même clé

    • groupByKey() : regroupe les valeurs par clé

    • sortByKey() : retourne un RDD trié par clé

    • join() : joint deux Pair RDD selon la clé

Principes de Big Data avec PySpark

Transformation reduceByKey()

  • La transformation reduceByKey() combine les valeurs ayant la même clé

  • Elle exécute des opérations en parallèle pour chaque clé du jeu de données

  • C'est une transformation, pas une action

regularRDD = sc.parallelize([("Messi", 23), ("Ronaldo", 34), 
                             ("Neymar", 22), ("Messi", 24)])
pairRDD_reducebykey = regularRDD.reduceByKey(lambda x,y : x + y)
pairRDD_reducebykey.collect()

[('Neymar', 22), ('Ronaldo', 34), ('Messi', 47)]
Principes de Big Data avec PySpark

Transformation sortByKey()

  • L'opération sortByKey() trie un Pair RDD par clé

  • Elle retourne un RDD trié par clé en ordre croissant ou décroissant

pairRDD_reducebykey_rev = pairRDD_reducebykey.map(lambda x: (x[1], x[0]))
pairRDD_reducebykey_rev.sortByKey(ascending=False).collect()

[(47, 'Messi'), (34, 'Ronaldo'), (22, 'Neymar')]
Principes de Big Data avec PySpark

Transformation groupByKey()

  • groupByKey() regroupe toutes les valeurs ayant la même clé dans le Pair RDD
airports = [("US", "JFK"),("UK", "LHR"),("FR", "CDG"),("US", "SFO")]
regularRDD = sc.parallelize(airports)
pairRDD_group = regularRDD.groupByKey().collect()
for cont, air in pairRDD_group:
  print(cont, list(air))

FR ['CDG'] US ['JFK', 'SFO'] UK ['LHR']
Principes de Big Data avec PySpark

Transformation join()

  • La transformation join() associe deux Pair RDD selon leur clé
RDD1 = sc.parallelize([("Messi", 34),("Ronaldo", 32),("Neymar", 24)])
RDD2 = sc.parallelize([("Ronaldo", 80),("Neymar", 120),("Messi", 100)])
RDD1.join(RDD2).collect()

[('Neymar', (24, 120)), ('Ronaldo', (32, 80)), ('Messi', (34, 100))]
Principes de Big Data avec PySpark

Passons à la pratique !

Principes de Big Data avec PySpark

Preparing Video For Download...