PySpark에서 Pair RDD 다루기

PySpark로 배우는 빅데이터 기초

Upendra Devisetty

Science Analyst, CyVerse

PySpark의 pair RDD 소개

  • 실제 데이터셋은 보통 키/값 쌍입니다

  • 각 행은 키이며 하나 이상의 값에 매핑됩니다

  • Pair RDD는 이런 데이터셋을 위한 특수 자료구조입니다

  • Pair RDD: 키는 식별자, 값은 데이터입니다

PySpark로 배우는 빅데이터 기초

pair RDD 생성하기

  • pair RDD를 만드는 두 가지 일반적 방법

    • 키-값 튜플 목록에서 생성
    • 일반 RDD에서 생성
  • 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]))
PySpark로 배우는 빅데이터 기초

pair RDD의 변환

  • 일반 변환은 모두 pair RDD에도 작동합니다

  • 개별 요소가 아니라 키-값 쌍에 작동하는 함수를 전달해야 합니다

  • pair RDD 변환 예

    • reduceByKey(func): 같은 키의 값을 합칩니다

    • groupByKey(): 같은 키의 값을 그룹화합니다

    • sortByKey(): 키로 정렬된 RDD를 반환합니다

    • join(): 키를 기준으로 두 pair RDD를 조인합니다

PySpark로 배우는 빅데이터 기초

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로 배우는 빅데이터 기초

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로 배우는 빅데이터 기초

groupByKey() 변환

  • groupByKey()는 pair RDD {{1}}에서 같은 키의 값을 모두 그룹화합니다
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로 배우는 빅데이터 기초

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로 배우는 빅데이터 기초

연습해 봅시다

PySpark로 배우는 빅데이터 기초

Preparing Video For Download...