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

Spark中行级DataFrame对比:两个DataFrame的行数与内容校验

嘿,我来帮你搞定这两个Spark DataFrame的行级对比和行数统计!下面分步骤给你具体实现方案:

第一步:快速获取行数统计

要知道两个DataFrame的行数差异,直接调用count()方法就能轻松搞定:

val sourceRowNum = source.count()
val targetRowNum = target.count()

println(s"Source DataFrame 行数:$sourceRowNum")
println(s"Target DataFrame 行数:$targetRowNum")

放到你的例子里,会输出Source DataFrame 行数:5和Target DataFrame 行数:11,一眼就能看到行数差。

第二步:行级差异深度对比

接下来我们分场景处理行级的差异:

场景1:找出仅在某一个DF中存在的行

用Spark的except方法可以直接取两个DF的差集,快速定位独有的行:

// 仅在Source里有、Target里没有的行
val onlyInSource = source.except(target)
println("仅在Source DataFrame中存在的行:")
onlyInSource.show()

// 仅在Target里有、Source里没有的行
val onlyInTarget = target.except(source)
println("仅在Target DataFrame中存在的行:")
onlyInTarget.show()

你的例子里,onlyInSource会是空的(因为Source的所有行Target都包含),而onlyInTarget会输出Target多出来的6行重复记录和那个mark1为80的Axx行。

场景2:对比同标识行的字段差异(假设name是唯一匹配键)

如果name是用来匹配行的唯一标识,我们可以通过join来对比具体字段的差异:

import org.apache.spark.sql.functions._

// 全外连接两个DF,对比mark1和mark2的差异
val fieldDiffDF = source.join(target, Seq("name"), "full_outer")
  // 生成mark1的差异描述
  .withColumn("mark1_diff", 
    when(source("mark1") =!= target("mark1"), 
      concat(lit("source: "), source("mark1"), lit(" | target: "), target("mark1"))
    ).otherwise(lit("无差异"))
  )
  // 生成mark2的差异描述
  .withColumn("mark2_diff", 
    when(source("mark2") =!= target("mark2"), 
      concat(lit("source: "), source("mark2"), lit(" | target: "), target("mark2"))
    ).otherwise(lit("无差异"))
  )
  // 过滤出有差异的行
  .filter(col("mark1_diff") =!= "无差异" || col("mark2_diff") =!= "无差异")

println("字段存在差异的行:")
fieldDiffDF.show()

这个代码会定位到Axx的mark1差异:source是70,target里有一条是80,结果里会清晰展示这个差异。

场景3:统计行的重复次数差异

如果需要看相同内容的行在两个DF里的出现次数差异,可以用分组统计后再对比:

// 统计Source中每行的出现次数
val sourceRowCounts = source.groupBy("name", "mark1", "mark2")
  .agg(count("*").alias("source_count"))

// 统计Target中每行的出现次数
val targetRowCounts = target.groupBy("name", "mark1", "mark2")
  .agg(count("*").alias("target_count"))

// 对比次数差异,空值填充为0
val countDiffDF = sourceRowCounts.join(targetRowCounts, Seq("name", "mark1", "mark2"), "full_outer")
  .na.fill(0, Seq("source_count", "target_count"))
  .filter(col("source_count") =!= col("target_count"))

println("行出现次数有差异的记录:")
countDiffDF.show()

运行后会看到:

  • Axx(70,75)在Source出现1次,Target出现2次
  • Axx(80,75)在Source出现0次,Target出现1次
    这些次数差异都会被清晰展示出来。

内容的提问来源于stack exchange,提问作者Teju Priya

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 08:21:23