【发布时间】:2016-12-30 11:50:02
【问题描述】:
我正在使用 MLlib 的矩阵分解向用户推荐项目。我有一个很大的隐式交互矩阵,包含 M=2000 万用户和 N=50k 个项目。在训练模型后,我想为每个用户获得一个简短的推荐列表(例如 200 个)。我在MatrixFactorizationModel 中尝试了recommendProductsForUsers,但它非常非常慢(运行了9 个小时,但还远未完成。我正在测试50 个执行程序,每个执行程序都有8g 内存)。这可能是意料之中的,因为recommendProductsForUsers 需要计算所有M*N 用户与项目的交互并为每个用户获得顶部。
我会尝试使用更多的执行器,但从我在 Spark UI 上的应用程序详细信息中看到,我怀疑即使我有 1000 个执行器它也能在数小时或一天内完成(9 小时后它仍然在 flatmap 中https://github.com/apache/spark/blob/master/mllib/src/main/scala/org/apache/spark/mllib/recommendation/MatrixFactorizationModel.scala#L279-L289,总共 10000 个任务,只完成了 ~200 个)
除了增加执行者的数量之外,我还可以调整其他什么来加快推荐过程吗?
这里是示例代码:
val data = input.map(r => Rating(r.getString(0).toInt, r.getString(1).toInt, r.getLong(2))).cache
val rank = 20
val alpha = 40
val maxIter = 10
val lambda = 0.05
val checkpointIterval = 5
val als = new ALS()
.setImplicitPrefs(true)
.setCheckpointInterval(checkpointIterval)
.setRank(rank)
.setAlpha(alpha)
.setIterations(maxIter)
.setLambda(lambda)
val model = als.run(ratings)
val recommendations = model.recommendProductsForUsers(200)
recommendations.saveAsTextFile(outdir)
【问题讨论】:
-
您确定 Spark 充分利用了 8g 内存吗?也许它真的经常命中磁盘缓存。
标签: scala apache-spark apache-spark-mllib collaborative-filtering