【问题标题】:Requirement failed: Reordering broke - Spark Cassandra要求失败:重新排序中断 - Spark Cassandra
【发布时间】:2019-09-18 23:40:02
【问题描述】:

代码 -

val configDetails2 = configDetails1
    .join(skuDetails, configDetails1.col("sku_num") === skuDetails.col("sku") &&
      configDetails1.col("ccn") === skuDetails.col("ccn"), "left_outer")
    .select(
      configDetails1.col("*"),
      skuDetails.col("part"),
      skuDetails.col("part_description"),
      skuDetails.col("part_qty"))
    .withColumn("item_name", when($"part".isNull, "DBNULL").otherwise($"part"))
    .withColumn("item_description", when($"part_description".isNull, "DBNULL").otherwise($"part_description"))
    .withColumn("item_qty", when($"part_qty".isNull, lit(0)).otherwise($"part_qty"))
    .drop("part", "part_description", "part_qty")

  val itemKey = configDetails2.select("item_name").rdd
  val itemMaster = itemKey
    .joinWithCassandraTable("dig_master", "item_master")
    .select("buyer", "cfg_name".as("cfg"), "item", "ms_name".as("scheduler")).map(_._2) 

错误 -

原因:java.lang.IllegalArgumentException:要求失败: 重新排序中断 ({ccn#98, sku_num#54, sku#223, part#224, ccn#243},ArrayBuffer(sku_num, ccn, sku, part, ccn)) 不是 ({ccn#98, ccn#222, sku_num#54, sku#223, part#224, ccn#243},ArrayBuffer(sku_num, ccn, sku, 部分, ccn, sku, 部分, ccn, sku_num, ccn, sku, 部分, ccn)))

  在 scala.Predef$.require(Predef.scala:224)  在 org.apache.spark.sql.cassandra.execution.DSEDirectJoinStrategy.apply(DSEDirectJoinStrategy.scala:69) 在 org.apache.spark.sql.catalyst.planning.QueryPlanner$$anonfun$1.apply(QueryPlanner.scala:62) 在 org.apache.spark.sql.catalyst.planning.QueryPlanner$$anonfun$1.apply(QueryPlanner.scala:62) 在 scala.collection.Iterator$$anon$12.nextCur(Iterator.scala:434) 在 scala.collection.Iterator$$anon$12.hasNext(Iterator.scala:440) 在 scala.collection.Iterator$$anon$12.hasNext(Iterator.scala:439) 在 org.apache.spark.sql.catalyst.planning.QueryPlanner.plan(QueryPlanner.scala:92) 在 org.apache.spark.sql.catalyst.planning.QueryPlanner$$anonfun$2$$anonfun$apply$2.apply(QueryPlanner.scala:77) 在 org.apache.spark.sql.catalyst.planning.QueryPlanner$$anonfun$2$$anonfun$apply$2.apply(QueryPlanner.scala:74) 在 scala.collection.TraversableOnce$$anonfun$foldLeft$1.apply(TraversableOnce.scala:157) 在 scala.collection.TraversableOnce$$anonfun$foldLeft$1.apply(TraversableOnce.scala:157) 在 scala.collection.Iterator$class.foreach(Ite​​rator.scala:893) scala.collection.AbstractIterator.foreach(Ite​​rator.scala:1336) 在 scala.collection.TraversableOnce$class.foldLeft(TraversableOnce.scala:157) 在 scala.collection.AbstractIterator.foldLeft(Iterator.scala:1336) 在 org.apache.spark.sql.catalyst.planning.QueryPlanner$$anonfun$2.apply(QueryPlanner.scala:74) 在 org.apache.spark.sql.catalyst.planning.QueryPlanner$$anonfun$2.apply(QueryPlanner.scala:66) 在 scala.collection.Iterator$$anon$12.nextCur(Iterator.scala:434) 在 scala.collection.Iterator$$anon$12.hasNext(Iterator.scala:440) 在 org.apache.spark.sql.catalyst.planning.QueryPlanner.plan(QueryPlanner.scala:92) 在 org.apache.spark.sql.catalyst.planning.QueryPlanner$$anonfun$2$$anonfun$apply$2.apply(QueryPlanner.scala:77) 在 org.apache.spark.sql.catalyst.planning.QueryPlanner$$anonfun$2$$anonfun$apply$2.apply(QueryPlanner.scala:74) 在 scala.collection.TraversableOnce$$anonfun$foldLeft$1.apply(TraversableOnce.scala:157) 在 scala.collection.TraversableOnce$$anonfun$foldLeft$1.apply(TraversableOnce.scala:157) 在 scala.collection.Iterator$class.foreach(Ite​​rator.scala:893) scala.collection.AbstractIterator.foreach(Ite​​rator.scala:1336) 在 scala.collection.TraversableOnce$class.foldLeft(TraversableOnce.scala:157) 在 scala.collection.AbstractIterator.foldLeft(Iterator.scala:1336) 在 org.apache.spark.sql.catalyst.planning.QueryPlanner$$anonfun$2.apply(QueryPlanner.scala:74) 在 org.apache.spark.sql.catalyst.planning.QueryPlanner$$anonfun$2.apply(QueryPlanner.scala:66) 在 scala.collection.Iterator$$anon$12.nextCur(Iterator.scala:434) 在 scala.collection.Iterator$$anon$12.hasNext(Iterator.scala:440) 在 org.apache.spark.sql.catalyst.planning.QueryPlanner.plan(QueryPlanner.scala:92) 在 org.apache.spark.sql.catalyst.planning.QueryPlanner$$anonfun$2$$anonfun$apply$2.apply(QueryPlanner.scala:77) 在 org.apache.spark.sql.catalyst.planning.QueryPlanner$$anonfun$2$$anonfun$apply$2.apply(QueryPlanner.scala:74) 在 scala.collection.TraversableOnce$$anonfun$foldLeft$1.apply(TraversableOnce.scala:157) 在 scala.collection.TraversableOnce$$anonfun$foldLeft$1.apply(TraversableOnce.scala:157) 在 scala.collection.Iterator$class.foreach(Ite​​rator.scala:893) scala.collection.AbstractIterator.foreach(Ite​​rator.scala:1336) 在 scala.collection.TraversableOnce$class.foldLeft(TraversableOnce.scala:157) 在 scala.collection.AbstractIterator.foldLeft(Iterator.scala:1336) 在 org.apache.spark.sql.catalyst.planning.QueryPlanner$$anonfun$2.apply(QueryPlanner.scala:74) 在 org.apache.spark.sql.catalyst.planning.QueryPlanner$$anonfun$2.apply(QueryPlanner.scala:66) 在 scala.collection.Iterator$$anon$12.nextCur(Iterator.scala:434) 在 scala.collection.Iterator$$anon$12.hasNext(Iterator.scala:440) 在 org.apache.spark.sql.catalyst.planning.QueryPlanner.plan(QueryPlanner.scala:92) 在 org.apache.spark.sql.execution.QueryExecution.sparkPlan$lzycompute(QueryExecution.scala:84) 在 org.apache.spark.sql.execution.QueryExecution.sparkPlan(QueryExecution.scala:80) 在 org.apache.spark.sql.execution.QueryExecution.executedPlan$lzycompute(QueryExecution.scala:89) 在 org.apache.spark.sql.execution.QueryExecution.executedPlan(QueryExecution.scala:89) 在 org.apache.spark.sql.execution.QueryExecution.toRdd$lzycompute(QueryExecution.scala:92) 在 org.apache.spark.sql.execution.QueryExecution.toRdd(QueryExecution.scala:92) 在 org.apache.spark.sql.Dataset.rdd$lzycompute(Dataset.scala:2590) 在 org.apache.spark.sql.Dataset.rdd(Dataset.scala:2587) core.CollabStandardConfig$.delayedEndpoint$core$CollabStandardConfig$1(CollabStandardConfig.scala:185)

无法找到此错误的具体引用。 任何帮助表示赞赏。

【问题讨论】:

  • 什么是 DSE 版本?
  • DSE 6.0.7 Spark 2.2.3.4

标签: scala apache-spark-sql datastax-enterprise cassandra-3.0


【解决方案1】:

您是否将 scala 版本 2.10 升级到 2.11?然后尝试以下选项,

 val itemKey = configDetails2.select("item_name").rdd
  val itemMaster = itemKey
    .joinWithCassandraTable("dig_master", "item_master")
    .select("buyer", "cfg_name".as("cfg"), "item", "ms_name".as("scheduler")).map(_._2) 

把上面的代码改成SQL join as data frame,而不是把Dataframe转成dataset。

【讨论】:

  • scala 版本是 2.11。
猜你喜欢
  • 2020-06-19
  • 2018-01-05
  • 1970-01-01
  • 2021-10-14
  • 2021-11-25
  • 2015-09-24
  • 1970-01-01
  • 2019-02-24
  • 2015-04-04
相关资源
最近更新 更多