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

Spark joinWith连接条件疑似失效?不同版本执行结果不一致

Spark不同版本joinWith结果差异问题

我使用如下Scala代码:

import org.apache.spark.sql._
import spark.implicits._

case class Conversions(
    date: String,
    currency: String,
    rate: String
)

val usdConversions: Dataset[Conversions] = Seq(
    Conversions("19990811", "EUR", "0.1"),
    Conversions("19990811", "KPN", "0.2"),
    Conversions("19990812", "EUR", "0.3"),
    Conversions("19990812", "KPN", "0.4"),
).toDS

val usdToEurConversions = usdConversions.where($"currency" === "EUR")

usdConversions.joinWith(usdToEurConversions, 
    usdConversions("date") === usdToEurConversions("date")
).show(false)

在Spark 3.1.1/3.1.2环境执行后得到笛卡尔积结果:

+--------------------+--------------------+
|                  _1|                  _2|
+--------------------+--------------------+
|{19990811, EUR, 0.1}|{19990811, EUR, 0.1}|
|{19990811, EUR, 0.1}|{19990812, EUR, 0.3}|
|{19990811, KPN, 0.2}|{19990811, EUR, 0.1}|
|{19990811, KPN, 0.2}|{19990812, EUR, 0.3}|
|{19990812, EUR, 0.3}|{19990811, EUR, 0.1}|
|{19990812, EUR, 0.3}|{19990812, EUR, 0.3}|
|{19990812, KPN, 0.4}|{19990811, EUR, 0.1}|
|{19990812, KPN, 0.4}|{19990812, EUR, 0.3}|
+--------------------+--------------------+

而在Spark 3.3.0环境执行得到预期的等值连接结果:

+--------------------+--------------------+
|_1                  |_2                  |
+--------------------+--------------------+
|{19990811, EUR, 0.1}|{19990811, EUR, 0.1}|
|{19990811, KPN, 0.2}|{19990811, EUR, 0.1}|
|{19990812, EUR, 0.3}|{19990812, EUR, 0.3}|
|{19990812, KPN, 0.4}|{19990812, EUR, 0.3}|
+--------------------+--------------------+

问题原因

这是Spark 3.1.x版本中的一个已知bug:当joinWith的两个Dataset源自同一个父Dataset(此处usdToEurConversions是usdConversions过滤得到的衍生Dataset)时,Spark的查询解析器无法正确识别连接条件中对同源列的引用,导致连接条件被忽略,最终执行了笛卡尔积连接。

该问题在Spark 3.2.0及后续版本中已被修复——Spark优化了Dataset lineage的解析逻辑,能够正确关联同源Dataset的列,确保连接条件正常生效。

解决方案

  1. 升级Spark版本:直接将Spark升级到3.2.0或更高版本(如你使用的3.3.0),即可避免该问题。
  2. 临时规避方案(无法升级时):
    • 重命名衍生Dataset的列,明确区分同源列:
      val usdToEurConversions = usdConversions.where($"currency" === "EUR")
        .withColumnRenamed("date", "eur_date")
      usdConversions.joinWith(usdToEurConversions, usdConversions("date") === usdToEurConversions("eur_date"))
        .show(false)
      
    • 将Dataset注册为临时视图,通过SQL查询重新构建衍生Dataset,打破原有的lineage关联:
      usdConversions.createOrReplaceTempView("conversions")
      val usdToEurConversions = spark.sql("SELECT date, currency, rate FROM conversions WHERE currency = 'EUR'")
        .as[Conversions]
      usdConversions.joinWith(usdToEurConversions, usdConversions("date") === usdToEurConversions("date"))
        .show(false)
      

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 03:27:34