【问题标题】:Processing multiple message in parallel with ActiveMQ使用 ActiveMQ 并行处理多条消息
【发布时间】:2014-12-27 17:42:55
【问题描述】:

我想使用简单的处理器/异步处理器作为目标并行处理队列中的消息。处理器每条消息需要一点时间,但可以单独处理每条消息,因此可以同时处理(在正常范围内)。

我很难找到示例,尤其是关于骆驼路线的 xml 配置。

到目前为止,我已经定义了一个线程池、路由和处理器:

<threadPool id="smallPool" threadName="MyProcessorThread" poolSize="5" maxPoolSize="50" maxQueueSize="100"/>
<route>
    <from uri="broker:queue:inbox" />
    <threads executorServiceRef="smallPool">
        <to uri="MyProcessor" />
    </threads>
</route>
<bean id="MyProcessor" class="com.example.java.MyProcessor" />

我的处理器看起来像:

public class MyProcessor implements Processor {
    @Override
    public void process(Exchange exchange) throws Exception {
        Message in = exchange.getIn();
        String msg = in.getBody(String.class);      
        System.out.println(msg);
        try {
            Thread.sleep(10 * 1000); // Do something in the background
        } catch (InterruptedException e) {}
        System.out.println("Done!");
    }
}

不幸的是,当我将消息发布到队列时,它们仍然被一一处理,每个延迟 10 秒(我的“后台任务”)。

谁能指出我使用定义的线程池处理消息的正确方向或解释我做错了什么?

【问题讨论】:

  • 你试过concurrentConsumers吗?在您的示例中,这将是 broker:queue:inbox?concurrentConsumers=5 (或其他)。我不记得每个消费者线程是否会使用自己的处理器实例,但除非您在处理器中启动一个新线程,否则无论如何您都需要多个线程,因为在再次调用“from”之前必须完成路由。
  • 您看过 Camel 负载均衡器吗? camel.apache.org/load-balancer.html
  • @Fortyrunner 我认为负载均衡器将用于和出站或到端点。在这种情况下,并发消费者将发挥作用。
  • 我使用的一个技巧是设置多个 SEDA 队列并将输入队列负载平衡到这些 SEDA 队列上。 SEDA 本身在单独的线程上运行。也许我应该看看concurrentConsumers。有不止一种方法可以做到!

标签: parallel-processing apache-camel activemq


【解决方案1】:

您应该使用 cmets 中所述的 concurrentConsumers 选项,

<route>
    <from uri="broker:queue:inbox?concurrentConsumers=5" />
    <to uri="MyProcessor" />
</route>

请注意,您还可以设置maxConcurrentConsumers 以使用最小/最大范围的并发消费者,因此 Camel 将根据负载自动增长/缩小。

在 JMS 文档中查看更多详细信息

【讨论】:

  • 有没有办法在运行时动态更改concurrentConsumers 属性?骆驼提供这样的能力?一种在运行时增加或减少消费者数量的方法。也许通过 JMX 或 HawtIO mgtm 控制台。
  • 我们有什么方法可以动态更改concurrentConcumers? @Tuelho 你有什么解决办法吗?
  • 您可以使用 JMX / 从 hawtio Web 控制台更改它
  • 我正在从 s3 存储桶读取数据 -> 基于令牌进行拆分 -> 流式传输 -> parallelProcessing()。我没有指定任何线程,但是,看起来它正在并行处理我的子消息。我是否正确解释它?骆驼 v2.19.1。提前致谢!
猜你喜欢
  • 2019-07-26
  • 1970-01-01
  • 2019-06-02
  • 1970-01-01
  • 2016-11-20
  • 2016-06-06
  • 1970-01-01
  • 2015-06-18
  • 2016-04-14
相关资源
最近更新 更多