Nền tảng Big Data với PySpark
Upendra Devisetty
Science Analyst, CyVerse


Các biến đổi RDD cơ bản
map(), filter(), flatMap(), và 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)
Các thao tác trả về giá trị sau khi tính toán trên RDD
Các RDD Action cơ bản
collect()
take(N)
first()
count()
collect() trả về toàn bộ phần tử của tập dữ liệu dưới dạng mảng
take(N) trả về mảng gồm N phần tử đầu tiên
RDD_map.collect()
[1, 4, 9, 16]
RDD_map.take(2)
[1, 4]
RDD_map.first()
[1]
RDD_flatmap.count()
5
Nền tảng Big Data với PySpark