Работа с парными RDD в PySpark

Основы Big Data с PySpark

Upendra Devisetty

Science Analyst, CyVerse

Введение в парные RDD в PySpark

  • Реальные наборы данных обычно представлены в виде пар «ключ/значение»

  • Каждая строка — это ключ, которому соответствует одно или несколько значений

  • Парный RDD — специальная структура для работы с такими данными

  • Парный RDD: ключ — идентификатор, значение — данные

Основы Big Data с PySpark

Создание парных RDD

  • Два основных способа создания парных RDD

    • Из списка кортежей «ключ — значение»
    • Из обычного RDD
  • Привести данные к виду «ключ/значение» для парного 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

Преобразования парных RDD

  • Все обычные преобразования применимы к парным RDD

  • Функции должны работать с парами «ключ — значение», а не с отдельными элементами

  • Примеры преобразований парных RDD

    • reduceByKey(func): объединяет значения с одинаковым ключом

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

    • sortByKey(): возвращает RDD, отсортированный по ключу

    • join(): объединяет два парных 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() сортирует парный 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() группирует все значения с одинаковым ключом в парном 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() объединяет два парных 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...