所以,这是我当前的设置:
<int-amqp:inbound-channel-adapter channel="input-channel" queue-names="probni" message-converter="jsonMessageConverter"
channel-transacted="true"
transaction-manager="dataSourceTransactionManager"/>
<int:chain input-channel="input-channel" output-channel="inputc1">
<int:service-activator ref="h1" method="handle" />
<int:service-activator ref="h2" method="handle" />
<int:service-activator ref="h3" method="handle" />
<int:splitter />
</int:chain>
<int:publish-subscribe-channel id="inputc1"/>
<int:claim-check-in input-channel="inputc1" output-channel="nullChannel" message-store="messageStore" order="1" />
<int:bridge input-channel="inputc1" output-channel="inputc2" order="2" />
<int:publish-subscribe-channel id="inputc2" task-executor="taskExecutor" />
<int-amqp:outbound-channel-adapter channel="inputc2" exchange-name="exch" amqp-template="rabbitTemplate" order="1" />
<int:service-activator input-channel="inputc2" output-channel="nullChannel"
expression="@messageStore.removeMessage(headers['id'])" order="2" />
图片是这样的:
我需要的是在任何消息转发到 int:brdige 或从 int:bridge(蓝色箭头 - 所以桥是仅当 jdbc 已提交时才应实际转发消息的边界组件)。
谢谢!
更新
这里是我为什么需要这种设置的描述:
用例:接收 amqp 消息,处理它并保存到 db 并将生成的 amqp 消息转发到管道。消息不应丢失(无论是到达消息还是传出消息,例如断电等)。 形成序列的多个消息可以到达具有相同组织的多个不同进程,如下所示。
我想如何解决它:
线程 1:
- 接收来自AMQP的消息并开始交易
- 在一系列服务激活器中处理 AMQP 消息(可能每个激活器都会写入数据库,但在提交之前不会保存更改)
- 将传出的 AMQP 消息写入数据库(完全序列化)- 在提交之前不会保存
- 提交
- 向 AMQP 确认传入消息
- 将传出消息转发到进程中的 THREAD2(不是 DB 指针,而是真实消息)
线程 2:
- 正在接收来自 THREAD1 的消息
- 尝试向AMQP发送消息
- 如果发送成功,则从 THREAD1 第 3 步中保存的数据库中删除传出的 AMQP 消息。
线程 3:
- 轮询数据库以获取要发送的消息(从 THREAD1 第 3 步开始)(每 10 秒)
- 如果发现任何新消息,将其标记为在下一次轮询中发送(同时 THREAD2 可以删除此消息)
- 如果在第二次轮询时(20 秒后)消息仍然存在 - 意味着 THREAD2 未能发送它,因此我们在这里生成新线程来执行与 THREAD2 相同的任务。
更新 2
试过这个设置,但有一些问题:
<int:transaction-synchronization-factory id="transactionSynchronizationFactory">
<int:after-commit expression="payload" channel="committed-channel" />
</int:transaction-synchronization-factory>
<int-amqp:inbound-channel-adapter channel="input-channel" queue-names="probni" message-converter="jsonMessageConverter"
channel-transacted="true"
transaction-manager="dataSourceTransactionManager" advice-chain="amqpMethodInterceptor"/>
和:
@Component
public class AmqpMethodInterceptor implements MethodInterceptor {
private TransactionSynchronizationFactory factory;
public AmqpMethodInterceptor(TransactionSynchronizationFactory factory){
this.factory = factory;
}
@Override
public Object invoke(MethodInvocation invocation) throws Throwable {
if (TransactionSynchronizationManager.isActualTransactionActive()) {
TransactionSynchronization synchronization = factory.create("123");
TransactionSynchronizationManager.registerSynchronization(synchronization);
}
Object result = invocation.proceed();
return result;
}
}
after-commit 被调用,但此时消息为空,所以似乎我没有什么可以转发到 committed-channel。知道如何做这部分吗?
最佳答案
只有 TransactionSynchronization
才有可能.
网桥实际上是 TX 的“边界”,但此处的提交确实发生在 send()
之后方法。只是因为你的下一个 channel 是 executor
, 因此,当前事务线程无事可做,只在发送之后执行提交,而不是之前。
为了您的目标,您应该实现 MethodInterceptor
建议注入(inject)<int-amqp:inbound-channel-adapter>
通过advice-chain
.并尝试利用 DefaultTransactionSynchronizationFactory
的逻辑与 ExpressionEvaluatingTransactionSynchronizationProcessor
,您将能够向 afterCommitChannel
发送消息.
您在 Advice
中的代码应该使用这个模板:
if (TransactionSynchronizationManager.isActualTransactionActive()) {
TransactionSynchronization synchronization = this.transactionSynchronizationFactory.create(key);
TransactionSynchronizationManager.registerSynchronization(synchronization);
}
哪里key
可以是任何唯一对象,用于区分与 TX 同步的资源。
更新
but the message is null at this point in time, so seems like I have nothing to forward to committed-channel.
这是真的,因为您使用了 TransactionSynchronizationFactory
以不同寻常的方式。
好吧,不管怎样,让我们尝试欺骗它,因为对我来说你走对了路。
factory.create("123");
这样做:
DefaultTransactionalResourceSynchronization synchronization = new DefaultTransactionalResourceSynchronization(key);
TransactionSynchronizationManager.bindResource(key, synchronization.getResourceHolder());
return synchronization;
重点是TransactionSynchronizationManager.bindResource()
.我的想法是在下游某个地方,在 TX 结束之前,这样做:
IntegrationResourceHolder holder =
(IntegrationResourceHolder) TransactionSynchronizationManager.getResource("123");
holder.setMessage(message);
我认为这甚至是可能的:
<int:outbound-channel-adapter
expression="T(org.springframework.transaction.support.TransactionSynchronizationManager).getResource('123').setMessage(#root)"/>
作为退出 TX 之前的最后一个端点。
关于java - Spring-Integration 事务在管道其余部分之前提交,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/41104769/