PySpark เบื้องต้น
Benjamin Schmidt
Data Engineer

RDD หรือ Resilient Distributed Datasets:
map() หรือ filter() และดึงผลลัพธ์หรือสร้าง RDD ด้วย collect() หรือ paralelize()# Initialize a Spark session from pyspark.sql import SparkSession spark = SparkSession.builder.appName("RDDExample").getOrCreate()# Create a DataFrame from a csv census_df = spark.read.csv("/census.csv")# Convert DataFrame to RDD census_rdd = census_df.rdd# Show the RDD's contents using collect() census_rdd.collect()
# Collect the entire DataFrame into a local Python list of Row objects
data_collected = df.collect()
# Print the collected data
for row in data_collected:
print(row)
```
map(): ใช้ฟังก์ชัน (รวมถึง lambda function ที่เขียนเอง) กับข้อมูลทั้งชุด เช่น:
rdd.map(map_function)collect(): รวบรวมข้อมูลจากทั่วคลัสเตอร์ เช่น:
rdd.collect()PySpark เบื้องต้น