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

如何监控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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 12:44:51