Principes de Big Data avec PySpark
Upendra Devisetty
Science Analyst, CyVerse
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
Deux façons courantes de créer des Pair RDD
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]))
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é
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)]
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')]
groupByKey() regroupe toutes les valeurs ayant la même clé dans le Pair RDDairports = [("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']
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