PySpark使用groupBy结合case when按类别年份分组统计各月nums总和问题咨询
问题原因
你编写的代码存在两个核心错误:
- 聚合逻辑顺序颠倒:
when条件判断需要放在sum函数内部,逐行判断月份后再对符合条件的nums求和,而非先求和再判断月份 - 分组提取的年份字段未设置别名,最终结果不会显示规范的
year列名
修正后的代码实现
from pyspark.sql.functions import when, year, month, sum, col new_sdf = cat_sdf.groupBy( "category", year("date").alias("year") # 给提取的年份字段设置别名 ).agg( sum(when(month("date") == 1, col("nums")).otherwise(0)).alias("jan_num"), sum(when(month("date") == 2, col("nums")).otherwise(0)).alias("feb_num"), sum(when(month("date") == 3, col("nums")).otherwise(0)).alias("mar_num"), sum(when(month("date") == 4, col("nums")).otherwise(0)).alias("apr_num"), # 剩余5-11月逻辑同上,依次补全即可 sum(when(month("date") == 12, col("nums")).otherwise(0)).alias("dec_num") )
更简便的pivot实现方案
如果不想手动编写12个月的判断逻辑,可以直接用Spark的pivot行转列方法实现,代码更简洁稳妥:
from pyspark.sql.functions import year, month, sum # 先提取年份、月份字段 tmp_sdf = cat_sdf.withColumn("year", year("date"))\ .withColumn("month", month("date")) # 分组后指定12个月份做pivot,求和后重命名列 new_sdf = tmp_sdf.groupBy("category", "year")\ .pivot("month", [1,2,3,4,5,6,7,8,9,10,11,12]) # 指定固定月份列表,避免无数据的月份列缺失 .sum("nums")\ .na.fill(0) # 无数据的月份值补0 # 按顺序重命名列 .toDF("category", "year", "jan_num", "feb_num", "mar_num", "apr_num", "may_num", "jun_num", "jul_num", "aug_num", "sep_num", "oct_num", "nov_num", "dec_num")
内容的提问来源于stack exchange,提问作者zeniifaye
相关产品推荐
相关产品推荐

