PySpark में Pair RDDs के साथ काम करना

PySpark के साथ Big Data Fundamentals

Upendra Devisetty

Science Analyst, CyVerse

PySpark में pair RDDs का परिचय

  • वास्तविक जीवन के datasets आमतौर पर key/value pairs होते हैं

  • हर row एक key होती है और एक या अधिक values से मैप होती है

  • ऐसे datasets के लिए Pair RDD एक विशेष data structure है

  • Pair RDD: Key पहचानकर्ता है और value डेटा है

PySpark के साथ Big Data Fundamentals

pair RDDs बनाना

  • Pair RDDs बनाने के दो सामान्य तरीके

    • key-value tuple की list से
    • एक regular RDD से
  • Paired RDD के लिए डेटा को key/value रूप में लाएँ

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 Fundamentals

pair RDDs पर transformations

  • सभी regular transformations, pair RDD पर काम करती हैं

  • आपको ऐसी functions पास करने होंगे जो key-value pairs पर operate करें, न कि individual elements पर

  • Paired RDD transformations के उदाहरण

    • reduceByKey(func): समान key की values को combine करें

    • groupByKey(): समान key की values को group करें

    • sortByKey(): key के आधार पर sorted RDD लौटाए

    • join(): key के आधार पर दो pair RDDs को join करें

PySpark के साथ Big Data Fundamentals

reduceByKey() transformation

  • reduceByKey() transformation समान key की values को जोड़ता है

  • यह dataset में हर key के लिए parallel operations चलाता है

  • यह 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)]
PySpark के साथ Big Data Fundamentals

sortByKey() transformation

  • sortByKey() operation pair RDD को key के अनुसार क्रमबद्ध करता है

  • यह key के आधार पर ascending या descending order में sorted 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 Fundamentals

groupByKey() transformation

  • groupByKey() pair RDD में समान key की सभी values को group करता है
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 Fundamentals

join() transformation

  • join() transformation दो pair RDDs को उनकी 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))]
PySpark के साथ Big Data Fundamentals

अभ्यास करते हैं

PySpark के साथ Big Data Fundamentals

Preparing Video For Download...