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的列,确保连接条件正常生效。
解决方案
- 升级Spark版本:直接将Spark升级到3.2.0或更高版本(如你使用的3.3.0),即可避免该问题。
- 临时规避方案(无法升级时):
- 重命名衍生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)
- 重命名衍生Dataset的列,明确区分同源列:
内容的提问来源于stack exchange,提问作者Flamma
相关产品推荐
相关产品推荐

