PySpark의 RDD 연산

PySpark로 배우는 빅데이터 기초

Upendra Devisetty

Science Analyst, CyVerse

PySpark 연산 개요

  • 변환은 새로운 RDD를 생성합니다.
  • 액션은 RDD에서 계산을 수행합니다.
PySpark로 배우는 빅데이터 기초

RDD 변환

  • 변환은 지연 평가를 따릅니다.

  • 기본 RDD 변환

    • map(), filter(), flatMap(), union()
PySpark로 배우는 빅데이터 기초

map() 변환

  • map() 변환은 RDD의 모든 요소에 함수를 적용합니다.

map

RDD = sc.parallelize([1,2,3,4])
RDD_map = RDD.map(lambda x: x * x)
PySpark로 배우는 빅데이터 기초

filter() 변환

  • filter 변환은 조건을 통과한 요소만 가진 새 RDD를 반환합니다.

filter

RDD = sc.parallelize([1,2,3,4])
RDD_filter = RDD.filter(lambda x: x > 2)
PySpark로 배우는 빅데이터 기초

flatMap() 변환

  • flatMap() 변환은 원본 RDD의 각 요소에 대해 여러 값을 반환합니다.

RDD = sc.parallelize(["hello world", "how are you"])
RDD_flatmap = RDD.flatMap(lambda x: x.split(" "))
PySpark로 배우는 빅데이터 기초

union() 변환

inputRDD = sc.textFile("logs.txt")
errorRDD = inputRDD.filter(lambda x: "error" in x.split())
warningsRDD = inputRDD.filter(lambda x: "warnings" in x.split())
combinedRDD = errorRDD.union(warningsRDD)
PySpark로 배우는 빅데이터 기초

RDD 액션

  • 계산을 실행하여 RDD에서 값을 반환합니다.

  • 기본 RDD 액션

    • collect()

    • take(N)

    • first()

    • count()

PySpark로 배우는 빅데이터 기초

collect() 및 take() 액션

  • collect()는 데이터셋의 모든 요소를 배열로 반환합니다.

  • take(N)은 앞의 N개 요소를 배열로 반환합니다.

RDD_map.collect()
[1, 4, 9, 16]
RDD_map.take(2)
[1, 4]
PySpark로 배우는 빅데이터 기초

first() 및 count() 액션

  • first()는 RDD의 첫 번째 요소를 반환합니다.
RDD_map.first()
[1]
  • count()는 RDD의 요소 수를 반환합니다.
RDD_flatmap.count()
5
PySpark로 배우는 빅데이터 기초

RDD 연산 실습해 봅시다

PySpark로 배우는 빅데이터 기초

Preparing Video For Download...