Spark DataFrame实现:跳过当前行计数,汇总同Key下其余计数
我来帮你搞定这个需求,其实核心思路特别清晰——先算出每个id对应的所有occurences总和,再用这个总和减去当前行的occurences,就能得到排除当前item后的其他计数值总和了。下面我给你分步骤拆解实现过程:
1. 确认输入与分组结果
首先先回顾下你的输入数据和分组后的中间结果:
输入数据定义
val df = sc.parallelize(Seq( ("0","car1", "success"), ("0","car1", "success"), ("0","car3", "success"), ("0","car2", "success"), ("1","car1", "success"), ("1","car2", "success"), ("0","car3", "success") )).toDF("id", "item", "status")
分组得到中间结果df2
val df2 = df.groupBy("id", "item").agg(count("item").alias("occurences"))
df2的展示结果如下:
+---+----+----------+
| id|item|occurences|
+---+----+----------+
| 0|car3| 2|
| 0|car2| 1|
| 0|car1| 2|
| 1|car2| 1|
| 1|car1| 1|
+---+----+----------+
2. 实现排除当前Item的计数值汇总
这里提供两种可行的实现方式,你可以根据数据规模选择:
方法1:分组聚合+关联(适合小数据量)
先单独计算每个id的总occurences,再和df2关联做减法:
import org.apache.spark.sql.functions._ // 第一步:计算每个id的总occurences val idTotalDF = df2.groupBy("id").agg(sum("occurences").alias("total_occurences")) // 第二步:关联df2和总计数表,计算排除当前item的总和 val resultDF = df2.join(idTotalDF, "id") .withColumn("other_occurences_sum", col("total_occurences") - col("occurences")) .select("id", "item", "occurences", "other_occurences_sum")
方法2:窗口函数(大数据量推荐,性能更优)
用窗口函数直接在df2上计算每个id的总occurences,避免关联操作:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.expressions.Window // 定义窗口:按id分组 val idWindow = Window.partitionBy("id") // 计算总计数并得到目标结果 val resultDF = df2.withColumn("total_occurences", sum("occurences").over(idWindow)) .withColumn("other_occurences_sum", col("total_occurences") - col("occurences")) .select("id", "item", "occurences", "other_occurences_sum")
3. 最终结果展示
两种方法得到的结果是一致的,如下所示:
+---+----+----------+---------------------+
| id|item|occurences|other_occurences_sum|
+---+----+----------+---------------------+
| 0|car3| 2| 3|
| 0|car2| 1| 4|
| 0|car1| 2| 3|
| 1|car2| 1| 1|
| 1|car1| 1| 1|
+---+----+----------+---------------------+
简单解释下:比如id=0的car3,总计数是2+1+2=5,减去自身的2,得到其他项的总和3,完全符合需求。
内容的提问来源于stack exchange,提问作者xem

