【问题标题】:Distributed Word2Vec Model Training using Apache Spark 2.0.0 and mllib使用 Apache Spark 2.0.0 和 mllib 进行分布式 Word2Vec 模型训练
【发布时间】:2017-02-06 21:32:14
【问题描述】:

我一直在尝试使用 spark 和 mllib 来训练 word2vec 模型,但我似乎没有在大型数据集上获得分布式机器学习的性能优势。我的理解是,如果我有 w 个工人,那么,如果我创建一个具有 n 个分区的 RDD,其中 n>w 并且我尝试通过调用 Word2Vec 的 fit 函数以 RDD 作为参数来创建一个 Word2Vec 模型,那么 spark 将分发数据统一地在这些 w 个 worker 上训练单独的 word2vec 模型,并在最后使用某种 reducer 函数从这些 w 个模型创建单个输出模型。这将减少计算时间,而不是 1 个块,w 个数据块将被同时处理。权衡是可能会发生一些精度损失,具体取决于最后使用的 reducer 函数。 Spark 中的 Word2Vec 是否真的以这种方式工作?如果确实如此,我可能需要使用可配置参数。

编辑

添加提出这个问题的原因。我在 10 台工作机器上运行 java spark word2vec 代码,并在查看文档后为 executor-memory、driver memory 和 num-executors 设置合适的值,用于映射到 rdd 分区的 2.5gb 输入文本文件,然后用作mllib word2vec 模型的训练数据。培训部分花费了数小时。工作节点的数量似乎对训练时间没有太大影响。相同的代码在较小的数据文件(大约 10 MB 的数量级)上成功运行

代码

SparkConf conf = new SparkConf().setAppName("SampleWord2Vec");
conf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer");
conf.registerKryoClasses(new Class[]{String.class, List.class});
JavaSparkContext jsc = new JavaSparkContext(conf);
JavaRDD<List<String>> jrdd = jsc.textFile(inputFile, 3).map(new Function<String, List<String>>(){            
        @Override
        public List<String> call(String s) throws Exception {
            return Arrays.asList(s.split(","));
        }        
});
jrdd.persist(StorageLevel.MEMORY_AND_DISK());
Word2Vec word2Vec = new Word2Vec()
      .setWindowSize(20)
      .setMinCount(20);

Word2VecModel model = word2Vec.fit(jrdd);
jrdd.unpersist(false);
model.save(jsc.sc(), outputfile);
jsc.stop();
jsc.close();

【问题讨论】:

  • 如果您分享您的代码以及有关如何运行 spark-submit 的更多详细信息,将会有所帮助。当你跑步时,你是否看到你所有的工人一直都在活动? Spark 历史 UI 将让您深入了解。您的代码可能没有性能并且您没有完全分发您的代码。 Spark ML 包括基于数据帧 API 的 JavaWord2Vec。这应该很快。
  • spark ml JavaWord2Vec(dataframes api) 是否应该比 mllib 版本 (javardd api) 更好。我放弃了 spark ml 版本,因为当我尝试迭代模型向量时它给出了一些编译错误。
  • 数据帧 API 背后的催化剂优化器性能更高,应该更容易。你不会迭代,这是使用 Spark 的一种可怕的糟糕方式。 ML 允许您构建管道,这些管道实际上对您选择的列的所有值执行功能映射。同样,代码会有所帮助。
  • 我已经用有问题的部分更新了问题。我已经删除了迭代模型向量的部分,但是模型训练步骤花费了太多时间。日志会打印 alpha 的值,因为它从 0.025 下降,并且进展非常缓慢。

标签: java apache-spark apache-spark-mllib word2vec


【解决方案1】:

从 cmets、答案和否决票来看,我想我无法正确地提出我的问题。但我想知道的答案是肯定的,可以在 spark 上并行训练你的 word2vec 模型。此功能的拉取请求是很久以前创建的:

https://github.com/apache/spark/pull/1719

