Spark中DataFrame对比咨询:基于df1与df2的技术问题
嘿,针对你提到的Spark里对比df1和df2这两个结构类似的DataFrame的需求,我整理了几种实用的方法,你可以根据具体场景来选:
1. 找出两个DataFrame的共同数据(交集)
如果想知道哪些行在两个DF里都存在,用intersect方法就行,它会返回去重后的共同行:
// Scala版本 val commonRows = df1.intersect(df2) commonRows.show()
# Python版本 common_rows = df1.intersect(df2) common_rows.show()
要是需要保留重复行(比如某个行在df1出现2次、df2出现3次,想返回2次),可以用Spark 2.4+支持的intersectAll替代intersect。
2. 找出两个DataFrame的差异行
这是最常用的对比需求,分两种实现方式:
方式A:用exceptAll(Spark 2.4+推荐)
exceptAll能精准找出某一个DF独有的行,还保留重复行的数量。我们可以分别找出df1和df2的独有行,再合并起来标记来源:
// Scala版本 import org.apache.spark.sql.functions.lit // df1有但df2没有的行 val df1Only = df1.exceptAll(df2) // df2有但df1没有的行 val df2Only = df2.exceptAll(df1) // 合并差异行并标记来源 val allDiff = df1Only.withColumn("source", lit("df1")) .union(df2Only.withColumn("source", lit("df2"))) allDiff.show()
# Python版本 from pyspark.sql.functions import lit df1_only = df1.exceptAll(df2) df2_only = df2.exceptAll(df1) all_diff = df1_only.withColumn("source", lit("df1")) \ .union(df2_only.withColumn("source", lit("df2"))) all_diff.show()
方式B:用全外关联(兼容低版本Spark)
如果你的Spark版本低于2.4,没有exceptAll,可以通过全外关联来实现:
// Scala版本 import org.apache.spark.sql.functions.lit val df1Marked = df1.withColumn("source", lit("df1")) val df2Marked = df2.withColumn("source", lit("df2")) // 全外关联后,过滤出两边不匹配的行 val diffRows = df1Marked.join(df2Marked, df1.columns, "full_outer") .where(df1Marked("WEEK").isNull || df2Marked("WEEK").isNull) diffRows.show()
3. 逐字段对比(定位具体差异字段)
如果想知道具体是哪些字段的值不一样,可以先通过唯一键(比如你的WEEK+DIM1+DIM2)关联两个DF,再逐字段对比标记差异:
// Scala版本 import org.apache.spark.sql.functions.{when, concat, lit, col} // 用唯一键关联两个DF val joined = df1.join(df2, Seq("WEEK", "DIM1", "DIM2"), "full_outer") // 生成每个字段的差异描述 val compared = joined .withColumn("T1_diff", when( (df1("T1").isNull && df2("T1").isNotNull) || (df1("T1").isNotNull && df2("T1").isNull) || df1("T1") =!= df2("T1"), concat(lit("df1: "), df1("T1"), lit(", df2: "), df2("T1")) )) .withColumn("T2_diff", when( (df1("T2").isNull && df2("T2").isNotNull) || (df1("T2").isNotNull && df2("T2").isNull) || df1("T2") =!= df2("T2"), concat(lit("df1: "), df1("T2"), lit(", df2: "), df2("T2")) )) // 过滤出有差异的行查看 compared.filter(col("T1_diff").isNotNull || col("T2_diff").isNotNull).show()
注意:这里特意处理了NULL值的情况,因为默认的!=不会把NULL和非NULL判定为差异,如果你不需要这个逻辑,可以简化判断条件。
4. 用测试工具做全量校验(适合单元测试)
如果是在做单元测试,想要快速验证两个DF是否完全一致,可以用第三方库spark-testing-base,它能直接对比DF的结构和数据,还会输出详细差异:
// Scala版本,需要先引入spark-testing-base依赖 import com.holdenkarau.spark.testing.DataFrameSuiteBase class DataFrameComparisonTest extends DataFrameSuiteBase { test("verify df1 equals df2") { // 严格校验结构和数据顺序 assertDataFrameEquals(df1, df2) // 忽略行顺序的校验 assertDataFrameApproximateEquals(df1, df2, 0.0) } }
内容的提问来源于stack exchange,提问作者Dark Shadows
相关产品推荐
相关产品推荐

