Resilient Distributed Datasets ใน PySpark

PySpark เบื้องต้น

Benjamin Schmidt

Data Engineer

การประมวลผลแบบขนานใน PySpark คืออะไร?

  • กระจายข้อมูลและการประมวลผลโดยอัตโนมัติไปยังหลาย node ในคลัสเตอร์
  • ประมวลผลชุดข้อมูลขนาดใหญ่แบบกระจายข้ามหลาย node
  • Worker node ประมวลผลข้อมูลแบบขนาน แล้วรวมผลเมื่อเสร็จงาน
  • ประมวลผลได้เร็วขึ้นในระดับขนาดใหญ่ (เช่น กิกะไบต์หรือเทราไบต์)

การประมวลผลแบบขนาน

PySpark เบื้องต้น

ทำความเข้าใจ RDD

RDD หรือ Resilient Distributed Datasets:

  • คอลเลกชันข้อมูลแบบกระจายทั่วคลัสเตอร์ พร้อมกู้คืนอัตโนมัติเมื่อ node ล้มเหลว
  • เหมาะกับข้อมูลขนาดใหญ่
  • ไม่สามารถเปลี่ยนแปลงได้โดยตรง แต่แปลงได้ด้วยการดำเนินการอย่าง map() หรือ filter() และดึงผลลัพธ์หรือสร้าง RDD ด้วย collect() หรือ paralelize()
PySpark เบื้องต้น

การสร้าง RDD

# 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()
PySpark เบื้องต้น

การใช้ 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)
```    
PySpark เบื้องต้น

RDD vs DataFrame

DataFrames

  • ระดับสูง: ออกแบบมาเพื่อความสะดวกในการใช้งาน
  • คำสั่งคล้าย SQL: ทำงานกับคิวรีแบบ SQL และดำเนินการซับซ้อนด้วยโค้ดน้อยกว่า
  • ข้อมูล Schema: มีคอลัมน์และชนิดข้อมูลเหมือนตาราง SQL

RDD

  • ระดับต่ำ: ยืดหยุ่นกว่า แต่ต้องใช้โค้ดมากขึ้นสำหรับงานซับซ้อน
  • ความปลอดภัยของชนิดข้อมูล: รักษาชนิดข้อมูลได้ แต่ไม่มีการปรับแต่งประสิทธิภาพแบบ DataFrame
  • ไม่มี Schema: ทำงานกับข้อมูลเชิงโครงสร้างอย่าง SQL หรือเชิงสัมพันธ์ได้ยากกว่า
  • รองรับการขยายขนาดขนาดใหญ่
  • โค้ดยาวมากเมื่อเทียบกับ DataFrame และไม่เหมาะกับงานวิเคราะห์ข้อมูล
PySpark เบื้องต้น

ฟังก์ชันและเมธอดที่ใช้บ่อย

  • map(): ใช้ฟังก์ชัน (รวมถึง lambda function ที่เขียนเอง) กับข้อมูลทั้งชุด เช่น: rdd.map(map_function)
  • collect(): รวบรวมข้อมูลจากทั่วคลัสเตอร์ เช่น: rdd.collect()
PySpark เบื้องต้น

มาฝึกกันเถอะ!

PySpark เบื้องต้น

Preparing Video For Download...