Spring Integration带Poller的QueueChannel链路追踪传播问题
Spring Integration带Poller的QueueChannel链路追踪全链路传播解决方案
问题场景
基于Spring Integration 6.1.4、Spring Boot 3.1.5、Micrometer Tracing(Brave桥接)1.1.6的集成流程中:
- HTTP请求触发MessagingGateway,经
sftpListingChannel获取SFTP文件列表 - 拆分列表后发送到带Poller的
sftpFetchingChannel执行文件下载 - 存在两个核心问题:
- Poller执行时会终止初始Trace ID,生成新ID,导致全链路追踪断裂
- 拆分后的子操作报错时,
errorChannel处理错误使用新Trace ID,无法关联原HTTP请求链路
- 已尝试
ObservationPropagationChannelInterceptor解决普通通道的追踪传播,但errorChannel仍存在追踪断裂或无追踪的问题
针对性解决方案
一、修复Poller线程的追踪上下文传播
Poller默认使用的TaskExecutor不会自动传播Micrometer观测上下文,必须替换为支持传播的实现,并增强通道拦截:
配置
ObservationPropagationTaskExecutor
包装自定义线程池,确保线程切换时传递追踪上下文:@Bean public TaskExecutor sftpFetchTaskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(5); executor.setMaxPoolSize(10); executor.setQueueCapacity(25); executor.setThreadNamePrefix("sftp-fetch-"); executor.initialize(); return new ObservationPropagationTaskExecutor(executor); }给Poller配置该Executor并添加观测传播通知
确保Poller执行全程保留追踪上下文:@Bean public ObservationPropagationAdvice observationPropagationAdvice() { return new ObservationPropagationAdvice(); } @Bean public PollerMetadata sftpFetchPoller() { return Pollers.fixedDelay(Duration.ofSeconds(1)) .taskExecutor(sftpFetchTaskExecutor()) .advice(observationPropagationAdvice()) .get(); }给所有业务通道添加
ObservationPropagationChannelInterceptor
确保消息在通道间传递时,追踪上下文被正确携带:@Bean public ChannelInterceptor observationPropagationInterceptor() { return new ObservationPropagationChannelInterceptor(); } @Bean public MessageChannel sftpListingChannel() { DirectChannel channel = new DirectChannel(); channel.addInterceptor(observationPropagationInterceptor()); return channel; } @Bean public QueueChannel sftpFetchingChannel() { QueueChannel channel = new QueueChannel(); channel.addInterceptor(observationPropagationInterceptor()); return channel; }
二、解决errorChannel的追踪关联问题
errorChannel的消息默认不会自动继承Poller线程的追踪上下文,需要手动处理上下文传递:
给errorChannel添加观测传播拦截器
确保错误消息处理时能恢复追踪上下文:@Bean public MessageChannel errorChannel() { DirectChannel channel = new DirectChannel(); channel.addInterceptor(observationPropagationInterceptor()); // 自定义错误处理逻辑 channel.subscribe(message -> { Throwable error = (Throwable) message.getPayload(); String traceId = message.getHeaders().get(TraceMessageHeaders.TRACE_ID, String.class); log.error("SFTP下载任务失败,关联Trace ID: {}", traceId, error); }); return channel; }配置Poller的ErrorHandler,手动传递追踪上下文到errorChannel
当Poller线程抛出异常时,将当前追踪信息附加到错误消息头:@Bean public ErrorHandler observationAwareErrorHandler(Tracer tracer, MessageChannel errorChannel) { return throwable -> { Span currentSpan = tracer.currentSpan(); if (currentSpan != null) { Message<Throwable> errorMessage = MessageBuilder.withPayload(throwable) .setHeader(TraceMessageHeaders.TRACE_ID, currentSpan.context().traceId()) .setHeader(TraceMessageHeaders.SPAN_ID, currentSpan.context().spanId()) .build(); errorChannel.send(errorMessage); } else { errorChannel.send(MessageBuilder.withPayload(throwable).build()); } }; }将自定义ErrorHandler绑定到Poller
@Bean public PollerMetadata sftpFetchPoller(ErrorHandler observationAwareErrorHandler) { return Pollers.fixedDelay(Duration.ofSeconds(1)) .taskExecutor(sftpFetchTaskExecutor()) .advice(observationPropagationAdvice()) .errorHandler(observationAwareErrorHandler) .get(); }
三、验证步骤
- 触发HTTP请求,查看日志中Trace ID是否在
sftpListingChannel->sftpFetchingChannel->errorChannel(报错时)全程一致 - 检查
/actuator/trace端点(需开启actuator),确认链路完整 - 验证错误日志中的Trace ID能关联到原HTTP请求的Trace ID
内容的提问来源于stack exchange,提问作者asceta
相关产品推荐
相关产品推荐

