【发布时间】:2017-04-11 02:06:41
【问题描述】:
总结
我的问题是关于 Apache Spark Streaming 如何通过改进并行化或将多个写入合并为一个更大的写入来处理需要很长时间的输出操作。在这种情况下,写入是对 Neo4J 的密码请求,但它可以应用于其他数据存储。
环境
我有一个用 Java 编写的 Apache Spark Streaming 应用程序,它写入 2 个数据存储:Elasticsearch 和 Neo4j。以下是版本:
- Java 8
- Apache Spark 2.11
- Neo4J 3.1.1
- Neo4J Java 螺栓驱动程序 1.1.2
Elasticsearch 的输出非常简单,因为我使用了 Elasticsearch-Hadoop for Apache Spark 库。
我们的直播
我们的输入是从 Kafka 接收到的关于特定主题的流,我通过 map 函数反序列化流的元素以创建 JavaDStream<[OurMessage]> dataStream。然后我对此消息进行转换以创建一个密码查询String cypherRequest(使用 OurMessage 到字符串的转换),该查询发送到管理 Bolt Driver 与 Neo4j 的连接的单例(我知道我应该使用连接池,但也许那是另一个问题)。密码查询根据 OurMessage 的内容生成许多节点和/或边。
代码如下所示。
dataStream.foreachRDD( rdd -> {
rdd.foreach( cypherQuery -> {
BoltDriverSingleton.getInstance().update(cypherQuery);
});
});
优化的可能性
关于如何提高吞吐量,我有两个想法:
- 我不确定 Spark Streaming 并行化是否会下降到 RDD 元素级别。意思是,RDD 的输出可以并行化(在 `stream.foreachRDD()` 中,但是可以并行化 RDD 的每个元素(在 `rdd.foreach()` 中)。如果是后者,那么 `reduce对我们的 `dataStream` 进行的`转换增加了 Spark 并行输出这些数据的能力(每个 JavaRDD 将只包含一个密码查询)?
- 即使改进了并行化,如果我可以实现某种生成器,它使用 RDD 的每个元素来创建单个密码查询,该查询添加来自所有元素的节点/边,而不是一个密码查询每个 RDD。但是,如果不使用另一个 kafka 实例,我怎么能做到这一点,这可能是矫枉过正?
我是不是想太多了?我试图研究太多,以至于我可能太深入了。
旁白:如果其中有任何完全错误,我提前道歉。你不知道你不知道什么,我刚刚开始使用 Apache Spark 和 Java 8 w/lambdas。正如 Spark 用户现在必须知道的那样,要么 Spark 由于其范式非常不同而具有陡峭的学习曲线,要么我是个白痴 :)。
感谢任何可以提供帮助的人;这是我很久以来的第一个 StackOverflow 问题,所以请留下反馈,我会根据需要回复并更正这个问题。
【问题讨论】:
-
我们需要一些关于 Neo4J 设置的信息。有没有设置索引?我已经能够通过触发文档很容易地击倒elasticsearch,因为它必须索引所有进来的东西。使用logstash会有所帮助,因为它会以它可以处理的速度缓冲和提供文档到elasticsearch。也就是说,如果没有看到每个查询对图表所做的修改类型,我们将无法提供太多帮助。证明问题的完整示例会有所帮助。我猜火花在这里是无关紧要的,一个紧密的循环会显示同样的问题。
-
我反对这种互斥性。 Spark 有一个陡峭的学习曲线,并且我是个白痴。
-
@Vidya 很公平:)
-
瓶颈是什么?是neo4j,是Singleton,是构建查询吗?另外,当您说 Singleton 时,这在分布式环境中意味着什么,例如,您的所有密码查询是否都通过 Spark 驱动程序?
-
RDD 具有函数
foreachPartition,它为每个分区提供了一个迭代器。您可以通过重新分区来增加并行度,然后使用foreachPartition和 BoltDriverSingleton 的“线程本地”实例(如果有意义的话)来构建小批量密码查询。希望你明白我的意思?
标签: java apache-spark neo4j streaming spark-streaming