Scala/Spark下DataFrame优先按primary_id关联无匹配降级按parent_id关联实现
问题原因
你原有代码的核心逻辑错误出在join条件的判断:(regionValuesDf.col("primary_id") === marketplaceDf.col("primary_id")).isNull 这个判断几乎不会触发。两个字段相等的判断返回Boolean类型,仅当两边字段存在null值时才会返回null,其余场景只会返回true/false,因此降级匹配parent_id的逻辑永远不会生效,自然出现0000000002的values为null的问题。
解决方案
不需要多次join,单次全外连接+窗口函数即可实现优先级匹配逻辑,代码如下:
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._ // 定义匹配条件:primary_id匹配 或 parent_id匹配 val joinCondition = regionValuesDf("primary_id") === marketplaceDf("primary_id") || regionValuesDf("primary_id") === marketplaceDf("parent_id") // 全外连接两个表,保留两边所有行 val joinedDf = regionValuesDf.join(marketplaceDf, joinCondition, "full_outer") // 定义窗口:按主键分组,按匹配优先级排序 val windowSpec = Window.partitionBy( coalesce(marketplaceDf("primary_id"), regionValuesDf("primary_id")) ).orderBy( // 自身primary_id匹配优先级最高,parent_id匹配次之 when(regionValuesDf("primary_id") === marketplaceDf("primary_id"), 1) .when(regionValuesDf("primary_id") === marketplaceDf("parent_id"), 2) .otherwise(3) ) // 取每个分组优先级最高的行,整理输出字段 val resultDf = joinedDf .withColumn("rn", row_number().over(windowSpec)) .filter(col("rn") === 1) .select( coalesce(marketplaceDf("primary_id"), regionValuesDf("primary_id")).alias("primary_id"), marketplaceDf("marketplace"), marketplaceDf("parent_id"), regionValuesDf("values") )
逻辑说明
- 先通过宽松的join条件拿到所有可能的匹配结果,允许一行Marketplace数据同时匹配到自身和parent_id对应的两条RegionValues记录
- 用窗口函数给每个匹配结果打优先级标签,自身primary_id匹配的优先级高于parent_id匹配
- 取每个主键分组下优先级最高的一条记录,即可得到符合预期的结果,同时保留RegionValues中存在但Marketplace不存在的主键行,完全匹配你给出的期望输出。
内容的提问来源于stack exchange,提问作者219CID
相关产品推荐
相关产品推荐

