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
相关产品推荐
相关产品推荐

