PySpark 入門
Benjamin Schmidt
Data Engineer
spark.sql() 執行# SQL 聚合查詢
spark.sql("""
SELECT Department, SUM(Salary) AS Total_Salary, AVG(Salary) AS Average_Salary
FROM employees
GROUP BY Department
""").show()
# 篩選薪資高於 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()
# 型別轉換範例 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()
# 使用 RDD 做聚合的範例 rdd = df.rdd.map(lambda row: (row["Department"], row["Salary"]))rdd_aggregated = rdd.reduceByKey(lambda x, y: x + y)print(rdd_aggregated.collect())
groupBy() 等全量操作explain() 檢視執行計畫並據此最佳化SUM()、AVERAGE() 等函式彙總資料PySpark 入門