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

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")
  )
逻辑说明
  1. 先通过宽松的join条件拿到所有可能的匹配结果,允许一行Marketplace数据同时匹配到自身和parent_id对应的两条RegionValues记录
  2. 用窗口函数给每个匹配结果打优先级标签,自身primary_id匹配的优先级高于parent_id匹配
  3. 取每个主键分组下优先级最高的一条记录,即可得到符合预期的结果,同时保留RegionValues中存在但Marketplace不存在的主键行,完全匹配你给出的期望输出。

内容的提问来源于stack exchange,提问作者219CID

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 03:36:03