PySpark 聚合

PySpark 入门

Benjamin Schmidt

Data Engineer

PySpark SQL 聚合概览

  • 常见 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 入门

聚合中的数据类型处理

# 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 入门

使用 RDD 做聚合

# 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 聚合最佳实践

  • 提前过滤:先缩小数据再做聚合
  • 处理数据类型:确保数据干净且类型正确
  • 避免全量操作:减少使用 groupBy() 等全量操作
  • 选对接口:多数场景优先用 DataFrame,因其优化更好
  • 监控性能:用 explain() 查看执行计划并优化
PySpark 入门

核心要点

  • PySpark SQL 聚合:使用 SUM()AVERAGE() 等汇总数据
  • DataFrame 与 SQL:组合使用以灵活处理数据
  • 数据类型处理:解决聚合时的类型不匹配
  • RDD vs DataFrame:权衡取舍并选对工具
PySpark 入门

Passons à la pratique !

PySpark 入门

Preparing Video For Download...