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

停止JMS/ActiveMQ入站集成流时如何强制提交事务

问题:Spring Integration动态集成流停止前事务提交与连接保持问题

背景

我们有一个动态注册(后续通过事件启动)的Spring Integration集成流,作用是作为ActiveMQ队列的独占消费者,读取一条消息后断开连接,后续可重新启动,通过控制总线启停流。

初始实现代码

private IntegrationFlow createTxxIntegrationFlow(Long taskId) {
        return IntegrationFlows 
            // INBOUND
                .from(Jms.messageDrivenChannelAdapter(jmsConnectionFactory)
                         .id("jmsInboundT" + taskId).autoStartup(false)
                         .destination(new ActiveMQQueue("T" + taskId + "?consumer.exclusive=true&consumer.prefetchSize=0")))

            // PROCESSINGS
                ...
                // several handlers and transformers... not relevant here
                
            // OUTBOUND
                .routeToRecipients(r -> r// write to ActiveMQ
                                         .recipientFlow(f -> f.handle(Jms.outboundAdapter(jmsConnectionFactory).destination(myDestination)))
                                         // stop only the inbound (using the control bus)
                                         //.recipientFlow(f -> f.handle((p, h) -> "@jmsInboundT"+taskId+".stop()").channel(controlChannel))
                                         // or stop the whole flow (using the control bus)
                                         .recipientFlow(f -> f.handle((p, h) -> "@T"+taskId+".stop()").channel(controlChannel))
                )
               .get();
}

初始问题现象

使用上述定义的流时,停止入站/流后会出现如下日志:

2022-12-16 11:19:20.601  WARN 6576 --- [erContainer#0-1] o.s.j.l.DefaultMessageListenerContainer  : Setup of JMS message listener invoker failed for destination 'queue://T65' - trying to recover. Cause: The Session is closed

能读取并处理一条ActiveMQ消息,然后停止流,但事务未提交,消息会重新回到队列中。

修改后的代码

为解决事务未提交问题,我们添加了事务管理器并手动提交事务:

// autowired
private final PlatformTransactionManager txManager;

public void handleTransaction() {
        var status = txManager.getTransaction(txDefinition);
        txManager.commit(status);
}

private IntegrationFlow createTxxIntegrationFlow(Long taskId) {
        return IntegrationFlows 
            // INBOUND
                .from(Jms.messageDrivenChannelAdapter(jmsConnectionFactory)
                         .id("jmsInboundT" + taskId).autoStartup(false)
                         .destination(new ActiveMQQueue("T" + taskId + "?consumer.exclusive=true&consumer.prefetchSize=0"))
                         .configureListenerContainer(c -> c.transactionManager(txManager)))

            // PROCESSINGS
                ...
                
            // OUTBOUND MESSAGE
                .routeToRecipients(r -> r// write to ActiveMQ
                                         .recipientFlow(f -> f.handle(Jms.outboundAdapter(jmsConnectionFactory).destination(myDestination)))
                                         // commit the transaction
                                         .recipientFlow(f -> f.handle(this, "handleTransaction"))
                                         // stop the inbound (using the control bus)
                                         //.recipientFlow(f -> f.handle((p, h) -> "@jmsInboundT"+taskId+".stop()").channel(controlChannel))
                                         // stop the whole flow (using the control bus)
                                         .recipientFlow(f -> f.handle((p, h) -> "@T"+taskId+".stop()").channel(controlChannel))
                )
               .get();
}

修改后新问题

修改后事务能正常提交,消息被消费且不会回到队列,入站/流也能停止,但应用会反复与Broker断开并重新连接:

2022-12-16 11:50:23.782  INFO 14908 --- [ActiveMQ Task-1] o.a.a.t.failover.FailoverTransport       : Successfully connected to tcp://localhost:61616
2022-12-16 11:50:25.852  INFO 14908 --- [ActiveMQ Task-1] o.a.a.t.failover.FailoverTransport       : Successfully connected to tcp://localhost:61616
2022-12-16 11:50:26.872  INFO 14908 --- [ActiveMQ Task-1] o.a.a.t.failover.FailoverTransport       : Successfully connected to tcp://localhost:61616
2022-12-16 11:50:27.891  INFO 14908 --- [ActiveMQ Task-1] o.a.a.t.failover.FailoverTransport       : Successfully connected to tcp://localhost:61616
2022-12-16 11:50:29.914  INFO 14908 --- [ActiveMQ Task-1] o.a.a.t.failover.FailoverTransport       : Successfully connected to tcp://localhost:61616

由于连接反复中断,我们无法再保持与Broker的独占连接。

提问

如何在停止流前有效强制提交事务,同时不先破坏连接?


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 08:30:58