java - Spring-Integration 事务在管道其余部分之前提交

标签 java spring transactions spring-integration

所以,这是我当前的设置:

<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" />

图片是这样的:

enter image description here

我需要的是在任何消息转发到 int:brdige 或从 int:bridge(蓝色箭头 - 所以桥是仅当 jdbc 已提交时才应实际转发消息的边界组件)。

谢谢!

更新

这里是我为什么需要这种设置的描述:

用例:接收 amqp 消息,处理它并保存到 db 并将生成的 amqp 消息转发到管道。消息不应丢失(无论是到达消息还是传出消息,例如断电等)。 形成序列的多个消息可以到达具有相同组织的多个不同进程,如下所示。

我想如何解决它:

线程 1:

  1. 接收来自AMQP的消息并开始交易
  2. 在一系列服务激活器中处理 AMQP 消息(可能每个激活器都会写入数据库,但在提交之前不会保存更改)
  3. 将传出的 AMQP 消息写入数据库(完全序列化)- 在提交之前不会保存
  4. 提交
  5. 向 AMQP 确认传入消息
  6. 将传出消息转发到进程中的 THREAD2(不是 DB 指针,而是真实消息)

线程 2:

  1. 正在接收来自 THREAD1 的消息
  2. 尝试向AMQP发送消息
  3. 如果发送成功,则从 THREAD1 第 3 步中保存的数据库中删除传出的 AMQP 消息。

线程 3:

  1. 轮询数据库以获取要发送的消息(从 THREAD1 第 3 步开始)(每 10 秒)
  2. 如果发现任何新消息,将其标记为在下一次轮询中发送(同时 THREAD2 可以删除此消息)
  3. 如果在第二次轮询时(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/

相关文章:

java - Proguard混淆后如何保留@PathVariable参数名称

asp.net - ASP .Net WorldPay 集成

c# - 更新并发事务 C# 的数量问题

ruby-on-rails-3 - 即使事务回滚,Rails 也会在事务中保存记录?

java - 在maven元素中我如何指定文件的相对路径

Java不会在for循环内打印

java - 简单的 Java 游戏中的 MVC?

java - 尝试对空对象引用调用虚拟方法 'double android.location.Location.getLatitude()'

java - 在类型 'contextMessage' 的对象上找不到字段或属性 'org.springframework.webflow.engine.impl.RequestControlContextImpl'

java - 当应用程序没有 web.xml 时如何强制重定向到 HTTPS