【问题标题】:Developing a spark streaming application开发火花流应用程序
【发布时间】:2014-09-10 21:23:01
【问题描述】:

所以我要解决的问题如下:

  • 我需要一个以特定频率发出消息的数据源
  • 有 N 个神经网络需要单独处理每条消息
  • 所有神经网络的输出都被聚合,只有当每个消息的所有 N 个输出都被收集时,才应该声明消息已完全处理
  • 最后,我应该测量消息被完全处理所花费的时间(从发出消息到收集该消息的所有 N 个神经网络输出之间的时间)

我很好奇如何使用火花流处理这样的任务。

我当前的实现使用 3 种类型的组件:一个自定义接收器和两个实现 Function 的类,一个用于神经网络,一个用于最终聚合器。

概括地说,我的应用程序是这样构建的:

JavaReceiverInputDStream<...> rndLists = jssc.receiverStream(new JavaRandomReceiver(...));

Function<JavaRDD<...>, Void> aggregator = new JavaSyncBarrier(numberOfNets);

for(int i = 0; i < numberOfNets; i++){
    rndLists.map(new NeuralNetMapper(neuralNetConfig)).foreachRDD(aggregator);
}

不过,我遇到的主要问题是它在本地模式下的运行速度比提交到 4 节点集群时要快。

我的实现一开始是错误的还是这里发生了其他事情?

这里还有一个完整的帖子http://apache-spark-user-list.1001560.n3.nabble.com/Developing-a-spark-streaming-application-td12893.html,其中详细介绍了前面提到的三个组件中的每一个的实现。

【问题讨论】:

    标签: apache-spark spark-streaming


    【解决方案1】:

    似乎有很多重复的对象实例化和序列化。后者可能会影响您在集群中的性能。

    您应该只尝试一次实例化您的神经网络。您必须确保它们是可序列化的。您应该使用flatMap 而不是多个maps + union。大致如下:

    // Initialize neural net first
    List<NeuralNetMapper> neuralNetMappers = new ArrayList<>(numberOfNets);
    for(int i = 0; i < numberOfNets; i++){
        neuralNetMappers.add(new NeuralNetMapper(neuralNetConfig));
    }
    
    // Then create a DStream applying all of them
    JavaDStream<Result> neuralNetResults = rndLists.flatMap(new FlatMapFunction<Item, Result>() {
        @Override
        public Iterable<Result> call(Item item) {
            List<Result> results = new ArrayList<>(numberOfNets);
            for (int i = 0; i < numberOfNets; i++) {
                results.add(neuralNetMappers.get(i).doYourNeuralNetStuff(item));
            }
            return results;
        }
    });
    
    // The aggregation stuff
    neuralNetResults.foreachRDD(aggregator);
    

    如果您负担得起以这种方式初始化网络,您可以节省大量时间。此外,您在链接帖子中包含的 union 内容似乎是不必要的,并且会影响您的表现:flatMap 可以。

    最后,为了进一步调整您在集群中的性能,you can use the Kryo serializer

    【讨论】:

    • 感谢您的回复,我确实开始在单个节点上使用该解决方案,它确实加快了速度。我尝试在多节点集群上运行它,但似乎没有任何性能提升。所以我的问题是我将如何以受益于多节点集群的方式编写它?一些关于 spark 如何分配工作负载的一般知识也会有所帮助,因为我没有找到这方面的很多细节。
    • 首先,您应该确定哪些阶段需要更长时间。您可以使用localhost:4040/stages 上的监控网络应用程序按阶段的使用量对阶段进行排序。比较在单节点和集群中运行的此输出。此外,您将从使用 Spark 1.1.0 中受益,因为它具有更好的流式调试功能。在该 Web 界面中,您还应该查看“流式传输”选项卡,您将在其中看到每秒正在处理的事件数。您可能应该使用它来衡量性能。
    • 您的集群是如何连接的?节点之间的不良连接将对性能产生严重的负面影响。
    • 再次感谢您的快速回复,连接应该不是问题,这些是连接在同一个 LAN 上的 4 台不同的机器,它们之间的延迟小于 1 毫秒。我会看看你提到的统计数据。与此同时,我觉得我应该提到事件当前以每秒 1 个的速度传入,批处理间隔也设置为 1 秒(所以基本上每秒 1 个事件)。理论上需要最多时间的部分是神经网络处理本身(理想情况下,我会达到数百个足够多的神经网络),这就是我希望并行化的部分。
    • 我在 spark 之前也有过一些 Storm 的经验,它有一个更直观的分配其工作的范式,这就是我在 spark 上遇到的主要问题,我无法理解它是如何并行化其应用程序的。
    猜你喜欢
    • 2015-11-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-08-09
    • 2019-03-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多