You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Apache Camel路由:队列清空后自动停止的实现方法及示例

在Apache Camel中实现队列清空后停止路由

当然可以实现这个需求,Apache Camel提供了多种灵活的方式来检测队列状态并控制路由启停,以下是具体方案和示例:

一、示例实现(以ActiveMQ队列为例)

通过自定义处理器结合Camel的路由控制API,我们可以在处理消息后检查队列是否为空,为空则停止路由:

import org.apache.camel.CamelContext;
import org.apache.camel.builder.RouteBuilder;
import org.apache.camel.component.jms.JmsComponent;
import org.apache.camel.impl.DefaultCamelContext;
import org.apache.camel.support.DefaultRouteController;
import org.apache.activemq.ActiveMQConnectionFactory;

public class QueueCleanupStopRoute {
    public static void main(String[] args) throws Exception {
        CamelContext context = new DefaultCamelContext();
        
        // 配置ActiveMQ连接
        ActiveMQConnectionFactory connFactory = new ActiveMQConnectionFactory("tcp://localhost:61616");
        context.addComponent("jms", JmsComponent.jmsComponentAutoAcknowledge(connFactory));
        
        context.addRoutes(new RouteBuilder() {
            @Override
            public void configure() throws Exception {
                from("jms:queue:myTargetQueue")
                    .process(exchange -> {
                        // 编写你的消息处理逻辑
                        String msg = exchange.getIn().getBody(String.class);
                        System.out.println("已处理消息: " + msg);
                        
                        // 查询当前队列剩余消息数
                        var session = connFactory.createConnection().createSession(false, 1);
                        long remainingMsgCount = session.createQueue("myTargetQueue").getBrowseableMessages().length;
                        
                        // 队列空则停止路由
                        if (remainingMsgCount == 0) {
                            DefaultRouteController routeController = context.getRouteController(DefaultRouteController.class);
                            routeController.stopRoute("queueProcessingRoute");
                            System.out.println("队列已清空,路由已停止");
                        }
                    })
                    .routeId("queueProcessingRoute");
            }
        });
        
        context.start();
        // 保持进程运行,直到路由停止
        Thread.sleep(Long.MAX_VALUE);
        context.stop();
    }
}

二、检查队列空与停止路由的核心方法

1. 查询队列剩余消息数

不同消息中间件的查询方式略有不同:

  • ActiveMQ:通过Queue.getBrowseableMessages().length获取当前队列的消息总数
  • RabbitMQ:使用AMQP客户端的QueueDeclareOk.getMessageCount()方法获取队列消息数
  • Kafka:可以通过AdminClient查询消费者组的未提交偏移量,结合主题分区消息数判断是否处理完成

2. 停止Camel路由

通过Camel的RouteController可以直接控制路由启停:

// 获取路由控制器实例
DefaultRouteController routeController = camelContext.getRouteController(DefaultRouteController.class);
// 根据路由ID停止指定路由
routeController.stopRoute("yourRouteId");

3. 更简洁的空闲检测方案

如果不需要精确检测队列是否为空,而是允许一段空闲时间后停止,可以使用Camel的idleConsumer配置:

from("jms:queue:myTargetQueue?idleConsumerTimeout=3000&idleConsumerStrategy=stop")
    .process(exchange -> {
        // 消息处理逻辑
    })
    .routeId("queueProcessingRoute");

这个配置会在队列连续3秒没有新消息时,自动停止当前消费者。如果需要完全停止路由,可以在空闲触发时结合RouteController调用停止方法。

内容的提问来源于stack exchange,提问作者nikolakoco

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.01 10:42:53