Praca z parami RDD w PySpark

Podstawy Big Data z PySpark

Upendra Devisetty

Science Analyst, CyVerse

Wprowadzenie do par RDD w PySpark

  • Rzeczywiste zbiory danych zazwyczaj mają postać par klucz/wartość

  • Każdy wiersz stanowi klucz i odwzorowuje jedną lub więcej wartości

  • Para RDD to specjalna struktura danych do pracy z tego rodzaju zbiorami

  • Para RDD: klucz jest identyfikatorem, wartość – danymi

Podstawy Big Data z PySpark

Tworzenie par RDD

  • Dwa typowe sposoby tworzenia par RDD

    • Z listy krotek klucz-wartość
    • Ze zwykłego RDD
  • Dane należy przekształcić do postaci klucz/wartość

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]))
Podstawy Big Data z PySpark

Transformacje na parach RDD

  • Wszystkie zwykłe transformacje działają na parach RDD

  • Funkcje muszą operować na parach klucz-wartość, nie na pojedynczych elementach

  • Przykłady transformacji na parach RDD

    • reduceByKey(func): łączy wartości o tym samym kluczu

    • groupByKey(): grupuje wartości o tym samym kluczu

    • sortByKey(): zwraca RDD posortowane według klucza

    • join(): łączy dwie pary RDD na podstawie klucza

Podstawy Big Data z PySpark

Transformacja reduceByKey()

  • Transformacja reduceByKey() łączy wartości o tym samym kluczu

  • Wykonuje operacje równolegle dla każdego klucza w zbiorze danych

  • Jest transformacją, nie akcją

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)]
Podstawy Big Data z PySpark

Transformacja sortByKey()

  • Operacja sortByKey() porządkuje parę RDD według klucza

  • Zwraca RDD posortowane rosnąco lub malejąco według klucza

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')]
Podstawy Big Data z PySpark

Transformacja groupByKey()

  • groupByKey() grupuje wszystkie wartości o tym samym kluczu w parze 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']
Podstawy Big Data z PySpark

Transformacja join()

  • Transformacja join() łączy dwie pary RDD na podstawie klucza
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))]
Podstawy Big Data z PySpark

Czas na ćwiczenia!

Podstawy Big Data z PySpark

Preparing Video For Download...