You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在PySpark中结合使用GroupBy、Having与Order By

问题:将SQL的GROUP BY/HAVING/ORDER BY逻辑转换为PySpark代码

需要把以下SQL查询转换为PySpark代码,实现分组、过滤聚合结果、排序,并将数据迁移到新DataFrame:

SELECT TABLE1.NAME, Count(TABLE1.NAME) AS COUNTOFNAME, 
Count(TABLE1.ATTENDANCE) AS COUNTOFATTENDANCE INTO SCHOOL_DATA_TABLE
FROM TABLE1
WHERE (((TABLE1.NAME) Is Not Null))
GROUP BY TABLE1.NAME
HAVING (((Count(TABLE1.NAME))>1) AND ((Count(TABLE1.ATTENDANCE))<>5))
ORDER BY Count(TABLE1.NAME) DESC;

用户尝试的代码执行失败,代码如下:

df2= df.select('NAME','ATTENDANCE')
df2=df2.groupBy('NAME').agg(count('NAME').alias('name1').agg(count('ATTENDANCE').alias('NEW_ATTENDANCE'))).filter((col('name1')>1) & (col('NEW_ATTENDANCE') !=5))

示例测试数据:

rdd = spark.sparkContext.parallelize([
    ('Aayush', 10),
    ('Aayush', 9),
    ('Shiva', 5 ),
    ('Alia', 6),
    ('Aayan', 11),
    ('Alia',9)])
df_1 = spark.createDataFrame(rdd, schema=['NAME','ATTENDANCE'])

正确的PySpark实现

1. 导入依赖函数

from pyspark.sql import functions as F
from pyspark.sql.functions import col

2. 分步实现逻辑

# 1. 过滤NAME不为空的行,选择需要的列
filtered_df = df_1.filter(col("NAME").isNotNull()).select("NAME", "ATTENDANCE")

# 2. 分组并执行两个聚合统计
aggregated_df = filtered_df.groupBy("NAME").agg(
    F.count("NAME").alias("COUNTOFNAME"),
    F.count("ATTENDANCE").alias("COUNTOFATTENDANCE")
)

# 3. 应用HAVING条件过滤聚合结果
having_filtered_df = aggregated_df.filter(
    (col("COUNTOFNAME") > 1) & (col("COUNTOFATTENDANCE") != 5)
)

# 4. 按聚合结果降序排序
final_df = having_filtered_df.orderBy(col("COUNTOFNAME").desc())

# 5. 对应SQL的INTO,将结果保存为表
final_df.write.saveAsTable("SCHOOL_DATA_TABLE")

3. 链式调用简化版

from pyspark.sql import functions as F

final_df = df_1.filter(col("NAME").isNotNull()) \
    .groupBy("NAME") \
    .agg(
        F.count("NAME").alias("COUNTOFNAME"),
        F.count("ATTENDANCE").alias("COUNTOFATTENDANCE")
    ) \
    .filter((col("COUNTOFNAME") > 1) & (col("COUNTOFATTENDANCE") != 5)) \
    .orderBy(F.col("COUNTOFNAME").desc())

# 保存结果表
final_df.write.saveAsTable("SCHOOL_DATA_TABLE")

4. 测试结果验证

执行后最终结果如下:

NAMECOUNTOFNAMECOUNTOFATTENDANCE
Aayush22
Alia22

错误代码问题分析

用户代码的核心错误是在agg方法中嵌套调用另一个agg,PySpark的agg需要一次性传入所有聚合表达式,不能链式嵌套调用,这会导致语法解析失败。


内容的提问来源于stack exchange,提问作者BigData Lover

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.11 19:30:41