不使用Delivery Delay为Spring JMS DefaultMessageListenerContainer添加消息消费延迟
消费端实现IBM MQ消息延迟处理方案
针对你无法使用MQ Delivery Delay、只能在消费阶段实现每条消息延迟处理的需求,这里提供两种可行的实现思路,基于DefaultMessageListenerContainer和IBM MQ消息的内置属性完成:
核心思路
IBM MQ的每条消息都会携带JMS_IBM_PutDate和JMS_IBM_PutTime属性,记录消息被放入队列的时间。我们可以通过这两个属性计算消息当前已在队列中停留的时长,再与设定的延迟时长对比,等待剩余时间后再执行业务逻辑。
方案1:在onMessage中直接等待(简单场景)
适合消息量不大、延迟时长较短的场景,直接在监听线程中等待剩余延迟时间后处理业务:
代码实现
import com.ibm.mq.jms.MQMessage; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import javax.jms.Message; import javax.jms.MessageListener; import java.time.Instant; import java.time.LocalDateTime; import java.time.ZoneId; import java.time.format.DateTimeFormatter; import java.time.temporal.ChronoUnit; public class DelayedMQMessageListener implements MessageListener { private static final Logger log = LoggerFactory.getLogger(DelayedMQMessageListener.class); // 设定延迟时长,示例为5分钟(单位:毫秒) private static final long DELAY_MILLIS = 5 * 60 * 1000; // MQ时间格式:putDate(yyyyMMdd) + putTime(HHmmssSSS) private static final DateTimeFormatter MQ_PUT_TIME_FORMATTER = DateTimeFormatter.ofPattern("yyyyMMddHHmmssSSS"); @Override public void onMessage(Message message) { try { MQMessage mqMessage = (MQMessage) message; // 获取消息的PUT时间属性 String putDate = mqMessage.getStringProperty("JMS_IBM_PutDate"); String putTime = mqMessage.getStringProperty("JMS_IBM_PutTime"); // 转换为Instant时间戳 String putDateTimeStr = putDate + putTime; LocalDateTime putLocalDateTime = LocalDateTime.parse(putDateTimeStr, MQ_PUT_TIME_FORMATTER); Instant putInstant = putLocalDateTime.atZone(ZoneId.systemDefault()).toInstant(); // 计算剩余等待时间 Instant now = Instant.now(); long elapsedMillis = ChronoUnit.MILLIS.between(putInstant, now); long waitMillis = DELAY_MILLIS - elapsedMillis; if (waitMillis > 0) { log.info("消息[{}]需等待{}毫秒后处理", mqMessage.getMessageID(), waitMillis); Thread.sleep(waitMillis); } // 执行业务处理逻辑 processMessage(message); // 手动确认消息(如果使用CLIENT_ACKNOWLEDGE模式) message.acknowledge(); } catch (InterruptedException e) { // 恢复线程中断状态,避免中断信号丢失 Thread.currentThread().interrupt(); log.error("消息等待延迟时被中断", e); // 根据业务需求决定是否放弃处理或触发重投 } catch (Exception e) { log.error("处理MQ消息失败", e); // 处理业务异常,比如触发MQ消息重投 } } private void processMessage(Message message) { // 你的业务处理逻辑 } }
容器配置调整
为避免等待阻塞导致消息堆积,需调整DefaultMessageListenerContainer的并发线程数:
import org.springframework.jms.listener.DefaultMessageListenerContainer; import javax.jms.ConnectionFactory; @Bean public DefaultMessageListenerContainer mqListenerContainer(ConnectionFactory connectionFactory) { DefaultMessageListenerContainer container = new DefaultMessageListenerContainer(); container.setConnectionFactory(connectionFactory); container.setDestinationName("YOUR_TARGET_QUEUE"); container.setMessageListener(new DelayedMQMessageListener()); // 设置并发消费者数量,根据消息量和延迟时长调整 container.setConcurrentConsumers(5); container.setMaxConcurrentConsumers(10); // 使用客户端确认模式,确保消息处理完成后再确认 container.setSessionAcknowledgeMode(javax.jms.Session.CLIENT_ACKNOWLEDGE); return container; }
方案2:异步延迟处理(高并发场景)
如果消息量较大、延迟时长较长,直接阻塞监听线程会导致队列堆积,此时可以将消息转发到异步延迟任务中处理,释放监听线程:
代码实现
import com.ibm.mq.jms.MQMessage; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; import javax.jms.Message; import javax.jms.MessageListener; import java.time.Instant; import java.time.LocalDateTime; import java.time.ZoneId; import java.time.format.DateTimeFormatter; import java.time.temporal.ChronoUnit; import java.util.concurrent.TimeUnit; public class AsyncDelayedMQMessageListener implements MessageListener { private static final Logger log = LoggerFactory.getLogger(AsyncDelayedMQMessageListener.class); private static final long DELAY_MILLIS = 5 * 60 * 1000; private static final DateTimeFormatter MQ_PUT_TIME_FORMATTER = DateTimeFormatter.ofPattern("yyyyMMddHHmmssSSS"); private final ThreadPoolTaskExecutor delayTaskExecutor; public AsyncDelayedMQMessageListener(ThreadPoolTaskExecutor delayTaskExecutor) { this.delayTaskExecutor = delayTaskExecutor; } @Override public void onMessage(Message message) { try { MQMessage mqMessage = (MQMessage) message; String putDate = mqMessage.getStringProperty("JMS_IBM_PutDate"); String putTime = mqMessage.getStringProperty("JMS_IBM_PutTime"); LocalDateTime putLocalDateTime = LocalDateTime.parse(putDate + putTime, MQ_PUT_TIME_FORMATTER); Instant putInstant = putLocalDateTime.atZone(ZoneId.systemDefault()).toInstant(); Instant now = Instant.now(); long elapsedMillis = ChronoUnit.MILLIS.between(putInstant, now); long waitMillis = DELAY_MILLIS - elapsedMillis; if (waitMillis <= 0) { // 已达到延迟时长,直接处理 processMessage(message); message.acknowledge(); return; } // 提交异步延迟任务,释放监听线程 delayTaskExecutor.schedule(() -> { try { processMessage(message); message.acknowledge(); } catch (Exception e) { log.error("延迟处理消息[{}]失败", mqMessage.getMessageID(), e); // 处理失败逻辑,比如触发重投 } }, waitMillis, TimeUnit.MILLISECONDS); } catch (Exception e) { log.error("转发消息到延迟任务失败", e); // 处理转发异常,比如拒绝消息让MQ重投 } } private void processMessage(Message message) { // 你的业务处理逻辑 } }
异步线程池配置
import org.springframework.context.annotation.Bean; import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; import java.util.concurrent.ThreadPoolExecutor; @Bean public ThreadPoolTaskExecutor delayTaskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(10); executor.setMaxPoolSize(20); executor.setQueueCapacity(100); executor.setThreadNamePrefix("delay-task-"); // 任务拒绝策略:当队列满时,由调用线程处理(避免消息丢失) executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); executor.initialize(); return executor; }
关键注意事项
- 属性可用性:确保IBM MQ JMS客户端已启用
JMS_IBM_PutDate和JMS_IBM_PutTime属性,默认情况下这些属性是自动暴露的,无需额外配置。 - 时间时区:转换时间时需注意服务器时区与MQ服务器时区一致,避免时间计算偏差。
- 消息重投:如果处理过程中发生异常,MQ可能会重新投递消息,此时消息的PUT时间仍是原始时间,等待逻辑会自动基于原始时间计算,符合需求。
- 线程资源:方案1中线程等待会占用监听线程,需根据消息量和延迟时长合理设置并发数;方案2的异步线程池需根据任务量调整参数,避免任务堆积。
内容的提问来源于stack exchange,提问作者Plainsage
相关产品推荐
相关产品推荐

