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

如何在@KafkaListener方法执行前获取含Header的Kafka消息?

在@KafkaListener执行前获取Kafka消息Header并处理MDC的方案

方案一:直接在监听方法参数中获取(最简单直接)

直接在@KafkaListener标注的方法里接收ConsumerRecord参数,在业务逻辑执行前处理Header和MDC。这种方式无需额外配置,适合简单场景。

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.slf4j.MDC;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;
import java.util.UUID;

@Component
public class KafkaMessageListener {

    @KafkaListener(topics = "your_topic")
    public void handleMessage(ConsumerRecord<String, String> record) {
        // 先处理Header写入MDC
        processHeadersToMDC(record.headers());
        
        // 业务逻辑代码
        System.out.println("处理消息:" + record.value());
        
        // 线程复用场景下,务必清理MDC避免污染
        MDC.clear();
    }

    private void processHeadersToMDC(org.apache.kafka.common.header.Headers headers) {
        // 获取或生成correlationId
        String correlationId = getHeaderValue(headers, "correlationId");
        if (correlationId == null) {
            correlationId = UUID.randomUUID().toString();
        }
        // 获取或生成requestId
        String requestId = getHeaderValue(headers, "requestId");
        if (requestId == null) {
            requestId = UUID.randomUUID().toString();
        }
        // 写入MDC
        MDC.put("correlationId", correlationId);
        MDC.put("requestId", requestId);
    }

    private String getHeaderValue(org.apache.kafka.common.header.Headers headers, String headerName) {
        return headers.lastHeader(headerName) != null 
                ? new String(headers.lastHeader(headerName).value()) 
                : null;
    }
}

方案二:使用ConsumerInterceptor(全局统一处理)

实现Spring Kafka的ConsumerInterceptor,在消息传递给监听方法前全局拦截处理,适合多个监听方法需要统一逻辑的场景。

1. 定义拦截器

import org.apache.kafka.clients.consumer.ConsumerInterceptor;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.common.TopicPartition;
import org.slf4j.MDC;
import java.util.Map;
import java.util.UUID;

public class MdcConsumerInterceptor implements ConsumerInterceptor<String, String> {

    @Override
    public ConsumerRecords<String, String> onConsume(ConsumerRecords<String, String> records) {
        // 遍历每条消息处理MDC(批量消费时需处理每条)
        records.forEach(record -> processHeadersToMDC(record.headers()));
        return records;
    }

    private void processHeadersToMDC(org.apache.kafka.common.header.Headers headers) {
        String correlationId = getHeaderValue(headers, "correlationId") == null 
                ? UUID.randomUUID().toString() 
                : getHeaderValue(headers, "correlationId");
        String requestId = getHeaderValue(headers, "requestId") == null 
                ? UUID.randomUUID().toString() 
                : getHeaderValue(headers, "requestId");
        
        MDC.put("correlationId", correlationId);
        MDC.put("requestId", requestId);
    }

    private String getHeaderValue(org.apache.kafka.common.header.Headers headers, String headerName) {
        return headers.lastHeader(headerName) != null 
                ? new String(headers.lastHeader(headerName).value()) 
                : null;
    }

    @Override
    public void onCommit(Map<TopicPartition, OffsetAndMetadata> offsets) {
        // 提交偏移量后清理MDC
        MDC.clear();
    }

    @Override
    public void close() {
        MDC.clear();
    }

    @Override
    public void configure(Map<String, ?> configs) {
        // 初始化配置,按需实现
    }
}

2. 配置拦截器到消费者工厂

import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
import java.util.HashMap;
import java.util.Map;
import static org.apache.kafka.clients.consumer.ConsumerConfig.*;

@Configuration
public class KafkaConsumerConfig {

    @Bean
    public ConsumerFactory<String, String> consumerFactory() {
        Map<String, Object> configProps = new HashMap<>();
        configProps.put(BOOTSTRAP_SERVERS_CONFIG, "kafka-server:9092");
        configProps.put(GROUP_ID_CONFIG, "your-group-id");
        configProps.put(KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
        configProps.put(VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
        // 注册拦截器
        configProps.put(INTERCEPTOR_CLASSES_CONFIG, "com.your.package.MdcConsumerInterceptor");
        return new DefaultKafkaConsumerFactory<>(configProps);
    }
}

方案三:使用Spring AOP拦截监听方法(灵活可控)

通过AOP拦截所有@KafkaListener标注的方法,在方法执行前从参数中提取消息Header处理MDC,适合已有代码无需修改方法参数的场景。

import org.aspectj.lang.JoinPoint;
import org.aspectj.lang.annotation.After;
import org.aspectj.lang.annotation.Aspect;
import org.aspectj.lang.annotation.Before;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.slf4j.MDC;
import org.springframework.stereotype.Component;
import java.util.UUID;

@Aspect
@Component
public class KafkaListenerMdcAspect {

    @Before("@annotation(org.springframework.kafka.annotation.KafkaListener)")
    public void beforeListenerExecution(JoinPoint joinPoint) {
        // 遍历方法参数,找到ConsumerRecord类型的参数
        for (Object arg : joinPoint.getArgs()) {
            if (arg instanceof ConsumerRecord) {
                ConsumerRecord<?, ?> record = (ConsumerRecord<?, ?>) arg;
                processHeadersToMDC(record.headers());
                break;
            }
        }
    }

    @After("@annotation(org.springframework.kafka.annotation.KafkaListener)")
    public void afterListenerExecution() {
        // 方法执行完毕后清理MDC
        MDC.clear();
    }

    private void processHeadersToMDC(org.apache.kafka.common.header.Headers headers) {
        String correlationId = getHeaderValue(headers, "correlationId") == null 
                ? UUID.randomUUID().toString() 
                : getHeaderValue(headers, "correlationId");
        String requestId = getHeaderValue(headers, "requestId") == null 
                ? UUID.randomUUID().toString() 
                : getHeaderValue(headers, "requestId");
        
        MDC.put("correlationId", correlationId);
        MDC.put("requestId", requestId);
    }

    private String getHeaderValue(org.apache.kafka.common.header.Headers headers, String headerName) {
        return headers.lastHeader(headerName) != null 
                ? new String(headers.lastHeader(headerName).value()) 
                : null;
    }
}

注意:无论使用哪种方案,都要在合适的时机清理MDC(比如方法结束后、偏移量提交后),避免线程池复用导致的MDC数据污染。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 20:50:32