PySpark로 배우는 빅데이터 기초
Upendra Devisetty
Science Analyst, CyVerse


기본 RDD 변환
map(), filter(), flatMap(), union()
RDD = sc.parallelize([1,2,3,4])
RDD_map = RDD.map(lambda x: x * x)

RDD = sc.parallelize([1,2,3,4])
RDD_filter = RDD.filter(lambda x: x > 2)

RDD = sc.parallelize(["hello world", "how are you"])
RDD_flatmap = RDD.flatMap(lambda x: x.split(" "))

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)
계산을 실행하여 RDD에서 값을 반환합니다.
기본 RDD 액션
collect()
take(N)
first()
count()
collect()는 데이터셋의 모든 요소를 배열로 반환합니다.
take(N)은 앞의 N개 요소를 배열로 반환합니다.
RDD_map.collect()
[1, 4, 9, 16]
RDD_map.take(2)
[1, 4]
RDD_map.first()
[1]
RDD_flatmap.count()
5
PySpark로 배우는 빅데이터 기초