停止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
相关产品推荐
相关产品推荐

