Big Data Fundamentals with PySpark
Upendra Devisetty
Science Analyst, CyVerse


RDD Transformations พื้นฐาน
map(), filter(), flatMap() และ union()map() ใช้ฟังก์ชันกับทุก element ใน RDD
RDD = sc.parallelize([1,2,3,4])
RDD_map = RDD.map(lambda x: x * x)
filter() คืน RDD ใหม่ที่มีเฉพาะ element ที่ผ่านเงื่อนไข
RDD = sc.parallelize([1,2,3,4])
RDD_filter = RDD.filter(lambda x: x > 2)
flatMap() คืนค่าหลายค่าสำหรับแต่ละ element ใน RDD ต้นทาง
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)
Actions คือการดำเนินการที่คืนค่าหลังประมวลผลบน RDD
RDD Actions พื้นฐาน
collect()
take(N)
first()
count()
collect() คืน element ทั้งหมดของชุดข้อมูลในรูปแบบ array
take(N) คืน array ที่มี N element แรกของชุดข้อมูล
RDD_map.collect()
[1, 4, 9, 16]
RDD_map.take(2)
[1, 4]
first() แสดง element แรกของ RDDRDD_map.first()
[1]
count() คืนจำนวน element ทั้งหมดใน RDDRDD_flatmap.count()
5
Big Data Fundamentals with PySpark