Spring Integration:绑定事务的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

