如何监控Spring Integration中Messaging Gateway的JmsOutboundGateway
监控Spring Integration中JmsOutboundGateway的实际调用情况
针对你将多个JmsOutboundGateway绑定到同一请求通道的场景,以下几种实用方案可实现监控实际处理消息的网关,同时能快速定位MQ连接异常的网关:
方案1:给每个网关添加请求处理通知(RequestHandlerAdvice)
通过自定义RequestHandlerAdvice,在网关处理消息的前后埋点,记录网关标识、处理状态,还能直接捕获连接异常。
步骤1:实现自定义监控通知
public class GatewayMonitoringAdvice extends AbstractRequestHandlerAdvice { private final String gatewayId; public GatewayMonitoringAdvice(String gatewayId) { this.gatewayId = gatewayId; } @Override protected Object doInvoke(ExecutionCallback callback, Object target, Message<?> message) throws Exception { // 记录网关开始处理消息 System.out.printf("网关[%s]开始处理消息,内容:%s%n", gatewayId, message.getPayload()); try { Object result = callback.execute(); // 处理成功时,给响应消息添加网关标识头 if (result instanceof Message<?>) { return MessageBuilder.fromMessage((Message<?>) result) .setHeader("processedByGateway", gatewayId) .build(); } System.out.printf("网关[%s]处理消息成功%n", gatewayId); return result; } catch (Exception e) { // 捕获异常,比如MQ连接断开的错误 System.err.printf("网关[%s]处理消息失败,异常:%s%n", gatewayId, e.getMessage()); throw e; } } }
步骤2:给每个JmsOutboundGateway绑定通知
修改网关Bean配置,通过adviceChain指定对应的监控通知:
@Bean @ServiceActivator(inputChannel = "inputLocalChannel", adviceChain = "firstGatewayAdvice") public JmsOutboundGateway firstGateway(){ JmsOutboundGateway gateway = new JmsOutboundGateway(); // 原有配置(如setConnectionFactory、setRequestDestination等) return gateway; } @Bean public GatewayMonitoringAdvice firstGatewayAdvice() { return new GatewayMonitoringAdvice("firstGateway"); } @Bean @ServiceActivator(inputChannel = "inputLocalChannel", adviceChain = "secondGatewayAdvice") public JmsOutboundGateway secondGateway(){ JmsOutboundGateway gateway = new JmsOutboundGateway(); // 原有配置 return gateway; } @Bean public GatewayMonitoringAdvice secondGatewayAdvice() { return new GatewayMonitoringAdvice("secondGateway"); }
步骤3:在回调中获取网关标识
在自定义回调里,从响应消息的Header中拿到处理网关的标识:
class CustomCallback implements ListenableFutureCallback<Message<String>> { @Override public void onSuccess(Message<String> result) { String gatewayId = result.getHeaders().get("processedByGateway", String.class); // 这里可将监控信息上报到监控系统(如Prometheus、ELK) System.out.printf("消息最终由网关[%s]处理完成%n", gatewayId); } @Override public void onFailure(Throwable ex) { // 结合通知里的异常日志,定位是哪个网关的连接出了问题 System.err.println("消息处理失败:" + ex.getMessage()); } }
方案2:自定义JmsOutboundGateway子类
通过继承JmsOutboundGateway,重写核心方法实现监控埋点:
public class MonitoredJmsOutboundGateway extends JmsOutboundGateway { private final String gatewayId; public MonitoredJmsOutboundGateway(String gatewayId) { this.gatewayId = gatewayId; } @Override protected Message<?> sendAndReceive(Message<?> requestMessage) throws MessagingException { System.out.printf("网关[%s]接收请求消息:%s%n", gatewayId, requestMessage.getPayload()); try { Message<?> response = super.sendAndReceive(requestMessage); if (response != null) { // 给响应添加网关标识 return MessageBuilder.fromMessage(response) .setHeader("processedByGateway", gatewayId) .build(); } System.out.printf("网关[%s]处理完成,无响应返回%n", gatewayId); return response; } catch (Exception e) { System.err.printf("网关[%s]处理请求失败,异常:%s%n", gatewayId, e.getMessage()); throw e; } } }
替换原有网关Bean的实现:
@Bean @ServiceActivator(inputChannel = "inputLocalChannel") public JmsOutboundGateway firstGateway(){ MonitoredJmsOutboundGateway gateway = new MonitoredJmsOutboundGateway("firstGateway"); // 原有配置 return gateway; } @Bean @ServiceActivator(inputChannel = "inputLocalChannel") public JmsOutboundGateway secondGateway(){ MonitoredJmsOutboundGateway gateway = new MonitoredJmsOutboundGateway("secondGateway"); // 原有配置 return gateway; }
方案3:利用消息Header传递网关标识
在每个网关的HeaderMapper中添加自定义标识,将网关ID写入消息头,后续从响应中读取:
@Bean @ServiceActivator(inputChannel = "inputLocalChannel") public JmsOutboundGateway firstGateway(){ JmsOutboundGateway gateway = new JmsOutboundGateway(); // 原有配置 // 自定义HeaderMapper,添加网关标识到JMS消息属性 gateway.setHeaderMapper(new DefaultJmsHeaderMapper() { @Override public void fromHeaders(MessageHeaders headers, jakarta.jms.Message jmsMessage) throws JMSException { super.fromHeaders(headers, jmsMessage); jmsMessage.setStringProperty("gatewayId", "firstGateway"); } @Override public MessageHeaders toHeaders(jakarta.jms.Message jmsMessage) throws JMSException { MessageHeaders headers = super.toHeaders(jmsMessage); // 将JMS属性映射回Spring Message Header return MessageHeaders.wrap( new HashMap<>(headers) {{ put("processedByGateway", jmsMessage.getStringProperty("gatewayId")); }} ); } }); return gateway; }
这种方式适合需要在JMS消息本身携带标识的场景,方便跨系统追踪。
内容的提问来源于stack exchange,提问作者Luke
相关产品推荐
相关产品推荐

