【问题标题】:Speed up collaborative filtering for large dataset in Spark MLLib加速 Spark MLLib 中大型数据集的协同过滤
【发布时间】:2016-12-30 11:50:02
【问题描述】:

我正在使用 MLlib 的矩阵分解向用户推荐项目。我有一个很大的隐式交互矩阵,包含 M=2000 万用户和 N=50k 个项目。在训练模型后,我想为每个用户获得一个简短的推荐列表(例如 200 个)。我在MatrixFactorizationModel 中尝试了recommendProductsForUsers,但它非常非常慢(运行了9 个小时,但还远未完成。我正在测试50 个执行程序,每个执行程序都有8g 内存)。这可能是意料之中的,因为recommendProductsForUsers 需要计算所有M*N 用户与项目的交互并为每个用户获得顶部。

我会尝试使用更多的执行器,但从我在 Spark UI 上的应用程序详细信息中看到,我怀疑即使我有 1000 个执行器它也能在数小时或一天内完成(9 小时后它仍然在 flatmaphttps://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


【解决方案1】:

@Jack Lei:你找到答案了吗? 我自己尝试了一些事情,但只起到了一点作用。

例如:我试过了

javaSparkContext.setCheckpointDir("checkpoint/");

这有助于避免重复计算。

还尝试为每个 Executor 添加更多内存和开销 spark 内存

--conf spark.driver.maxResultSize=5g --conf spark.yarn.executor.memoryOverhead=4000

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2014-09-05
    • 2015-05-23
    • 1970-01-01
    • 2017-11-09
    • 1970-01-01
    • 2017-10-21
    • 2016-10-16
    • 1970-01-01
    相关资源
    最近更新 更多