ActiveMQ Artemis管理队列消息数持续增长问题咨询
我用ActiveMQ Artemis管理API获取队列状态和消费者数量,因为缓存了ClientRequestor,发现“activemq.management.*”队列的消息数一直在涨,就算配置了过期延迟,消息也没被丢弃。我本来以为消息被消费后会消失,这是为啥?
private ServerLocator locator; private ClientSessionFactory defaultFactory; private ClientSession session; private ClientRequestor requestor; public ManagementHelper(String defaultURL) { this.locator = ActiveMQClient.createServerLocator(defaultURL); this.defaultFactory = locator.createSessionFactory(); this.session = factory.createSession(this.username, this.password, false, true, true, locator.isPreAcknowledge(), locator.getAckBatchSize()); this.requestor = new ClientRequestor(session, "activemq.management"); } @Scheduled(fixedRateString = "20000") public void getConsumerNbos(ClientSessionFactory factory, ServerLocator locator) { ClientMessage message = session.createMessage(false); ManagementHelper.putOperationInvocation(message, ResourceNames.BROKER, listAllConsumersAsJSON); session.start(); ClientMessage replyConsumer = requestor.request(message); String resultJSON = (String) ManagementHelper.getResult(replyConsumer, String.class); ClientMessage message2 = session.createMessage(false); ManagementHelper.putOperationInvocation(message2, ResourceNames.BROKER, MANAGEMENT_OPERATION_QUEUES); ClientMessage replyQueueNames = requestor.request(message2); Object[] objQueueNames = (Object[]) ManagementHelper.getResult(replyQueueNames); }
原因分析
回复消息未确认:你通过
ClientRequestor.request()拿到回复消息后,没有对这些消息执行确认操作。即便会话配置了autoCommitAcks=true,部分场景下如果不显式触发确认(比如调用acknowledge()),消息不会被自动标记为已消费,会滞留在ClientRequestor创建的临时队列(即你看到的activemq.management.*队列)中。会话参数顺序错误:创建
ClientSession时参数顺序有误,导致会话的预确认(preAcknowledge)和事务(xa)模式配置不符合预期。错误的配置会让服务器无法正确识别消息的消费状态,进而无法清理已处理的消息。重复启动会话:定时任务中每次都调用
session.start(),但会话仅需启动一次即可正常工作。重复启动虽不直接导致堆积,但会引发不必要的资源操作,可能间接影响消息处理流程。
解决方案
确认回复消息:处理完回复结果后,显式确认消息,确保服务器清理已消费的消息:
ClientMessage replyConsumer = requestor.request(message); String resultJSON = (String) ManagementHelper.getResult(replyConsumer, String.class); replyConsumer.acknowledge(); // 确认回复消息 // 处理第二个请求的回复 ClientMessage replyQueueNames = requestor.request(message2); Object[] objQueueNames = (Object[]) ManagementHelper.getResult(replyQueueNames); replyQueueNames.acknowledge(); // 确认回复消息修正会话参数顺序:按照API定义的正确顺序创建
ClientSession,确保配置符合预期:this.session = defaultFactory.createSession( this.username, this.password, false, // autoCommitSends true, // autoCommitAcks locator.isPreAcknowledge(), // preAcknowledge false, // xa(无需XA事务时设为false) locator.getAckBatchSize() );移除重复的会话启动:将
session.start()移到构造方法中,仅初始化时启动一次:public ManagementHelper(String defaultURL) { this.locator = ActiveMQClient.createServerLocator(defaultURL); this.defaultFactory = locator.createSessionFactory(); this.session = defaultFactory.createSession(this.username, this.password, false, true, locator.isPreAcknowledge(), false, locator.getAckBatchSize()); this.requestor = new ClientRequestor(session, "activemq.management"); session.start(); // 仅初始化时启动一次 }清理冗余参数:定时任务
getConsumerNbos中的factory和locator参数未被使用,建议移除,避免代码混淆。
内容的提问来源于stack exchange,提问作者Cat

