Spring JMS与ActiveMQ Artemis:同步消息超时的临时队列管理问题
解决Artemis中JmsTemplate同步通信超时后无效临时队列堆积问题
针对你遇到的JmsTemplate#sendAndReceive超时后,监听器重建无消费者临时队列导致资源堆积的问题,结合Artemis的特性,提供以下几个可行方案:
方案1:监听器发送响应前检查目标队列状态
在监听器处理完请求准备发送响应时,先通过Artemis的客户端API查询响应队列(即生产者创建的临时队列)是否存在,或是否有活跃消费者。如果队列已被删除或无消费者,直接放弃发送响应。
示例代码:
@JmsListener(destination = "your-business-queue") public void processRequest(Message requestMsg, Session session) throws JMSException { // 1. 处理业务逻辑 String responseContent = handleBusinessLogic(requestMsg); // 2. 获取响应目标队列 Destination replyTo = requestMsg.getJMSReplyTo(); if (replyTo == null) { return; } String replyQueueName = ((Queue) replyTo).getQueueName(); // 3. 查询Artemis队列状态 try (ClientSession artemisSession = ((ActiveMQConnection) session.getConnection()).createSession()) { QueueQuery queueQuery = artemisSession.queueQuery(replyQueueName); // 队列不存在或无消费者时,终止响应发送 if (!queueQuery.isExists() || queueQuery.getConsumerCount() == 0) { return; } } // 4. 正常发送响应 Message responseMsg = session.createTextMessage(responseContent); try (MessageProducer producer = session.createProducer(replyTo)) { producer.send(responseMsg); } }
注意:需要确保监听器所在的连接有足够权限查询队列状态,可在Artemis的broker.xml中配置对应的权限。
方案2:生产者超时后主动通知监听器取消响应
在发送请求时附带唯一请求ID,生产者超时后发送"取消响应"指令到专门的控制队列,监听器维护已取消的请求ID集合,处理完业务后先校验该集合,若请求已被取消则不发送响应。
生产者端代码
public String sendRequestWithTimeout(String requestContent, long timeoutMillis) throws JMSException { String requestId = UUID.randomUUID().toString(); JmsTemplate jmsTemplate = getJmsTemplate(); try { return (String) jmsTemplate.sendAndReceive("your-business-queue", session -> { TextMessage requestMsg = session.createTextMessage(requestContent); requestMsg.setStringProperty("REQUEST_ID", requestId); return requestMsg; }); } catch (JmsTimeoutException e) { // 超时后发送取消指令 jmsTemplate.send("cancel-request-queue", session -> { Message cancelMsg = session.createMessage(); cancelMsg.setStringProperty("REQUEST_ID", requestId); return cancelMsg; }); throw e; } }
监听器端代码
// 用Guava Cache维护已取消的请求ID,自动过期清理 private final Cache<String, Boolean> cancelledRequests = CacheBuilder.newBuilder() .expireAfterWrite(1, TimeUnit.HOURS) // 过期时间根据业务处理时长调整 .concurrencyLevel(4) .build(); // 监听取消指令队列 @JmsListener(destination = "cancel-request-queue") public void handleCancelRequest(Message cancelMsg) throws JMSException { String requestId = cancelMsg.getStringProperty("REQUEST_ID"); if (requestId != null) { cancelledRequests.put(requestId, true); } } // 处理业务请求的监听器 @JmsListener(destination = "your-business-queue") public void processRequest(Message requestMsg, Session session) throws JMSException { String requestId = requestMsg.getStringProperty("REQUEST_ID"); if (requestId == null) { // 无请求ID的消息按原有逻辑处理 } // 处理业务逻辑 String responseContent = handleBusinessLogic(requestMsg); // 检查请求是否已被取消 if (cancelledRequests.getIfPresent(requestId) != null) { return; } // 发送响应 Destination replyTo = requestMsg.getJMSReplyTo(); if (replyTo != null) { TextMessage responseMsg = session.createTextMessage(responseContent); try (MessageProducer producer = session.createProducer(replyTo)) { producer.send(responseMsg); } } }
方案3:配置Artemis临时队列自动过期清理
利用Artemis的队列过期特性,给临时队列设置较短的过期时间,即使被监听器重建,也会在超时后被broker自动删除,避免长期堆积。
方式1:全局配置临时队列
修改Artemis的broker.xml,给所有临时地址(默认前缀为temp.)设置过期和自动删除:
<address-settings> <!-- 匹配所有临时地址 --> <address-setting match="temp.#"> <expiry-delay>30000</expiry-delay> <!-- 30秒后过期,根据业务处理时长调整 --> <auto-delete-queues>true</auto-delete-queues> <!-- 队列过期后自动删除 --> <auto-delete-addresses>true</auto-delete-addresses> <!-- 地址无队列时自动删除 --> </address-setting> </address-settings>
方式2:自定义JmsTemplate创建临时队列时设置过期
如果需要更细粒度的控制,可以自定义JmsTemplate的临时队列创建逻辑,设置过期时间:
public class CustomJmsTemplate extends JmsTemplate { @Override protected TemporaryQueue createTemporaryQueue(Session session) throws JMSException { TemporaryQueue tempQueue = super.createTemporaryQueue(session); // 转换为Artemis临时队列设置过期 if (tempQueue instanceof ActiveMQTemporaryQueue) { ((ActiveMQTemporaryQueue) tempQueue).setExpiryDelay(30000); // 30秒过期 } return tempQueue; } }
方案4:自定义同步请求响应逻辑,替代JmsTemplate#sendAndReceive
完全自定义同步通信流程,手动创建临时队列和消费者,超时后不仅关闭消费者,还显式删除临时队列,确保broker彻底清理资源,避免监听器后续重建。
示例代码:
public String customSendAndReceive(String requestContent, long timeoutMillis) throws JMSException, InterruptedException { ConnectionFactory connectionFactory = getConnectionFactory(); try (Connection connection = connectionFactory.createConnection()) { connection.start(); Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE); // 创建临时队列和消费者 TemporaryQueue replyQueue = session.createTemporaryQueue(); MessageConsumer consumer = session.createConsumer(replyQueue); // 发送请求 MessageProducer producer = session.createProducer(session.createQueue("your-business-queue")); TextMessage requestMsg = session.createTextMessage(requestContent); requestMsg.setJMSReplyTo(replyQueue); producer.send(requestMsg); producer.close(); // 等待响应,超时后清理资源 Message responseMsg = consumer.receive(timeoutMillis); if (responseMsg == null) { // 超时:关闭消费者并删除临时队列 consumer.close(); replyQueue.delete(); // 显式删除队列 throw new JmsTimeoutException("Request timed out"); } consumer.close(); replyQueue.delete(); return ((TextMessage) responseMsg).getText(); } }
内容的提问来源于stack exchange,提问作者Jonathan
相关产品推荐
相关产品推荐

