Робота з Pair RDD у PySpark

Основи Big Data з PySpark

Upendra Devisetty

Science Analyst, CyVerse

Вступ до pair RDD у PySpark

  • Реальні набори даних зазвичай мають пари ключ/значення

  • Кожен рядок — це ключ, що відображається на одне чи кілька значень

  • Pair RDD — спеціальна структура даних для таких наборів

  • Pair RDD: ключ — ідентифікатор, значення — дані

Основи Big Data з PySpark

Створення pair RDD

  • Два поширені способи створити pair RDD

    • Зі списку кортежів ключ-значення
    • Зі звичайного RDD
  • Перетворіть дані у формат ключ/значення для paired 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]))
Основи Big Data з PySpark

Перетворення для pair RDD

  • Усі звичайні перетворення працюють із pair RDD

  • Потрібно передавати функції, що оперують парами ключ-значення, а не окремими елементами

  • Приклади перетворень paired RDD

    • reduceByKey(func): обʼєднує значення з однаковим ключем

    • groupByKey(): групує значення з однаковим ключем

    • sortByKey(): повертає RDD, відсортований за ключем

    • join(): обʼєднує два pair RDD за ключем

Основи Big Data з PySpark

Перетворення reduceByKey()

  • Перетворення reduceByKey() обʼєднує значення з однаковим ключем

  • Воно виконує паралельні операції для кожного ключа в наборі даних

  • Це перетворення, а не дія

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)]
Основи Big Data з PySpark

Перетворення sortByKey()

  • Операція sortByKey() впорядковує pair RDD за ключем

  • Повертає RDD, відсортований за ключем за зростанням або спаданням

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')]
Основи Big Data з PySpark

Перетворення groupByKey()

  • groupByKey() групує всі значення з однаковим ключем у 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']
Основи Big Data з PySpark

Перетворення join()

  • Перетворення join() обʼєднує два pair RDD за їхнім ключем
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))]
Основи Big Data з PySpark

Давайте потренуємось!

Основи Big Data з PySpark

Preparing Video For Download...