如何在@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
相关产品推荐
相关产品推荐

