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

如何在Spring Integration中为消息添加分组ID并监控文件处理步骤

How to Add Correlation/Trace IDs to Spring Integration Messages for Logging

Great question! Tracking individual file processing flows in Spring Integration is super important for debugging and observability, and adding a consistent correlation/trace ID to each message is exactly the way to go. Here's how you can implement this effectively:

1. Inject a Unique Correlation ID into Message Headers

First, you need to attach a unique identifier to every message when it enters your integration flow. This ID will follow the message through all processing steps.

Option 1: Use Spring Integration's Built-in Correlation Header

Spring Integration provides a dedicated header for this purpose: IntegrationMessageHeaderAccessor.CORRELATION_ID. It's designed specifically for message correlation, so it's a standard, recommended choice.

You can inject this ID at the entry point of your flow (e.g., a file inbound adapter) using a transformer:

@Bean
public GenericTransformer<File, Message<File>> traceIdInjector() {
    return file -> MessageBuilder.withPayload(file)
            .setHeader(IntegrationMessageHeaderAccessor.CORRELATION_ID, UUID.randomUUID().toString())
            .build();
}

Option 2: Custom Trace ID Header

If you prefer a custom header name (like traceId), just replace the constant with your own string:

.setHeader("traceId", UUID.randomUUID().toString())

2. Log the ID Using MDC for Consistent Tracking

To make the ID appear in every log entry related to the message, use MDC (Mapped Diagnostic Context) from SLF4J. MDC lets you attach context data that's automatically included in all log statements for the current thread.

Create a Global Channel Interceptor

A global interceptor will handle MDC setup/teardown for every message passing through any channel in your flow:

@Bean
public GlobalChannelInterceptor traceIdLoggingInterceptor() {
    return new ChannelInterceptor() {
        @Override
        public Message<?> preSend(Message<?> message, MessageChannel channel) {
            // Get or generate the trace ID (fallback to a new UUID if missing)
            String traceId = message.getHeaders()
                    .getOrDefault(IntegrationMessageHeaderAccessor.CORRELATION_ID, UUID.randomUUID().toString())
                    .toString();
            
            // Add the ID to MDC so logs automatically include it
            MDC.put("traceId", traceId);
            
            // If the message didn't have the ID already, inject it to ensure consistency
            if (!message.getHeaders().containsKey(IntegrationMessageHeaderAccessor.CORRELATION_ID)) {
                return MessageBuilder.fromMessage(message)
                        .setHeader(IntegrationMessageHeaderAccessor.CORRELATION_ID, traceId)
                        .build();
            }
            return message;
        }

        @Override
        public void afterSendCompletion(Message<?> message, MessageChannel channel, boolean sent, Exception ex) {
            // Clean up MDC to avoid leaking context between different messages
            MDC.remove("traceId");
        }
    };
}

Update Your Logging Configuration

Modify your logger configuration (e.g., Logback) to include the traceId from MDC. Here's an example logback.xml pattern:

<appender name="CONSOLE" class="ch.qos.logback.core.ConsoleAppender">
    <encoder>
        <pattern>%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{36} - *Trace ID: %X{traceId}* - %msg%n</pattern>
    </encoder>
</appender>

3. Handle Async Processing (If Applicable)

If your flow uses asynchronous channels (e.g., ExecutorChannel), MDC context won't automatically transfer to the async thread by default. To fix this, wrap your task executor to carry over MDC data:

@Bean
public AsyncTaskExecutor mdcAwareTaskExecutor() {
    return new AsyncTaskExecutor() {
        private final ExecutorService delegate = Executors.newFixedThreadPool(10);

        @Override
        public void execute(Runnable task) {
            Map<String, String> mdcContext = MDC.getCopyOfContextMap();
            delegate.execute(() -> {
                if (mdcContext != null) {
                    MDC.setContextMap(mdcContext);
                }
                try {
                    task.run();
                } finally {
                    MDC.clear();
                }
            });
        }

        // Implement other AsyncTaskExecutor methods (submit, etc.) using the same pattern
    };
}

4. Query Logs by Trace ID

Once everything is set up, you can use your log analysis tool (ELK, Splunk, etc.) to search for a specific traceId value. This will return all log entries related to that single file/message's processing flow, from start to finish—making it easy to debug issues or verify end-to-end processing.


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 03:54:32