在 PySpark 中处理键值 RDD

使用 PySpark 的大数据基础

Upendra Devisetty

Science Analyst, CyVerse

PySpark 中的成对 RDD 简介

  • 真实数据通常是键/值对

  • 每行是一个键,对应一个或多个值

  • 成对 RDD 是处理此类数据的特殊数据结构

  • 成对 RDD:键为标识符,值为数据

使用 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]))
使用 PySpark 的大数据基础

成对 RDD 的转换

  • 所有常规转换都可用于成对 RDD

  • 需传入对键值对而非单个元素操作的函数

  • 成对 RDD 转换示例

    • reduceByKey(func):合并相同键的值

    • groupByKey():分组相同键的值

    • sortByKey():按键排序返回 RDD

    • join():按键连接两个成对 RDD

使用 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)]
使用 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')]
使用 PySpark 的大数据基础

groupByKey() 转换

  • groupByKey() 将成对 RDD {{1}} 中相同键的所有值分组
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']
使用 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))]
使用 PySpark 的大数据基础

Passons à la pratique !

使用 PySpark 的大数据基础

Preparing Video For Download...