在java中,spark mllib中有一个用于Word2Vec对象的setter方法(setNumPartitions)。这允许您在多个执行器上并行训练您的 word2vec 模型。 根据上面提到的拉取请求中的 cmets:

为了让我们的实现更具可扩展性,我们分别训练每个分区,并在每次迭代后合并每个分区的模型。为了使模型更准确,可能需要多次迭代。” p>

希望这对某人有所帮助。

【讨论】:

  • 你得到了一些基准吗?我也对 gensim、原始 word2vec、spark 的比较感兴趣。(请注意,与其他两个相比,spark 使用的是 skipgram 模型和 cbow)
  • 我遇到了同样的问题——即使使用 DataFrame,Spark w2v 默认使用一个执行器进行训练。正如您所说,必须使用 setNumPartitions 使其并行训练。谢谢你指出这一点。我个人认为这是一个糟糕的默认值设置。
【解决方案2】:

我看不出您的代码有任何本质上的错误。但是,我强烈建议您考虑使用数据帧 API。举个例子,下面是一个经常被抛出的小图表:

另外,我不知道您是如何“迭代”数据框的元素的(它们实际上并不是这样工作的)。这是来自Spark online docs 的示例:

您有一个大致的想法...但您必须首先将您的数据并行化为一个数据框。将您的 javardd 转换为 DataFrame 非常简单。

DataFrame fileDF = sqlContext.createDataFrame(jrdd, Model.class);

Spark 运行有向无环图 (DAG) 代替 MR,但概念是相同的。在您的数据上运行 'fit() 确实会在工作人员的数据上运行,然后减少到单个模型。但是这个模型本身会分布在内存中,直到你决定把它写下来。

但是,作为一个试验,通过 NLTK 或 Word2Vec 的原生 C++ 二进制文件运行同一个文件需要多长时间?

最后一个想法……你坚持到内存和磁盘有什么原因吗? Spark 有一个原生的.cache(),默认情况下会保存在内存中。 Spark 的强大之处在于对内存中保存的数据进行机器学习……内存中的大数据。如果您坚持到磁盘,即使使用 kryo,您也会在磁盘 I/O 上造成瓶颈。恕我直言,首先要尝试的是摆脱这种情况并坚持记忆。如果性能有所提高,那就太好了,但是通过 DataFrame 依靠 Catalyst 的强大功能,您会发现性能的飞跃。

我们没有讨论的一件事是您的集群。考虑一下每个节点有多少内存......每个节点有多少核心......您的集群是否与其他需要资源的应用程序虚拟化(像大多数虚拟主机一样过度配置)......是您在云中的集群?共享还是专用?

您是否查看过 Spark 的 UI 以分析代码的运行时操作?当您在模型拟合时对工人运行 top 时,您会看到什么?你能看到完整的 CPU 利用率吗?您是否尝试过指定 --executor-cores 以确保充分利用 CPU?

我已经多次看到所有工作都在一个工作节点上的一个核心上完成。拥有这些信息会很有帮助。

在对性能进行故障排除时,需要查看很多地方,包括 Spark 配置文件本身!

【讨论】:

  • 我坚持到内存和磁盘,因为程序无法将 jrdd 缓存到内存中。当我遇到这个问题时,我已经更改了默认设置(仅限内存)。同样作为基准,相同的文件在 python 中的 gensim 上运行半小时,单台机器比上面使用的 10 台机器更强大(4 倍内核数,相同的 RAM)。我认为我们正在谈论更多关于配置级别的设置。我想知道当有人调用它时,Spark 如何训练 word2vec 模型,即它是否拆分数据,为这些拆分创建单独的模型并将它们简化为单个模型?
猜你喜欢
  • 2018-03-15
  • 1970-01-01
  • 1970-01-01
  • 2015-08-28
  • 2017-05-17
  • 1970-01-01
  • 2017-05-30
  • 1970-01-01
  • 2020-09-07
相关资源
最近更新 更多