การทำงานกับ Pair RDDs ใน PySpark

Big Data Fundamentals with PySpark

Upendra Devisetty

Science Analyst, CyVerse

แนะนำ Pair RDDs ใน PySpark

  • ชุดข้อมูลในชีวิตจริงมักอยู่ในรูปแบบ key/value

  • แต่ละแถวคือ key ที่แมปกับค่าหนึ่งค่าขึ้นไป

  • Pair RDD คือโครงสร้างข้อมูลพิเศษสำหรับชุดข้อมูลประเภทนี้

  • Pair RDD: Key คือตัวระบุ และ value คือข้อมูล

Big Data Fundamentals with PySpark

การสร้าง Pair RDDs

  • วิธีสร้าง Pair RDD มี 2 แบบทั่วไป

    • จาก list ของ tuple แบบ key-value
    • จาก RDD ทั่วไป
  • จัดข้อมูลให้อยู่ในรูป key/value สำหรับ Pair 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 Fundamentals with PySpark

Transformations บน Pair RDDs

  • transformation ทั่วไปทุกแบบใช้กับ Pair RDD ได้

  • ต้องส่งฟังก์ชันที่ทำงานกับ key-value pair แทนการทำงานกับแต่ละ element

  • ตัวอย่าง transformation สำหรับ Pair RDD

    • reduceByKey(func): รวมค่าที่มี key เดียวกัน

    • groupByKey(): จัดกลุ่มค่าที่มี key เดียวกัน

    • sortByKey(): คืนค่า RDD ที่เรียงตาม key

    • join(): เชื่อม Pair RDD สองชุดตาม key

Big Data Fundamentals with PySpark

transformation reduceByKey()

  • transformation reduceByKey() รวมค่าที่มี key เดียวกัน

  • ทำงานแบบ parallel สำหรับแต่ละ key ในชุดข้อมูล

  • เป็น transformation ไม่ใช่ 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)]
Big Data Fundamentals with PySpark

transformation sortByKey()

  • operation sortByKey() เรียงลำดับ Pair RDD ตาม key

  • คืนค่า RDD ที่เรียงตาม key แบบจากน้อยไปมากหรือมากไปน้อย

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 Fundamentals with PySpark

transformation groupByKey()

  • groupByKey() จัดกลุ่มค่าทั้งหมดที่มี key เดียวกันใน 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 Fundamentals with PySpark

transformation join()

  • transformation join() เชื่อม Pair RDD สองชุดตาม key
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 Fundamentals with PySpark

มาฝึกกันเถอะ!

Big Data Fundamentals with PySpark

Preparing Video For Download...