การ aggregation ใน PySpark

PySpark เบื้องต้น

Benjamin Schmidt

Data Engineer

ภาพรวม PySpark SQL aggregations

  • การ aggregation ใน SQL ทั่วไปใช้งานได้กับ spark.sql()
    # SQL aggregation query
    spark.sql("""
      SELECT Department, SUM(Salary) AS Total_Salary, AVG(Salary) AS Average_Salary
      FROM employees
      GROUP BY Department
    """).show()
    
PySpark เบื้องต้น

การผสาน DataFrame และการดำเนินการ SQL

# Filter salaries over 3000
filtered_df = df.filter(df.Salary > 3000)

# Register filtered DataFrame as a view
filtered_df.createOrReplaceTempView("filtered_employees")

# Aggregate using SQL on the filtered view spark.sql(""" SELECT Department, COUNT(*) AS Employee_Count FROM filtered_employees GROUP BY Department """).show()
PySpark เบื้องต้น

การจัดการชนิดข้อมูลใน aggregations

# Example of type casting
data = [("HR", "3000"), ("IT", "4000"), ("Finance", "3500")]
columns = ["Department", "Salary"]
df = spark.createDataFrame(data, schema=columns)

# Convert Salary column to integer df = df.withColumn("Salary", df["Salary"].cast("int")) # Perform aggregation df.groupBy("Department").sum("Salary").show()
PySpark เบื้องต้น

การใช้ RDDs สำหรับ aggregations

# Example of aggregation with RDDs
rdd = df.rdd.map(lambda row: (row["Department"], row["Salary"]))

rdd_aggregated = rdd.reduceByKey(lambda x, y: x + y)
print(rdd_aggregated.collect())
PySpark เบื้องต้น

แนวทางปฏิบัติที่ดีสำหรับ PySpark aggregations

  • กรองข้อมูลก่อน: ลดขนาดข้อมูลก่อนทำ aggregation
  • จัดการชนิดข้อมูล: ตรวจสอบให้ข้อมูลสะอาดและมีชนิดที่ถูกต้อง
  • หลีกเลี่ยงการดำเนินการที่ใช้ข้อมูลทั้งชุด: ลดการใช้ groupBy() ให้น้อยที่สุด
  • เลือก interface ที่เหมาะสม: ใช้ DataFrames เป็นหลักเพราะมีการเพิ่มประสิทธิภาพในตัว
  • ติดตามประสิทธิภาพ: ใช้ explain() เพื่อตรวจสอบและปรับแผนการประมวลผล
PySpark เบื้องต้น

สรุปประเด็นสำคัญ

  • PySpark SQL Aggregations: ฟังก์ชันอย่าง SUM() และ AVERAGE() สำหรับสรุปข้อมูล
  • DataFrames และ SQL: การผสานทั้งสองแนวทางเพื่อความยืดหยุ่นในการจัดการข้อมูล
  • การจัดการชนิดข้อมูล: แก้ไขปัญหาชนิดข้อมูลไม่ตรงกันระหว่าง aggregation
  • RDDs กับ DataFrames: เข้าใจข้อดีข้อเสียและเลือกใช้ให้เหมาะสม
PySpark เบื้องต้น

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

PySpark เบื้องต้น

Preparing Video For Download...