【问题标题】:Why does Spark fail with "Detected cartesian product for INNER join between logical plans"?为什么 Spark 会因“检测到逻辑计划之间的 INNER 连接的笛卡尔积”而失败?
【发布时间】:2017-11-29 10:46:48
【问题描述】:

我正在使用 Spark 2.1.0

当我执行以下代码时,我从 Spark 收到错误消息。为什么?如何解决?

val i1 = Seq(("a", "string"), ("another", "string"), ("last", "one")).toDF("a", "b")
val i2 = Seq(("one", "string"), ("two", "strings")).toDF("a", "b")
val i1Idx = i1.withColumn("sourceId", lit(1))
val i2Idx = i2.withColumn("sourceId", lit(2))
val input = i1Idx.union(i2Idx)
val weights = Seq((1, 0.6), (2, 0.4)).toDF("sourceId", "weight")
weights.join(input, "sourceId").show

错误:

scala> weights.join(input, "sourceId").show
org.apache.spark.sql.AnalysisException: Detected cartesian product for INNER join between logical plans
Project [_1#34 AS sourceId#39, _2#35 AS weight#40]
+- Filter (((1 <=> _1#34) || (2 <=> _1#34)) && (_1#34 = 1))
   +- LocalRelation [_1#34, _2#35]
and
Union
:- Project [_1#0 AS a#5, _2#1 AS b#6]
:  +- LocalRelation [_1#0, _2#1]
+- Project [_1#10 AS a#15, _2#11 AS b#16]
   +- LocalRelation [_1#10, _2#11]
Join condition is missing or trivial.
Use the CROSS JOIN syntax to allow cartesian products between these relations.;
  at org.apache.spark.sql.catalyst.optimizer.CheckCartesianProducts$$anonfun$apply$19.applyOrElse(Optimizer.scala:1011)
  at org.apache.spark.sql.catalyst.optimizer.CheckCartesianProducts$$anonfun$apply$19.applyOrElse(Optimizer.scala:1008)
  at org.apache.spark.sql.catalyst.trees.TreeNode$$anonfun$3.apply(TreeNode.scala:288)
  at org.apache.spark.sql.catalyst.trees.TreeNode$$anonfun$3.apply(TreeNode.scala:288)
  at org.apache.spark.sql.catalyst.trees.CurrentOrigin$.withOrigin(TreeNode.scala:70)
  at org.apache.spark.sql.catalyst.trees.TreeNode.transformDown(TreeNode.scala:287)
  at org.apache.spark.sql.catalyst.trees.TreeNode$$anonfun$transformDown$1.apply(TreeNode.scala:293)
  at org.apache.spark.sql.catalyst.trees.TreeNode$$anonfun$transformDown$1.apply(TreeNode.scala:293)
  at org.apache.spark.sql.catalyst.trees.TreeNode$$anonfun$5.apply(TreeNode.scala:331)
  at org.apache.spark.sql.catalyst.trees.TreeNode.mapProductIterator(TreeNode.scala:188)
  at org.apache.spark.sql.catalyst.trees.TreeNode.transformChildren(TreeNode.scala:329)
  at org.apache.spark.sql.catalyst.trees.TreeNode.transformDown(TreeNode.scala:293)
  at org.apache.spark.sql.catalyst.trees.TreeNode$$anonfun$transformDown$1.apply(TreeNode.scala:293)
  at org.apache.spark.sql.catalyst.trees.TreeNode$$anonfun$transformDown$1.apply(TreeNode.scala:293)
  at org.apache.spark.sql.catalyst.trees.TreeNode$$anonfun$5.apply(TreeNode.scala:331)
  at org.apache.spark.sql.catalyst.trees.TreeNode.mapProductIterator(TreeNode.scala:188)
  at org.apache.spark.sql.catalyst.trees.TreeNode.transformChildren(TreeNode.scala:329)
  at org.apache.spark.sql.catalyst.trees.TreeNode.transformDown(TreeNode.scala:293)
  at org.apache.spark.sql.catalyst.trees.TreeNode.transform(TreeNode.scala:277)
  at org.apache.spark.sql.catalyst.optimizer.CheckCartesianProducts.apply(Optimizer.scala:1008)
  at org.apache.spark.sql.catalyst.optimizer.CheckCartesianProducts.apply(Optimizer.scala:993)
  at org.apache.spark.sql.catalyst.rules.RuleExecutor$$anonfun$execute$1$$anonfun$apply$1.apply(RuleExecutor.scala:85)
  at org.apache.spark.sql.catalyst.rules.RuleExecutor$$anonfun$execute$1$$anonfun$apply$1.apply(RuleExecutor.scala:82)
  at scala.collection.IndexedSeqOptimized$class.foldl(IndexedSeqOptimized.scala:57)
  at scala.collection.IndexedSeqOptimized$class.foldLeft(IndexedSeqOptimized.scala:66)
  at scala.collection.mutable.WrappedArray.foldLeft(WrappedArray.scala:35)
  at org.apache.spark.sql.catalyst.rules.RuleExecutor$$anonfun$execute$1.apply(RuleExecutor.scala:82)
  at org.apache.spark.sql.catalyst.rules.RuleExecutor$$anonfun$execute$1.apply(RuleExecutor.scala:74)
  at scala.collection.immutable.List.foreach(List.scala:381)
  at org.apache.spark.sql.catalyst.rules.RuleExecutor.execute(RuleExecutor.scala:74)
  at org.apache.spark.sql.execution.QueryExecution.optimizedPlan$lzycompute(QueryExecution.scala:73)
  at org.apache.spark.sql.execution.QueryExecution.optimizedPlan(QueryExecution.scala:73)
  at org.apache.spark.sql.execution.QueryExecution.sparkPlan$lzycompute(QueryExecution.scala:79)
  at org.apache.spark.sql.execution.QueryExecution.sparkPlan(QueryExecution.scala:75)
  at org.apache.spark.sql.execution.QueryExecution.executedPlan$lzycompute(QueryExecution.scala:84)
  at org.apache.spark.sql.execution.QueryExecution.executedPlan(QueryExecution.scala:84)
  at org.apache.spark.sql.Dataset.withTypedCallback(Dataset.scala:2791)
  at org.apache.spark.sql.Dataset.head(Dataset.scala:2112)
  at org.apache.spark.sql.Dataset.take(Dataset.scala:2327)
  at org.apache.spark.sql.Dataset.showString(Dataset.scala:248)
  at org.apache.spark.sql.Dataset.show(Dataset.scala:636)
  at org.apache.spark.sql.Dataset.show(Dataset.scala:595)
  at org.apache.spark.sql.Dataset.show(Dataset.scala:604)
  ... 48 elided

【问题讨论】:

  • @JacekLaskowski 我刚刚从官方网站下载了 Spark 2.1.0,它提出了同样的问题(在本地 shell 中)。 Spark 2.1.1 可以正常工作。
  • 已确认。我也可以用 2.1.0 重现它。是的,2.1.1 工作正常。 Scala 无关紧要,因为我使用的是使用 Scala 2.11.8 构建的官方版本(这就是我将其作为“噪音”删除的原因)。

标签: scala apache-spark apache-spark-sql


【解决方案1】:

开启flag后可以触发inner join

spark.conf.set("spark.sql.crossJoin.enabled", "true")

你也可以使用交叉连接。

weights.crossJoin(input)

或将别名设置为

weights.join(input, input("sourceId")===weights("sourceId"), "cross")

您可以找到更多关于 issue SPARK-6459 的信息,据说它已在 2.1.1 中修复

由于您已经使用了 2.1.1,因此该问题应该已经解决了。

【讨论】:

  • 这两种选择都不适合我……同样的例外。您指出的问题应该已经在我正在使用的 Spark 版本中修复。
  • 我在同一版本中遇到了同样的异常,并且别名对我有用。
  • 您要执行哪个连接?做一个特定的右、左、右外、左外等连接。
  • 我需要进行“内部”连接,但它不起作用。如果我指定“left_outer”没有问题
  • 我不需要交叉连接。我需要一个内部连接。
【解决方案2】:

tl;dr 升级到 Spark 2.1.1。这是 Spark 中已修复的问题。

(我真的希望我也可以向您展示在 2.1.1 中修复该问题的确切更改)

【讨论】:

  • 如果我们无法升级到 Spark 2.1.1,是否有解决方法?
  • 您知道问题是如何解决的吗?我有 spark 2.2,但我仍然遇到这个问题,即 spark 将常规连接误解为笛卡尔积。
  • 我还在 Spark 2.2 上遇到这个问题,在一个实际上没有笛卡尔积的连接上(或者一个应该简单推断的连接)。除非有记录在案的修复程序,否则更新不是一个很好的解决方案。
【解决方案3】:

对我来说:

Dataset<Row> ds1 = sparkSession.read().load("/tmp/data");
Dataset<Row> ds2 = ds1;
ds1.join(ds2, ds1.col("name").equalTo(ds2.col("name"))) // got "Detected cartesian product for INNER join between logical plans"

Dataset<Row> ds1 = sparkSession.read().load("/tmp/data");
Dataset<Row> ds2 = sparkSession.read().load("/tmp/data");
ds1.join(ds2, ds1.col("name").equalTo(ds2.col("name"))) // running properly without errors

我使用的是 Spark 2.1.0。

【讨论】:

    【解决方案4】:

    在 SPARK 版本中出现此错误: 2.3.0.cloudera3

    通过对数据帧进行别名来解决。

    例如将失败的数据帧重新分配给另一个数据帧,并将名称别名为另一个数据帧。

    val dataFrame = inDataFrame.alias("dataFrame")

    希望这会有所帮助。

    【讨论】:

      【解决方案5】:

      您可以在命令之上使用它 SET spark.sql.crossJoin.enabled = TRUE;

      如果您的查询很复杂,可以重新构造查询以获得更好的结果

      【讨论】:

      • 在这种情况下我们如何重构查询?
      猜你喜欢
      • 2023-03-08
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-11-14
      • 1970-01-01
      • 1970-01-01
      • 2013-11-16
      • 2011-11-11
      相关资源
      最近更新 更多