Apache Spark按条件分类统计记录计数的扩展技术问题
Apache Spark:结合类型与子类型分组并按条件分类计数
看起来你已经有了基础的分组方案,现在需要进一步按自定义条件拆分类别来统计,我来给你演示一个实用的实现思路,基于你提供的sales数据集:
首先先回顾下你的数据集:
val sales = Seq( ("Warsaw", 2016, "facebook","share",100), ("Warsaw", 2017, "facebook","like",200), ("Boston", 2015,"twitter","share",50), ("Boston", 2016,"facebook","share",150), ("Toronto", 2017,"twitter","like",50) ).toDF("city", "year","media","action","amount")
假设我们想要的输出是按城市(city)和年份(year)分组,然后统计以下4个类别的总数量:
- Facebook总互动量(包含like和share)
- Twitter总互动量(包含like和share)
- 全平台分享总数(share)
- 全平台点赞总数(like)
实现代码
我们可以用groupBy结合agg里的sum(when(...))来实现条件计数/求和,这样能灵活定义各类别:
import org.apache.spark.sql.functions.{sum, when, col} val result = sales .groupBy("city", "year") .agg( sum(when(col("media") === "facebook", col("amount"))).alias("facebook_total"), sum(when(col("media") === "twitter", col("amount"))).alias("twitter_total"), sum(when(col("action") === "share", col("amount"))).alias("total_shares"), sum(when(col("action") === "like", col("amount"))).alias("total_likes") ) .orderBy("city", "year") result.show()
输出结果
运行后你会得到这样的结构化结果:
+-------+----+-------------+--------------+------------+-----------+ | city|year|facebook_total|twitter_total|total_shares|total_likes| +-------+----+-------------+--------------+------------+-----------+ | Boston|2015| null| 50| 50| null| | Boston|2016| 150| null| 150| null| |Toronto|2017| null| 50| null| 50| | Warsaw|2016| 100| null| 100| null| | Warsaw|2017| 200| null| null| 200| +-------+----+-------------+--------------+------------+-----------+
灵活调整类别
如果你的条件更复杂(比如要统计"Facebook的点赞数"这类子类型组合),只需要修改when里的条件即可,比如:
sum(when(col("media") === "facebook" && col("action") === "like", col("amount"))).alias("facebook_likes")
另外,如果想要把空值替换为0,可以在sum后面加.otherwise(0),比如:
sum(when(col("media") === "facebook", col("amount")).otherwise(0)).alias("facebook_total")
这样输出里的null就会变成0,更美观。
内容的提问来源于stack exchange,提问作者ds_user
相关产品推荐
相关产品推荐

