【发布时间】:2015-05-27 13:43:51
【问题描述】:
我正在尝试使用与 stanford NLP (3.4.1) 集成的 spark/scala 来处理数百万个数据。 由于我使用的是社交媒体数据,因此我必须使用 NLP 进行主题生成(pos 标记)和 Sentiment 计算。
我必须分别处理 Twitter 数据和 NON Twitter 数据。所以我有两个类处理 Twitter/Non Twitter
我正在使用每个类的 lasy val 初始化来加载 stanfordNLP
features: Seq[String] = Seq("tokenize","ssplit","pos","parse","sentiment")
val props = new Properties()
props.put("annotators", features.mkString(", "))
props.put("pos.model", "tagger/gate-EN-twitter.model")
props.put("parse.model", "tagger/englishSR.ser.gz");
val pipeline = new StanfordCoreNLP(props)
注意:对于上面的 Twitter,我使用不同的 pos 模型和 shift reduce 解析模型进行解析。我使用 shift reduce 解析器的原因是为了一些垃圾 默认 PCFG 模型需要大量时间来处理并获得一些异常。Shift reduce 解析器在加载时大约需要 15 秒,在运行时处理速度更快 数据。
非推特类
features: Seq[String] = Seq("tokenize","ssplit","pos","parse","sentiment")
val props = new Properties()
props.put("annotators", features.mkString(", "))
props.put("parse.model", "tagger/englishSR.ser.gz");
这里我使用的是默认的 pos 模型和 shift reduce 解析器
问题:
目前我们使用 8 个 6 核节点运行,我可以使用 48 个分区运行。用于处理数百万数据 使用上述具有较小分区的配置,它对我来说很好。
8 个节点和 6 个核心,我们几乎有 48 个分区,如果我使用 42 个分区运行,大约需要 1 小时才能完成处理。
使用当前配置,我需要将其扩展到 200 个分区
8 个节点和 6 个核心,我们几乎有 48 个分区,如果我们以 200 个分区数运行,大约需要 2 小时,最后抛出一些异常,说一个节点丢失 或 java.lang.IllegalArgumentException: annotator "sentiment" requires annotator "binarized_trees" etc etc.
问题在于,如果我们将分区数量扩大到 200 个,使用 8 个节点和 6 个核心,而我们只有 48 个核心。
我怀疑它是因为加载移位减少了每个分区的解析器加载。我想一次加载这个类,然后进行广播,但standforndNLP 类不可序列化,所以我无法广播。
我们需要扩展到 200 个分区的原因是它运行速度快,处理这些数据的时间更短。
【问题讨论】:
-
您是否尝试过将 NLP jar 作为 submit-spark 的一部分传递(或将其与您的应用程序捆绑在一起)?本质上,调度程序应该在远程 JVM 中启动任务之前将额外的 jars 发送给工作人员
-
是的,我正在捆绑应用程序 jar。但我无法通过更多的分区来扩展应用程序
标签: apache-spark nlp stanford-nlp sentiment-analysis