Arbeta med par-RDD:er i PySpark

Grunderna i Big Data med PySpark

Upendra Devisetty

Science Analyst, CyVerse

Introduktion till par-RDD:er i PySpark

  • Verkliga datamängder består ofta av nyckel/värde-par

  • Varje rad är en nyckel och mappas till ett eller flera värden

  • Par-RDD är en särskild datastruktur för den här typen av datamängder

  • Par-RDD: nyckeln är identifieraren och värdet är datat

Grunderna i Big Data med PySpark

Skapa par-RDD:er

  • Två vanliga sätt att skapa par-RDD:er

    • Från en lista med nyckel/värde-tupler
    • Från en vanlig RDD
  • Forma datat till nyckel/värde-format för par-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]))
Grunderna i Big Data med PySpark

Transformationer på par-RDD:er

  • Alla vanliga transformationer fungerar på par-RDD:er

  • Funktionerna måste operera på nyckel/värde-par, inte på enskilda element

  • Exempel på transformationer för par-RDD:er

    • reduceByKey(func): Kombinerar värden med samma nyckel

    • groupByKey(): Grupperar värden med samma nyckel

    • sortByKey(): Returnerar en RDD sorterad efter nyckel

    • join(): Slår ihop två par-RDD:er baserat på nyckel

Grunderna i Big Data med PySpark

Transformationen reduceByKey()

  • Transformationen reduceByKey() kombinerar värden med samma nyckel

  • Den kör parallella operationer för varje nyckel i datamängden

  • Det är en transformation, inte en 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)]
Grunderna i Big Data med PySpark

Transformationen sortByKey()

  • Operationen sortByKey() sorterar en par-RDD efter nyckel

  • Den returnerar en RDD sorterad i stigande eller fallande ordning

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')]
Grunderna i Big Data med PySpark

Transformationen groupByKey()

  • groupByKey() grupperar alla värden med samma nyckel i par-RDD:n
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']
Grunderna i Big Data med PySpark

Transformationen join()

  • Transformationen join() slår ihop två par-RDD:er baserat på nyckel
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))]
Grunderna i Big Data med PySpark

Nu kör vi en övning!

Grunderna i Big Data med PySpark

Preparing Video For Download...