如何正确合并两个Spark DataFrame并生成统计汇总表?
解决方案:合并Spark DataFrame实现新增/移除统计
首先先修正你原始代码里的一个小语法错误(addStuffDf的第三个元组引号缺失),正确的初始化代码应该是:
val addStuffDf = Seq( ("A", "2018-03-22", 5), ("A", "2018-03-24", 1), ("B", "2018-03-24", 3) ).toDF("user", "dt", "count") val removedStuffDf = Seq( ("C", "2018-03-25", 10), ("A", "2018-03-24", 5), ("B", "2018-03-25", 1) ).toDF("user", "dt", "count")
针对你的需求,我们有两种高效的实现方式:
方法一:使用Full Outer Join + 空值填充
这种方式逻辑直观,适合数据量不大的场景:
- 先给两个DataFrame的
count列重命名,分别对应added和removed - 按
user和dt做全外连接(full outer join),确保所有记录都被保留 - 把连接后出现的
null值填充为0,因为没有对应操作的记录意味着数量为0
代码实现:
// 重命名count列 val addRenamed = addStuffDf.withColumnRenamed("count", "added") val removeRenamed = removedStuffDf.withColumnRenamed("count", "removed") // 全外连接并填充空值 val resultDf = addRenamed.join(removeRenamed, Seq("user", "dt"), "full_outer") .na.fill(0, Seq("added", "removed"))
方法二:使用Union + Pivot(更灵活扩展)
如果后续可能新增其他操作类型(比如updated),这种方式扩展性更强:
- 给两个DataFrame添加一个标识列
operation,分别标记为added和removed - 将两个DataFrame合并(union)
- 按
user和dt分组,通过pivot将operation的不同值转为列,聚合count的总和 - 同样填充空值为0
代码实现:
import org.apache.spark.sql.functions.lit // 添加操作类型标识 val addWithType = addStuffDf.withColumn("operation", lit("added")) val removeWithType = removedStuffDf.withColumn("operation", lit("removed")) // 合并、分组透视并填充空值 val resultDf = addWithType.union(removeWithType) .groupBy("user", "dt") .pivot("operation", Seq("added", "removed")) // 指定列顺序可选,避免随机排序 .sum("count") .na.fill(0, Seq("added", "removed"))
两种方法最终都会生成你需要的结构,比如针对你提供的测试数据,结果会包含:
A, 2018-03-22, 5, 0A, 2018-03-24, 1, 5B, 2018-03-24, 3, 0B, 2018-03-25, 0, 1C, 2018-03-25, 0, 10
内容的提问来源于stack exchange,提问作者Vasiliy Galkin
相关产品推荐
相关产品推荐

