如何在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. 测试结果验证
执行后最终结果如下:
| NAME | COUNTOFNAME | COUNTOFATTENDANCE |
|---|---|---|
| Aayush | 2 | 2 |
| Alia | 2 | 2 |
错误代码问题分析
用户代码的核心错误是在agg方法中嵌套调用另一个agg,PySpark的agg需要一次性传入所有聚合表达式,不能链式嵌套调用,这会导致语法解析失败。
内容的提问来源于stack exchange,提问作者BigData Lover
相关产品推荐
相关产品推荐

