在 PySpark 中使用 Pair RDD

使用 PySpark 的 Big Data 基礎

Upendra Devisetty

Science Analyst, CyVerse

PySpark 的 pair RDD 入門

  • 真實資料集通常是鍵/值配對

  • 每一列是鍵,對應一個或多個值

  • Pair RDD 是處理此類資料集的特殊結構

  • Pair RDD:鍵是識別子,值是資料

使用 PySpark 的 Big Data 基礎

建立 pair RDD

  • 建立 pair RDD 的兩種常見方式

    • 從鍵值 tuple 的清單
    • 從一般 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]))
使用 PySpark 的 Big Data 基礎

pair RDD 的轉換

  • 一般的轉換都能作用於 pair RDD

  • 需傳入對鍵值配對運作的函式,而非單一元素

  • paired RDD 的轉換範例

    • reduceByKey(func):合併同鍵的值

    • groupByKey():群組同鍵的值

    • sortByKey():回傳依鍵排序的 RDD

    • join():依鍵連接兩個 pair RDD

使用 PySpark 的 Big Data 基礎

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 的 Big Data 基礎

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')]
使用 PySpark 的 Big Data 基礎

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']
使用 PySpark 的 Big Data 基礎

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))]
使用 PySpark 的 Big Data 基礎

一起來練習吧!

使用 PySpark 的 Big Data 基礎

Preparing Video For Download...