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

Spring Integration:绑定事务的IntegrationFlow Poller致TraceId变更问题

IntegrationFlow Poller绑定事务时TraceId异常变更问题

问题现象

为Spring Integration的IntegrationFlow Poller绑定事务后,出现TraceId在事务提交前后不一致的情况,日志显示事务内的TraceId与Poller初始的TraceId不同:

2025-01-05T21:33:21.726+04:00[ INFO 1396 --- [transaction-sync-traceid] [ scheduling-1] [677ac2610ebc3f65f920ea37a7aa2cb8-f920ea37a7aa2cb8] [0;39m[36mc.i.s.c.PollerWithTransactionFlowConfig :Supplied Message is: Good Day
2025-01-05T21:33:21.733+04:00[ INFO 1396 --- [transaction-sync-traceid] [ scheduling-1] [677ac26181d6b7f5f53ebcee6f7e914a-f53ebcee6f7e914a] [0;39m[36mc.i.sample.config.TransactionConfig :Transaction Committed, but traceId is different here !!

根因分析

Spring Integration的Poller在绑定事务时,默认的事务同步机制没有传递Poller任务本身的Trace上下文到事务执行流程中。事务启动时会重新生成新的TraceId,导致前后TraceId不一致。

解决方案

方案1:自定义事务同步器传递Trace上下文

通过自定义TransactionSynchronization,在事务开启前保存当前Trace上下文,事务执行阶段恢复该上下文,确保TraceId一致:

import org.springframework.transaction.support.TransactionSynchronizationAdapter;
import brave.Tracer;
import brave.Span;

public class TracePropagationTransactionSync extends TransactionSynchronizationAdapter {
    private final Span currentSpan;

    public TracePropagationTransactionSync(Span currentSpan) {
        this.currentSpan = currentSpan;
    }

    @Override
    public void beforeCommit(boolean readOnly) {
        Tracer.current().withSpan(currentSpan);
    }
}

在Poller配置中注入该同步器:

@Bean
public IntegrationFlow pollerWithTransactionFlow(PlatformTransactionManager transactionManager, Tracer tracer) {
    return IntegrationFlows.from(() -> "Good Day",
            spec -> spec.poller(Pollers.fixedDelay(1000)
                    .transactional(transactionManager)
                    .transactionSynchronization(new TracePropagationTransactionSync(tracer.currentSpan()))))
            .handle(msg -> {
                // 业务逻辑处理
            })
            .get();
}

方案2:使用Trace感知的任务执行器

借助Spring Cloud Sleuth提供的TraceTaskExecutor包装Poller的任务执行器,确保Trace上下文在事务线程中传播:

import org.springframework.cloud.sleuth.instrument.async.TraceTaskExecutor;
import java.util.concurrent.Executors;

@Bean
public TaskExecutor traceAwareTaskExecutor() {
    return new TraceTaskExecutor(Executors.newFixedThreadPool(5));
}

@Bean
public IntegrationFlow pollerFlow(PlatformTransactionManager transactionManager) {
    return IntegrationFlows.from(() -> "Good Day",
            spec -> spec.poller(Pollers.fixedDelay(1000)
                    .taskExecutor(traceAwareTaskExecutor())
                    .transactional(transactionManager)))
            .handle(msg -> {
                // 业务逻辑处理
            })
            .get();
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 03:26:09