PySpark 聚合

PySpark 入門

Benjamin Schmidt

Data Engineer

PySpark SQL 聚合總覽

  • 常見 SQL 聚合可用 spark.sql() 執行
    # SQL 聚合查詢
    spark.sql("""
      SELECT Department, SUM(Salary) AS Total_Salary, AVG(Salary) AS Average_Salary
      FROM employees
      GROUP BY Department
    """).show()
    
PySpark 入門

結合 DataFrame 與 SQL 操作

# 篩選薪資高於 3000
filtered_df = df.filter(df.Salary > 3000)

# 將篩選後的 DataFrame 註冊為檢視
filtered_df.createOrReplaceTempView("filtered_employees")

# 於篩選後的檢視使用 SQL 做聚合 spark.sql(""" SELECT Department, COUNT(*) AS Employee_Count FROM filtered_employees GROUP BY Department """).show()
PySpark 入門

聚合中的資料型別處理

# 型別轉換範例
data = [("HR", "3000"), ("IT", "4000"), ("Finance", "3500")]
columns = ["Department", "Salary"]
df = spark.createDataFrame(data, schema=columns)

# 將 Salary 欄位轉為整數 df = df.withColumn("Salary", df["Salary"].cast("int")) # 執行聚合 df.groupBy("Department").sum("Salary").show()
PySpark 入門

用 RDD 進行聚合

# 使用 RDD 做聚合的範例
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 聚合最佳實務

  • 及早篩選:在聚合前先縮小資料量
  • 處理型別:確保資料乾淨且型別正確
  • 避免全量操作:盡量少用 groupBy() 等全量操作
  • 選對介面:多數情境優先用 DataFrame,享受最佳化
  • 監控效能:用 explain() 檢視執行計畫並據此最佳化
PySpark 入門

重點回顧

  • PySpark SQL 聚合:使用 SUM()AVERAGE() 等函式彙總資料
  • DataFrame 與 SQL:彈性結合兩種方式處理資料
  • 型別處理:在聚合時解決型別不相容問題
  • RDD 與 DataFrame:理解取捨,選擇合適工具
PySpark 入門

一起來練習吧!

PySpark 入門

Preparing Video For Download...