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

Spring JMS自定义事件监听器:队列消息增减监控实现问询

实现JMS队列入队/出队事件监控(结合Prometheus)

1. 定义自定义事件类型

先创建携带队列名称的基础事件类,以及入队、出队的具体事件子类:

// 基础队列事件类,封装队列名称信息
public class QueueMessageEvent extends ApplicationEvent {
    private final String queueName;

    public QueueMessageEvent(Object source, String queueName) {
        super(source);
        this.queueName = queueName;
    }

    public String getQueueName() {
        return queueName;
    }
}

// 消息入队事件
public class MessageEnqueuedEvent extends QueueMessageEvent {
    public MessageEnqueuedEvent(Object source, String queueName) {
        super(source, queueName);
    }
}

// 消息出队事件
public class MessageDequeuedEvent extends QueueMessageEvent {
    public MessageDequeuedEvent(Object source, String queueName) {
        super(source, queueName);
    }
}

2. 发布入队/出队事件

2.1 封装消息发送逻辑,触发入队事件

避免直接使用JmsTemplate,封装发送方法,消息发送完成后立即发布入队事件:

@Component
public class QueueMessageSender {
    private final JmsTemplate jmsTemplate;
    private final ApplicationEventPublisher eventPublisher;

    // 构造注入依赖
    public QueueMessageSender(JmsTemplate jmsTemplate, ApplicationEventPublisher eventPublisher) {
        this.jmsTemplate = jmsTemplate;
        this.eventPublisher = eventPublisher;
    }

    public void sendToQueue(String queueName, Object payload) {
        jmsTemplate.convertAndSend(queueName, payload);
        // 发布入队事件
        eventPublisher.publishEvent(new MessageEnqueuedEvent(this, queueName));
    }
}

2.2 监听消费完成,触发出队事件

有两种实现方式,按需选择:

方式一:在@JmsListener方法中手动触发

处理完消息业务逻辑后,直接发布出队事件:

@Component
public class QueueMessageConsumer {
    private final ApplicationEventPublisher eventPublisher;

    public QueueMessageConsumer(ApplicationEventPublisher eventPublisher) {
        this.eventPublisher = eventPublisher;
    }

    @JmsListener(destination = "RETRY_QUEUE")
    public void processRetryMessage(Message message) {
        // 业务处理逻辑
        // ...

        // 处理成功后发布出队事件
        eventPublisher.publishEvent(new MessageDequeuedEvent(this, "RETRY_QUEUE"));
    }

    @JmsListener(destination = "SUCCESS_QUEUE")
    public void processSuccessMessage(Message message) {
        // 业务处理逻辑
        // ...

        eventPublisher.publishEvent(new MessageDequeuedEvent(this, "SUCCESS_QUEUE"));
    }
}

方式二:用AOP批量拦截@JmsListener方法

如果有大量监听方法,用AOP可以避免重复代码,仅在方法成功执行后触发事件:

@Aspect
@Component
public class JmsListenerCompletionAspect {
    private final ApplicationEventPublisher eventPublisher;

    public JmsListenerCompletionAspect(ApplicationEventPublisher eventPublisher) {
        this.eventPublisher = eventPublisher;
    }

    @AfterReturning("@annotation(org.springframework.jms.annotation.JmsListener)")
    public void afterListenerExecution(JoinPoint joinPoint) {
        // 获取当前监听方法绑定的队列名称
        MethodSignature signature = (MethodSignature) joinPoint.getSignature();
        JmsListener jmsListener = signature.getMethod().getAnnotation(JmsListener.class);
        String queueName = jmsListener.destination();

        // 发布出队事件
        eventPublisher.publishEvent(new MessageDequeuedEvent(this, queueName));
    }
}

3. 配置Prometheus Gauge指标

创建带队列标签的Gauge,区分不同队列的监控指标:

@Configuration
public class QueueMetricsConfig {
    @Bean
    public Map<String, Gauge> queueSizeGauges(MeterRegistry meterRegistry) {
        Map<String, Gauge> gaugeMap = new HashMap<>();
        // 初始化业务中用到的队列Gauge
        List<String> queueNames = Arrays.asList("RETRY_QUEUE", "SUCCESS_QUEUE");
        
        queueNames.forEach(queueName -> {
            Gauge gauge = Gauge.builder("jms_queue_size", "当前JMS队列消息数量")
                    .tag("queue", queueName)
                    .register(meterRegistry);
            gaugeMap.put(queueName, gauge);
        });
        
        return gaugeMap;
    }
}

4. 监听事件更新Gauge值

创建事件监听器,根据入队/出队事件更新对应队列的Gauge计数:

@Component
public class QueueMetricsUpdater {
    private final Map<String, Gauge> queueSizeGauges;

    public QueueMetricsUpdater(Map<String, Gauge> queueSizeGauges) {
        this.queueSizeGauges = queueSizeGauges;
    }

    @EventListener
    public void onMessageEnqueued(MessageEnqueuedEvent event) {
        updateGauge(event.getQueueName(), true);
    }

    @EventListener
    public void onMessageDequeued(MessageDequeuedEvent event) {
        updateGauge(event.getQueueName(), false);
    }

    private void updateGauge(String queueName, boolean isEnqueue) {
        Gauge gauge = queueSizeGauges.get(queueName);
        if (gauge != null) {
            if (isEnqueue) {
                gauge.inc();
            } else {
                gauge.dec();
            }
        }
    }
}

注意事项

  • 增量更新为近似值:如果队列存在初始消息、或外部系统直接操作队列,指标会和实际队列大小有偏差。若需精确值,需调用JMS服务器的API(如ActiveMQ的QueueBrowser)定时拉取。
  • 异常处理:如果消息消费失败抛出异常,@AfterReturning不会触发,因此不会执行dec(),符合“消息未被移除则不减少计数”的逻辑。
  • 动态队列支持:如果有动态创建的队列,可在发送消息时动态初始化对应的Gauge,避免gaugeMap中找不到对应条目。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 17:27:14