Spark Scala中JOIN后同名前缀列的动态对比问题
解决Spark中Join后对比同名列并生成差异列的问题
看起来你在尝试对比Join后两张表的同名列,生成记录差异列名的字段,但代码里有几个关键问题导致没达到预期效果,我来一步步帮你修正:
原代码的核心问题
- 错误地用字符串对比列:你在
when条件里写的s"${prefixR}${x._1}"是字符串字面量,不是对r侧列的引用。Spark会把列值和这个字符串做对比,当然永远不相等,所以结果全是逗号。 - 差异行的获取逻辑不对:你先用
union-except-intersect获取对称差,再join的方式会丢失部分场景(比如两边都有同一个customer_id但列值不同的行),而且inner join会过滤掉单边存在的行,不符合对比所有差异的需求。 - concat方式不合理:用
reduce(concat(_,_))会把所有逗号拼接在一起,不管列是否有差异,无法准确收集差异列名。
修正后的实现方案
我们直接对两张表做关联(推荐用full_outer来覆盖所有差异场景),然后逐个对比同名列,把有差异的列名用逗号拼接起来:
import spark.implicits._ import org.apache.spark.sql.functions._ import org.apache.spark.sql.Column val dfCur = sc.parallelize(Seq( (1,"2019-01-01","2018-01-01",1,2), (7,"2019-01-01","2019-01-01",100,200), (3,"2019-01-01","2019-01-03",5,6) )).toDF("customer_id", "report_date", "date", "value_1", "value_2") val dfRaw = sc.parallelize(Seq( (2,"2019-01-01","2019-01-01",1,2), (7,"2019-01-01","2019-01-01",100,300), (3,"2019-01-01","2019-01-03",5,6) )).toDF("customer_id", "report_date", "date", "value_1", "value_2") // 定义关联键,根据业务需求调整,这里假设是customer_id+report_date+date val joinKeys = Seq("customer_id", "report_date", "date") // 获取需要对比的列:排除关联键,或者保留全部列(根据需求) val compareCols = dfCur.columns.filter(!joinKeys.contains(_)) // 生成每个列的差异判断表达式:列值不同则返回列名,否则返回null val diffExpressions = compareCols.map { colName => when(col(s"l.$colName") =!= col(s"r.$colName"), lit(colName)).otherwise(lit(null)) } // 用concat_ws把非null的差异列名用逗号拼接,空值则返回空字符串 val diffColumn = concat_ws(",", diffExpressions: _*) // 执行关联并添加差异列 val result = dfCur.as("l") .join(dfRaw.as("r"), joinKeys, "full_outer") // full_outer覆盖所有单边存在或双边值不同的行 .withColumn("XYZ2", diffColumn) result.show(false)
执行结果解析
运行后你会得到:
- customer_id=7的行:
value_2列值不同,所以XYZ2是value_2 - customer_id=2的行:仅在
dfRaw存在,l侧列全为null,所以XYZ2是value_1,value_2 - customer_id=1的行:仅在
dfCur存在,r侧列全为null,所以XYZ2是value_1,value_2 - customer_id=3的行:所有列值相同,
XYZ2为空字符串
可选调整
如果你只关心双边都存在但列值不同的行,可以把join类型改成inner;如果需要保留关联键的差异,可以把compareCols改成dfCur.columns(不排除关联键)。
内容的提问来源于stack exchange,提问作者Ged
相关产品推荐
相关产品推荐

