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

不使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 12:05:53