【问题标题】:reduceByKey is not a memberreduceByKey 不是会员
【发布时间】:2016-04-26 20:33:44
【问题描述】:

您好,我的代码只是从文档中获取字数。在生成输出之前,我还需要使用映射来查找数据值。这是代码。

   requests
    .filter(_.description.exists(_.length > 0))
    .flatMap { case request =>
      broadcastDataMap.value.get(request.requestId).map {
        data =>
          val text = Seq(
            data.name,
            data.taxonym,
            data.pluralTaxonym,
            request.description.get
          ).mkString(" ")
          getWordCountsInDocument(text).map { case (word, count) =>
            (word, Map(request.requestId -> count))
          }
      }
    }
    .reduceByKey(mergeMap)

错误信息是

reduceByKey is not a member of org.apache.spark.rdd.RDD[scala.collection.immutable.Map[String,scala.collection.immutable.Map[Int,Int]]]

我该如何解决这个问题?我确实需要调用 getWordCountsInDocument。谢谢!

【问题讨论】:

  • 你需要获取 PairRDD。尝试在 reduceByKey 之前使用 .map()

标签: scala apache-spark rdd


【解决方案1】:

reduceByKey 是 PairRDDFunctions 的成员,基本上它以RDD[(K, V)] 的形式隐式添加到 RDD。您可能需要将结构展平为RDD[String, Map[Int,Int]]

如果您可以为您的输入提供类型(requestsbroadcastDataMapmergeMap),我们或许可以为该转换提供一些帮助。

根据提供的类型,假设 getWordCountsInDocument 的返回类型是一些 Collection[(word, count: Int)]

变化:

broadcastDataMap.value.get(request.requestId).map {

broadcastDataMap.value.get(request.requestId).flatMap {

应该解决问题。

【讨论】:

  • 谢谢。 #1:requests是request的RDD,有(requestId, description) #2:broadcastDataMap是Map[requestID,Data(name,taxonym,pluralTaxonym)]的广播 #3:mergeMap是一个函数,它接受两个Map[Int , Int] 并返回一个 Map[Int, Int]: private def mergeMap(map1: Map[Int, Int], map2: Map[Int, Int]): Map[Int, Int] = { (map1 ++ map2) .map { case (key, _) => (key, map1.getOrElse(key, 0) + map2.getOrElse(key, 0)) } }
猜你喜欢
  • 2015-05-04
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-05-17
  • 1970-01-01
  • 1970-01-01
  • 2018-01-20
  • 2015-06-25
相关资源
最近更新 更多