【问题标题】:Consuming from Camel queue every x minutes每 x 分钟从 Camel 队列中消费
【发布时间】:2019-11-14 00:33:44
【问题描述】:
尝试实现一种方法,让我的消费者每 30 分钟左右从队列接收消息。
就上下文而言,我的错误队列中有 20 条消息,直到 x 分钟过去,然后我的路由消耗队列上的所有消息并继续“休眠”,直到又过去了 30 分钟。
不确定实现这一点的最佳方法,我尝试过 spring @Scheduled、camel timer 等,但都没有达到我的预期。我一直试图让它与路由策略一起使用,但在正确的功能中没有骰子。它似乎立即从队列中消耗。
路由策略是正确的路径还是有其他东西可以使用?
【问题讨论】:
标签:
spring-boot
apache-camel
【解决方案1】:
从队列中读取的路由总是尽可能快地读取任何消息。
您可以做的一件事是启动/停止或暂停使用消息的路由,因此请进行这种设置:
route 1: error_q_reader, which goes from('jms').
route 2: a timed route that fires every 20 mins
路由 2 可以使用 control bus 组件来启动路由。
from('timer?20mins') // or whatever the timer syntax is...
.to("controlbus:route?routeId=route1&action=start")
这里的棘手部分是知道何时停止路线。你让它运行5分钟吗?一旦消息全部消耗,您想停止它吗?可能有一种方法可以运行另一条可以检查队列深度的路由(比如每 1 分钟左右),如果它是 0,则关闭route 1,你可能会让它工作,但我可以向你保证这会变得混乱您尝试处理许多异步操作。
您还可以尝试一些更奇特的东西,例如自定义QueueBrowseStrategy,它可以在队列中没有消息时触发关闭route 1 的事件。
【解决方案2】:
我构建了一个客户 bean 来排空队列并关闭,但这不是一个非常优雅的解决方案,我很想找到一个更好的解决方案。
public class TriggeredPollingConsumer {
private ConsumerTemplate consumer;
private Endpoint consumerEndpoint;
private String endpointUri;
private ProducerTemplate producer;
private static final Logger logger = Logger.getLogger( TriggeredPollingConsumer.class );
public TriggeredPollingConsumer() {};
public TriggeredPollingConsumer( ConsumerTemplate consumer, String endpoint, ProducerTemplate producer ) {
this.consumer = consumer;
this.endpointUri = endpoint;
this.producer = producer;
}
public void setConsumer( ConsumerTemplate consumer) {
this.consumer = consumer;
}
public void setProducer( ProducerTemplate producer ) {
this.producer = producer;
}
public void setConsumerEndpoint( Endpoint endpoint ) {
consumerEndpoint = endpoint;
}
public void pollConsumer() throws Exception {
long count = 0;
try {
if ( consumerEndpoint == null ) consumerEndpoint = consumer.getCamelContext().getEndpoint( endpointUri );
logger.debug( "Consuming: " + consumerEndpoint.getEndpointUri() );
consumer.start();
producer.start();
while ( true ) {
logger.trace("Awaiting message: " + ++count );
Exchange exchange = consumer.receive( consumerEndpoint, 60000 );
if ( exchange == null ) break;
logger.trace("Processing message: " + count );
producer.send( exchange );
consumer.doneUoW( exchange );
logger.trace("Processed message: " + count );
}
producer.stop();
consumer.stop();
logger.debug( "Consumed " + (count - 1) + " message" + ( count == 2 ? "." : "s." ) );
} catch ( Throwable t ) {
logger.error("Something went wrong!", t );
throw t;
}
}
}
您配置 bean,然后从计时器调用 bean 方法,并配置直接路由来处理队列中的条目。
from("timer:...")
.beanRef("consumerBean", "pollConsumer");
from("direct:myRoute")
.to(...);
然后它将读取队列中的所有内容,并在一分钟内没有条目到达时立即停止。您可能想减少分钟,但我发现一秒钟意味着如果 JMS 有点慢,它会在排空队列的中途超时。
我也一直在研究 sjms-batch 组件,以及它如何与 pollEnrich 模式一起使用,但到目前为止我还没有能够让它工作。
【解决方案3】:
我通过在微服务方法中将我的应用程序用作 CronJob 解决了这个问题,并赋予它优雅地自行关闭的能力,我们可以设置属性 camel.springboot.duration-max-idle-seconds。因此,您的 JMS 消费者路线保持简单。
另一种方法是声明一个路由来控制 JMS 消费者路由的“生命周期”(启动、睡眠和恢复)。
我强烈建议您使用第一种方法。
【解决方案4】:
如果你使用ActiveMQ,你可以利用它的Scheduler feature。
您可以延迟在代理上传递消息,只需将 JMS 属性 AMQ_SCHEDULED_DELAY 设置为延迟的毫秒数。骆驼路线很容易
.setHeader("AMQ_SCHEDULED_DELAY", 60000)
这并不完全符合您的要求,因为它不是每 30 分钟排空一个队列,而是将每条单独的消息延迟 30 分钟。
请注意,您必须在代理配置中启用schedulerSupport。否则延迟属性将被忽略。
<broker brokerName="localhost" dataDirectory="${activemq.data}" schedulerSupport="true">
...
</broker>
【解决方案5】:
你应该考虑Aggregation EIP
from(URI_WAITING_QUEUE)
.aggregate(new GroupedExchangeAggregationStrategy())
.constant(true)
.completionInterval(TIMEOUT)
.to(URI_PROCESSING_BATCH_OF_EXCEPTIONS);
此示例描述了以下规则:所有传入URI_WAITING_QUEUE 的对象将被分组到List。 constant(true) 是一个分组条件(没有任何)。并且每个TIMEOUT 周期(以毫秒为单位)所有分组对象都将被传递到URI_PROCESSING_BATCH_OF_EXCEPTIONS 队列中。
所以URI_PROCESSING_BATCH_OF_EXCEPTIONS 队列将处理要处理的对象的List。可以申请Split EIP进行拆分,一一处理。