【发布时间】: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