【问题标题】:RDD accessing values in another RDDRDD 访问另一个 RDD 中的值
【发布时间】:2015-04-02 20:27:21
【问题描述】:

我有一个 RDD 需要访问另一个 RDD 的数据。但是,我总是收到Task not Serializable 错误。我已经扩展了Serializable 类,但是它没有用。代码是:

val oldError = rddOfRatings.aggregate(0.0)((accum, rating) =>
accum + calcError(rating.rating,
us.lookup(rating.user)(0),
it.lookup(rating.product)(0)).abs, _+_ ) / rddSize

其中usitrddOfRatings 是另一个RDD。我不明白的是,如果RDD 是不可变的,那么为什么不允许我允许从另一个RDD 访问RDD?问题似乎出在 usit 上,因为当我将它们删除以用于本地集合时,它工作正常。

谢谢。

【问题讨论】:

    标签: scala apache-spark rdd


    【解决方案1】:

    RDD 确实是不可序列化的,因为它们必须捕获变量(例如 SparkContext)。为了解决这个问题,将三个 RDD 连接在一起,您的累加器闭包中将拥有所有必要的值。

    【讨论】:

      【解决方案2】:

      rdd.lookup1 是一项昂贵的操作,即使可以,您也可能不想这样做。

      此外,“序列化”RDD 没有意义,因为 RDD 只是对数据的引用,而不是数据本身。

      此处采用的方法可能取决于这些数据集的大小。如果usit RDD 的大小与rddOfRatings 大致相同(考虑到预期的查找,这就是它的样子),最好的方法是预先加入它们。

      // 请注意,我不知道您的收藏的实际结构,所以以此为例说明

      val ratingErrorByUser = us.map(u => (u.id, u.error))
      val ratingErrorByProduct = it.map(i=> (i.id, i.error)) 
      val ratingsBykey = rddOfRatings.map(r=> (r.user, (r.product, r.rating)))
      val ratingsWithUserError = ratingsByKey.join(ratingErrorByUser)
      val ratingsWithProductError = ratingsWithUserError.map{case (userId, ((prodId, rating),userErr))} => (prodId,(rating, userErr))}
      val allErrors = ratingsWithProductError.join(ratingErrorByProduct)
      val totalErr = allErrors.map{case (prodId,((rating, userErr),prodErr)) => calcError(userErr, math.abs(prodErr), rating)}.reduce(_+_)
      val total = totalErr / rddOfRatings.count
      

      使用Spark DataFrame API 可能会容易得多

      1 如果必须查找(在这种情况下看起来不像!),请查看 Spark IndexedRdd

      【讨论】:

      • 谢谢。它启发了我改变我的代码并删除lookup。 IndexedRdd 很有趣。
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-08-23
      • 2016-12-16
      • 2016-12-23
      • 1970-01-01
      • 1970-01-01
      • 2018-02-05
      相关资源
      最近更新 更多