Spring Integration对接IBM MQ非阻塞消费无响应问题排查咨询
问题分析
- 每次调用
/event接口时都会重新执行jmsReactiveSource()方法,创建全新的IntegrationFlow实例和对应的JMS监听器。JMS消息驱动监听器默认只会消费启动后进入队列的消息,你先调用/pub发送消息再启动监听器,自然消费不到之前存入队列的消息,接口会一直挂起等待新消息。 - 流配置中使用的
MessageChannels.queue()是内存阻塞队列,监听器收到消息后会先存入该队列,若没有提前订阅就会导致消息积压,无法推给后续的反应式发布者。 - 默认的
Jms.messageDrivenChannelAdapter底层基于JMS规范的阻塞监听器实现,需要独占线程持续监听MQ,不符合你非阻塞降资源占用的需求。
修复方案
- 将
jmsReactiveSource注册为单例Bean,服务启动时就初始化并启动JMS监听器,提前接收队列中的消息:
@Bean public Publisher<Message<String>> jmsReactiveSource(ConnectionFactory ibmConnectionFactory) { return IntegrationFlows .from(Jms.messageDrivenChannelAdapter(ibmConnectionFactory) .destination("DEV.QUEUE.1") .autoStartup(true) .sessionAcknowledgeMode(Session.AUTO_ACKNOWLEDGE)) // 改用直通通道,消息直接向下游推送,不需要额外内存队列中转 .channel(MessageChannels.direct()) .log(org.springframework.integration.handler.LoggingHandler.Level.DEBUG) .toReactivePublisher(); }
- 调整
/event接口实现,注入单例的发布者实例,订阅后取第一条到达的消息返回:
@Autowired private Publisher<Message<String>> jmsReactiveSource; @GetMapping("/event") public Mono<String> getEvent() { return Flux.from(jmsReactiveSource) .next() .map(Message::getPayload); }
- 若需要完全的非阻塞实现,替换现有阻塞JMS客户端为IBM MQ官方反应式客户端
com.ibm.mq:mq-jakarta-reactive,完全适配Spring WebFlux的非阻塞线程模型,不需要额外的阻塞监听线程。
内容的提问来源于stack exchange,提问作者Don Donovan
相关产品推荐
相关产品推荐

