【问题标题】:Spark Cassandra Connector: for comprehension error (type mismatch)Spark Cassandra 连接器:用于理解错误(类型不匹配)
【发布时间】:2016-03-31 09:41:08
【问题描述】:

问题

这可能是由于我缺乏 Scala 知识,但似乎在 for 理解 中添加另一个级别应该可以正常工作。如果第一行 for comprehension 被注释掉,代码就可以工作。我最终想要一个 Set[Int] 而不是“1 到 2”,但它可以说明问题。 for 的前两行不需要类型说明符,但我包含它以表明我已经尝试了显而易见的方法。

工具/罐子

  • IntelliJ 2016.1
  • Java 8
  • Scala 2.10.5
  • 卡桑德拉 3.x
  • spark-assembly-1.6.0-hadoop2.6.0.jar(预建)
  • spark-cassandra-connector_2.10-1.6.0-M1-SNAPSHOT.jar(预建)
  • spark-cassandra-connector-assembly-1.6.0-M1-SNAPSHOT.jar(我构建)

代码

case class NotifHist(intnotifhistid:Int, eventhistids:Seq[Int], yosemiteid:String, initiatorname:String)
case class NotifHistSingle(intnotifhistid:Int, inteventhistid:Int, dataCenter:String, initiatorname:String)

object SparkCassandraConnectorJoins {
  def joinQueryAfterMakingExpandedRdd(sc:SparkContext, orgNodeId:Int) {

  val notifHist:RDD[NotifHistSingle] = for {
    orgNodeId:Int <- 1 to 2   // comment out this line and it works
    notifHist:NotifHist <- sc.cassandraTable[NotifHist](keyspace, "notifhist").where("intorgnodeid = ?", orgNodeId)
    eventHistId <- notifHist.eventhistids
  } yield NotifHistSingle(notifHist.intnotifhistid, eventHistId, notifHist.yosemiteid, notifHist.initiatorname)
  ...etc...
 }

编译输出

Information:3/29/16 8:52 AM - Compilation completed with 1 error and 0 warnings in 1s 507ms
       /home/jpowell/Projects/SparkCassandraConnector/src/com/mir3/spark/SparkCassandraConnectorJoins.scala
**Error:(88, 21) type mismatch;
 found   : scala.collection.immutable.IndexedSeq[Nothing]
 required: org.apache.spark.rdd.RDD[com.mir3.spark.NotifHistSingle]
      orgNodeId:Int <- 1 to 2
                    ^**

稍后

@slouc 感谢您的全面回答。我正在使用 for comprehension 的 语法糖来保持第二条语句的状态以填充 NotifHistSingle ctor 中的元素,所以我看不到如何让等效的地图/平面图工作。因此,我采用了以下解决方案:

def joinQueryAfterMakingExpandedRdd(sc:SparkContext, orgNodeIds:Set[Int]) {

  def notifHistForOrg(orgNodeId:Int): RDD[NotifHistSingle] = {
    for {
      notifHist <- sc.cassandraTable[NotifHist](keyspace, "notifhist").where("intorgnodeid = ?", orgNodeId)
      eventHistId <- notifHist.eventhistids
    } yield NotifHistSingle(notifHist.intnotifhistid, eventHistId, notifHist.yosemiteid, notifHist.initiatorname)
  }
  val emptyTable:RDD[NotifHistSingle] = sc.emptyRDD[NotifHistSingle]
  val notifHistForAllOrgs:RDD[NotifHistSingle] = orgNodeIds.foldLeft(emptyTable)((accum, oid) => accum ++ notifHistForOrg(oid))
}

【问题讨论】:

    标签: scala apache-spark spark-cassandra-connector


    【解决方案1】:

    对于理解实际上是语法糖;下面真正发生的是一系列链接的flatMap 调用,最后有一个map 替换yield。 Scala 编译器会像这样翻译每一个 for 理解。如果您在 for comprehension 中使用 if 条件,它们将被转换为过滤器,如果您没有产生任何内容,则使用 foreach。如需更多信息,请参阅here

    所以,解释一下你的情况 - 这个:

    val notifHist:RDD[NotifHistSingle] = for {
      orgNodeId:Int <- 1 to 2   // comment out this line and it works
      notifHist:NotifHist <- sc.cassandraTable[NotifHist](keyspace, "notifhist").where("intorgnodeid = ?", orgNodeId)
      eventHistId <- notifHist.eventhistids
    } yield NotifHistSingle(...)
    

    实际上是由编译器翻译成这样的:

    val notifHist:RDD[NotifHistSingle] = (1 to 2)
      .flatMap(x => sc.cassandraTable[NotifHist](keyspace, "notifhist").where("intorgnodeid = ?", x)
      .flatMap(x => x.eventhistids)
      .map(x => NotifHistSingle(...))
    

    如果您包含 1 to 2 行,则会收到错误消息,因为这会使您的 for 理解对序列(更准确地说是向量)进行操作。因此,当调用flatMap() 时,编译器希望您跟进一个函数,将向量的每个元素转换为GenTraversableOnce。如果您仔细查看您的 for 表达式的类型(大多数 IDE 只需将鼠标悬停在它上面就会显示它),您可以自己看到它:

    def flatMap[B, That](f: A => GenTraversableOnce[B])(implicit bf: CanBuildFrom[Repr, B, That]): That
    

    这就是问题所在。编译器不知道如何使用返回CassandraRDD 的函数来flatMap 向量1 to 10。它需要一个返回GenTraversableOnce 的函数。如果您删除 1 to 2 行,那么您将删除此限制。

    底线 - 如果你想使用 for 理解并从中产生值,你必须遵守类型规则。将由非序列且不能转换为序列的元素组成的序列展平是不可能的。

    您始终可以使用map 而不是flatMap,因为地图的限制较少(它需要A =&gt; B 而不是A =&gt; GenTraversableOnce[B])。这意味着您将获得一个序列,其中每个元素都是一组结果(每个查询一个组),而不是在一个巨大的序列中获取所有结果。您还可以尝试使用这些类型,尝试从查询结果中获取 GenTraversableOnce(例如调用 sc.cassandraTable().where().toArray 或其他东西;我并没有真正使用 Cassandra,所以我不知道)。

    【讨论】:

      猜你喜欢
      • 2015-03-12
      • 2020-04-13
      • 1970-01-01
      • 2017-01-30
      • 2020-10-17
      • 2017-04-19
      • 1970-01-01
      • 2020-05-08
      • 1970-01-01
      相关资源
      最近更新 更多