Làm việc với Pair RDD trong PySpark

Nền tảng Big Data với PySpark

Upendra Devisetty

Science Analyst, CyVerse

Giới thiệu Pair RDD trong PySpark

  • Dữ liệu thực tế thường ở dạng cặp khóa/giá trị

  • Mỗi hàng là một khóa ánh xạ tới một hoặc nhiều giá trị

  • Pair RDD là cấu trúc dữ liệu chuyên xử lý dạng này

  • Pair RDD: Khóa là định danh, giá trị là dữ liệu

Nền tảng Big Data với PySpark

Tạo pair RDD

  • Hai cách phổ biến để tạo pair RDD

    • Từ danh sách tuple khóa-giá trị
    • Từ một RDD thường
  • Đưa dữ liệu về dạng khóa/giá trị cho 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]))
Nền tảng Big Data với PySpark

Biến đổi trên pair RDD

  • Mọi biến đổi thông thường đều áp dụng cho pair RDD

  • Cần truyền hàm thao tác trên cặp khóa-giá trị thay vì phần tử đơn lẻ

  • Ví dụ về biến đổi trên paired RDD

    • reduceByKey(func): Gộp giá trị cùng khóa

    • groupByKey(): Nhóm giá trị cùng khóa

    • sortByKey(): Trả về RDD sắp theo khóa

    • join(): Nối hai pair RDD theo khóa

Nền tảng Big Data với PySpark

Biến đổi reduceByKey()

  • Biến đổi reduceByKey() gộp các giá trị có cùng khóa

  • Chạy song song theo từng khóa trong tập dữ liệu

  • Đây là transformation, không phải action

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)]
Nền tảng Big Data với PySpark

Biến đổi sortByKey()

  • sortByKey() sắp xếp pair RDD theo khóa

  • Trả về RDD sắp xếp tăng dần hoặc giảm dần theo khóa

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')]
Nền tảng Big Data với PySpark

Biến đổi groupByKey()

  • groupByKey() gom tất cả giá trị có cùng khóa trong 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']
Nền tảng Big Data với PySpark

Biến đổi join()

  • Biến đổi join() nối hai pair RDD dựa trên khóa
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))]
Nền tảng Big Data với PySpark

Ayo berlatih!

Nền tảng Big Data với PySpark

Preparing Video For Download...