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
相关产品推荐
相关产品推荐

