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

Spring Boot 2.1.9→2.7.18升级后JMS Listener意外关闭问题排查

问题排查与解决方案:升级Spring Boot后JMS Topic监听因SSE连接断开停止

问题背景

项目从Java 8升级到Java 21,同时将Spring Boot从2.1.9.RELEASE升级至2.7.18(对应Spring Integration JMS从5.1.8.RELEASE升级到5.5.20)。项目通过SSE推送JMS Topic消息,升级后出现异常:旧的SSE连接断开后,新消息写入该断开连接失败时,会导致JMS消息监听端点停止,后续所有新连接都无法接收Topic消息。

关键代码片段

SSE控制器代码:

return Flux.from(noticesUsersTopicJmsReactiveSource())
        .filter(m -> filterOutMyMessages(m, userName))
        .map(Message::getPayload);

JMS响应式源配置:

return IntegrationFlows
        .from(Jms.messageDrivenChannelAdapter(noticesUsersTopicDefaultMessageListenerContainer()))
        .channel(MessageChannels.publishSubscribe())
        .toReactivePublisher();

DefaultMessageListenerContainer配置:

@Bean
public DefaultMessageListenerContainer noticesUsersTopicDefaultMessageListenerContainer() {
    DefaultMessageListenerContainer d = new DefaultMessageListenerContainer();
    d.setConnectionFactory(topicConnectionFactory);
    d.setPubSubDomain(true);
    d.setDestinationName(noticesUsersTopic);        
    return d;
}

版本差异表现

  • Spring Boot 2.1.9(Spring Integration 5.1.8):旧SSE连接断开后,消息写入失败仅关闭该连接,JMS监听不受影响,新连接可正常接收消息。
  • Spring Boot 2.7.18(Spring Integration 5.5.20):旧SSE连接写入失败时,JMS消息驱动端点被停止,日志显示ChannelPublishingJmsMessageListener$GatewayDelegate和JmsMessageDrivenEndpoint停止,后续无法接收Topic消息。

原因分析

  1. Reactive流错误传播机制变更:Spring Integration 5.5中,Reactive订阅者的错误(如SSE连接断开导致的写入失败)会向上传播至JMS消息驱动端点,触发端点停止逻辑;而5.1版本中错误被隔离在订阅者层面,不会影响上游JMS监听。
  2. PublishSubscribeChannel默认行为:默认情况下,PublishSubscribeChannel的单个订阅者抛出错误会影响整个通道的消息处理,若错误未被捕获,会传递到JMS适配器导致其停止。
  3. JMS端点错误处理缺失:未配置自定义错误处理策略,默认逻辑会在未捕获错误发生时停止端点。

解决方案

1. 在Flux层捕获错误,阻止向上传播

在控制器返回的Flux中添加错误处理,拦截SSE连接断开导致的写入错误,避免错误传递到JMS源:

return Flux.from(noticesUsersTopicJmsReactiveSource())
        .filter(m -> filterOutMyMessages(m, userName))
        .map(Message::getPayload)
        // 捕获错误并记录,不中断流
        .onErrorResume(e -> {
            log.warn("SSE消息发送失败: {}", e.getMessage());
            return Flux.empty();
        });

或使用onErrorContinue跳过出错消息,继续处理后续消息:

return Flux.from(noticesUsersTopicJmsReactiveSource())
        .filter(m -> filterOutMyMessages(m, userName))
        .map(Message::getPayload)
        .onErrorContinue((e, obj) -> {
            log.warn("消息{}的SSE推送失败: {}", obj, e.getMessage());
        });

2. 配置PublishSubscribeChannel忽略订阅者失败

自定义PublishSubscribeChannel并设置ignoreFailures为true,让单个订阅者的错误不影响其他订阅者和通道本身:

@Bean
public PublishSubscribeChannel noticesPublishSubscribeChannel() {
    PublishSubscribeChannel channel = new PublishSubscribeChannel();
    channel.setIgnoreFailures(true);
    return channel;
}

// 在IntegrationFlow中使用自定义通道
return IntegrationFlows
        .from(Jms.messageDrivenChannelAdapter(noticesUsersTopicDefaultMessageListenerContainer()))
        .channel(noticesPublishSubscribeChannel())
        .toReactivePublisher();

3. 为JMS端点配置错误处理

为JMS消息驱动通道适配器添加错误通道或自定义错误处理器,防止端点因错误停止:

// 配置错误通道
return IntegrationFlows
        .from(Jms.messageDrivenChannelAdapter(noticesUsersTopicDefaultMessageListenerContainer())
                .errorChannel("jmsErrorChannel"))
        .channel(MessageChannels.publishSubscribe())
        .toReactivePublisher();

// 错误通道处理流程
@Bean
public IntegrationFlow jmsErrorFlow() {
    return IntegrationFlows.from("jmsErrorChannel")
            .handle(message -> {
                Throwable error = (Throwable) message.getPayload();
                log.error("JMS消息处理错误: {}", error.getMessage(), error);
            })
            .get();
}

// 或为DefaultMessageListenerContainer设置ErrorHandler
@Bean
public DefaultMessageListenerContainer noticesUsersTopicDefaultMessageListenerContainer() {
    DefaultMessageListenerContainer d = new DefaultMessageListenerContainer();
    d.setConnectionFactory(topicConnectionFactory);
    d.setPubSubDomain(true);
    d.setDestinationName(noticesUsersTopic);
    d.setErrorHandler(t -> {
        log.error("JMS监听器错误: {}", t.getMessage(), t);
    });
    return d;
}

4. 配置JMS容器自动恢复

为DefaultMessageListenerContainer设置recoveryInterval,让容器在停止后自动尝试恢复:

d.setRecoveryInterval(5000); // 每5秒尝试恢复连接

验证步骤

实施方案后测试以下场景:

  1. 打开浏览器建立SSE连接(连接1),确认能接收Topic消息。
  2. 刷新浏览器建立新连接(连接2)。
  3. 发送Topic消息,确认连接2正常接收,连接1写入失败但JMS监听端点不停止。
  4. 再次发送Topic消息,确认连接2仍能接收,后续新连接也可正常接收消息。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 15:26:30