Lucrul cu Pair RDD-uri în PySpark

Fundamentele Big Data cu PySpark

Upendra Devisetty

Science Analyst, CyVerse

Introducere în pair RDD-uri în PySpark

  • Seturile de date reale sunt de obicei perechi cheie/valoare

  • Fiecare rând este o cheie și corespunde uneia sau mai multor valori

  • Pair RDD este o structură specială pentru astfel de seturi de date

  • Pair RDD: cheia este identificatorul, valoarea este datele

Fundamentele Big Data cu PySpark

Crearea pair RDD-urilor

  • Două modalități comune de a crea pair RDD-uri

    • Dintr-o listă de tuple cheie-valoare
    • Dintr-un RDD obișnuit
  • Conversia datelor în formă cheie/valoare pentru 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]))
Fundamentele Big Data cu PySpark

Transformări pe pair RDD-uri

  • Toate transformările obișnuite funcționează pe pair RDD

  • Funcțiile trebuie să opereze pe perechi cheie-valoare, nu pe elemente individuale

  • Exemple de transformări pe pair RDD

    • reduceByKey(func): Combină valorile cu aceeași cheie

    • groupByKey(): Grupează valorile cu aceeași cheie

    • sortByKey(): Returnează un RDD sortat după cheie

    • join(): Unește două pair RDD-uri după cheie

Fundamentele Big Data cu PySpark

Transformarea reduceByKey()

  • Transformarea reduceByKey() combină valorile cu aceeași cheie

  • Execută operații paralele pentru fiecare cheie din set

  • Este o transformare, nu o acțiune

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)]
Fundamentele Big Data cu PySpark

Transformarea sortByKey()

  • Operația sortByKey() ordonează pair RDD-ul după cheie

  • Returnează un RDD sortat ascendent sau descendent după cheie

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')]
Fundamentele Big Data cu PySpark

Transformarea groupByKey()

  • groupByKey() grupează toate valorile cu aceeași cheie din 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']
Fundamentele Big Data cu PySpark

Transformarea join()

  • Transformarea join() unește două pair RDD-uri pe baza cheii
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))]
Fundamentele Big Data cu PySpark

Să exersăm!

Fundamentele Big Data cu PySpark

Preparing Video For Download...