You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何正确合并两个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 + 空值填充

这种方式逻辑直观,适合数据量不大的场景:

  1. 先给两个DataFrame的count列重命名,分别对应added和removed
  2. 按user和dt做全外连接(full outer join),确保所有记录都被保留
  3. 把连接后出现的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),这种方式扩展性更强:

  1. 给两个DataFrame添加一个标识列operation,分别标记为added和removed
  2. 将两个DataFrame合并(union)
  3. 按user和dt分组,通过pivot将operation的不同值转为列,聚合count的总和
  4. 同样填充空值为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, 0
  • A, 2018-03-24, 1, 5
  • B, 2018-03-24, 3, 0
  • B, 2018-03-25, 0, 1
  • C, 2018-03-25, 0, 10

内容的提问来源于stack exchange,提问作者Vasiliy Galkin

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.21 06:46:03