PySpark DataFrame分组计数后按天数计算均值的实现方法
解决PySpark计算分组计数日均的问题
嘿,这个需求其实很好实现,咱们分两种情况来处理,既能满足你硬编码天数的需求,也能适配动态计算天数的场景:
情况1:已知固定天数(比如你说的3天)
你已经通过groupBy('category').count()得到了每个分类的总计数,现在只需要对count列做除法运算,再根据你的目标结果取整就可以了。从你给出的目标结果来看,应该是向上取整(或者四舍五入也能得到同样结果),代码如下:
from pyspark.sql.functions import col, ceil, round # 先拿到你已经分组好的res DataFrame res = df.groupBy('category').count() # 方法1:向上取整(完全匹配你的目标结果) res_daily_avg = res.withColumn("count", ceil(col("count") / 3).cast("integer")) res_daily_avg.show() # 方法2:四舍五入到整数(和目标结果一致,适合需要近似值的场景) res_daily_avg = res.withColumn("count", round(col("count") / 3, 0).cast("integer")) res_daily_avg.show()
运行后就能得到你想要的结果:
+--------+-----+ |category|count| +--------+-----+ | cat2| 1| | cat3| 1| | cat1| 2| +--------+-----+
情况2:动态计算天数(更灵活,避免硬编码)
如果你的原始DataFrame里有日期列(比如叫date),建议不要硬写3,而是动态计算数据包含的总天数,这样后续数据天数变化时不用修改代码:
from pyspark.sql.functions import col, ceil, countDistinct # 先计算数据中的唯一天数 total_days = df.select(countDistinct("date")).first()[0] # 再计算每个分类的日均计数并取整 res = df.groupBy('category').count() res_daily_avg = res.withColumn("count", ceil(col("count") / total_days).cast("integer")) res_daily_avg.show()
这样不管数据是3天还是其他天数,都能自动适配啦~
内容的提问来源于stack exchange,提问作者User12345
相关产品推荐
相关产品推荐

