【问题标题】:RDD transformations and actions can only be invoked by the driverRDD 转换和操作只能由驱动程序调用
【发布时间】:2016-02-10 18:14:30
【问题描述】:

错误:

org.apache.spark.SparkException: RDD 转换和动作只能由驱动程序调用,不能在其他转换内部调用;例如,rdd1.map(x => rdd2.values.count() * x) 是无效的,因为值转换和计数操作不能在 rdd1.map 转换中执行。有关详细信息,请参阅 SPARK-5063。

def computeRatio(model: MatrixFactorizationModel, test_data: org.apache.spark.rdd.RDD[Rating]): Double = {
  val numDistinctUsers = test_data.map(x => x.user).distinct().count()
  val userRecs: RDD[(Int, Set[Int], Set[Int])] = test_data.groupBy(testUser => testUser.user).map(u => {
    (u._1, u._2.map(p => p.product).toSet, model.recommendProducts(u._1, 20).map(prec => prec.product).toSet)
  })
  val hitsAndMiss: RDD[(Int, Double)] = userRecs.map(x => (x._1, x._2.intersect(x._3).size.toDouble))

  val hits = hitsAndMiss.map(x => x._2).sum() / numDistinctUsers

  return hits
}

我正在使用MatrixFactorizationModel.scala中的方法,我必须映射用户,然后调用该方法来获取每个用户的结果。通过这样做,我引入了嵌套映射,我认为这会导致问题:

我知道这个问题实际上发生在:

val userRecs: RDD[(Int, Set[Int], Set[Int])] = test_data.groupBy(testUser => testUser.user).map(u => {
  (u._1, u._2.map(p => p.product).toSet, model.recommendProducts(u._1, 20).map(prec => prec.product).toSet)
})

因为在映射时我打电话给model.recommendProducts

【问题讨论】:

  • 你的问题是什么?这个问题非常重要(而且它非常概念化,否则火花可能会混乱),如果我是正确的,它在 UC 提供的spark's course 中进行了讨论伯克利

标签: scala mapreduce apache-spark apache-spark-mllib


【解决方案1】:

MatrixFactorizationModel 是一个分布式模型,因此您不能简单地从操作或转换中调用它。与您在这里所做的最接近的事情是这样的:

import org.apache.spark.rdd.RDD
import org.apache.spark.mllib.recommendation.{MatrixFactorizationModel, Rating}

def computeRatio(model: MatrixFactorizationModel, testUsers: RDD[Rating]) = {
  val testData = testUsers.map(r => (r.user, r.product)).groupByKey
  val n = testData.count

  val recommendations = model
     .recommendProductsForUsers(20)
     .mapValues(_.map(r => r.product))

  val hits = testData
    .join(recommendations)
    .values
    .map{case (xs, ys) => xs.toSet.intersect(ys.toSet).size}
    .sum

  hits / n
}

注意事项:

  • distinct 是一项昂贵的操作,在这里完全过时了,因为您可以从分组数据中获取相同的信息
  • 而不是 groupBy 后跟投影 (map),先投影后分组。如果您只需要产品 ID,则没有理由转移完整评分。

【讨论】:

  • 有一个问题我不认为它是答案,但它与 mllib 有人在recommendProductsForUsers stackoverflow.com/questions/33646889/… 上也有同样的问题@
  • 我认为它可能与 mllib 的版本有关,该方法是在我运行 1.3 的服务器上的 1.4 上引入的,因此该方法在那里不可用,因为它会引发非成员错误
  • 如果是这样,我不确定这里是否有任何合理的解决方案。如果 testUsers 中的元素数量足够小以适合驱动程序,您可以映射到本地结构。
  • 我也是这么想的,如果我明白你的意思是做一个collect 并映射过去?就像我在做的那样?我试图在我的代码中实现推荐产品ForUsers,但问题是recommendAll 是一个私有方法,除了所有元素都可用你能想到的任何其他解决方案吗?
  • 没错。将test_data: RDD[Rating] 替换为test_data: Seq[Rating]。关于其他选项,我没有想到任何明智的选择。几乎所有必需的方法都是私有的,因此如果不修改源代码和重建,您将无处可去。有什么理由使用 1.3?
猜你喜欢
  • 1970-01-01
  • 2012-04-25
  • 1970-01-01
  • 2015-09-02
  • 2012-09-23
  • 1970-01-01
  • 2019-04-02
  • 2013-07-13
  • 1970-01-01
相关资源
最近更新 更多