如何在Java Spring中记录包含分区与偏移量的原始Kafka消息?
我太懂你这种需求了——想把Kafka消息的原始内容和它的核心元数据(主题、分区、偏移量)打包在同一条日志里,方便后续排查问题,但试过的两种路子都走不通:用RecordInterceptor拿不到反序列化前的原始消息,自定义Deserializer又摸不到分区和偏移量。别愁,给你两个落地性很强的解决方案:
方案一:直接在@KafkaListener中接收ConsumerRecord
这是最直接的方式,跳过中间的自动反序列化,直接拿到最原始的消息载体,同时一次性获取所有元数据。
你只需要把@KafkaListener方法的参数改成ConsumerRecord<byte[], byte[]>,这样就能拿到原始的字节数组(也就是你说的原始JSON消息体),然后从ConsumerRecord对象里轻松提取主题、分区、偏移量这些信息,最后统一记录日志,再自己处理反序列化逻辑即可。
示例代码:
import org.apache.kafka.clients.consumer.ConsumerRecord; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Component; import com.fasterxml.jackson.databind.ObjectMapper; import java.nio.charset.StandardCharsets; @Component public class KafkaMessageConsumer { private static final Logger log = LoggerFactory.getLogger(KafkaMessageConsumer.class); private final ObjectMapper objectMapper = new ObjectMapper(); @KafkaListener(topics = "${kafka.target.topic}") public void consumeRawMessage(ConsumerRecord<byte[], byte[]> record) { // 提取Kafka元数据 String topic = record.topic(); int partition = record.partition(); long offset = record.offset(); // 将原始字节数组转为UTF-8编码的JSON字符串 String rawJsonMessage = new String(record.value(), StandardCharsets.UTF_8); // 记录包含所有信息的单条日志 log.info("Received Kafka message - Topic: {}, Partition: {}, Offset: {}, Raw Message: {}", topic, partition, offset, rawJsonMessage); // 这里可以继续将原始JSON反序列化为你的业务对象 try { YourBusinessObject businessObj = objectMapper.readValue(rawJsonMessage, YourBusinessObject.class); // 执行你的业务逻辑 } catch (Exception e) { log.error("Failed to deserialize raw Kafka message", e); } } }
这个方案的好处是逻辑直观,完全由你掌控消息的处理流程,日志内容想怎么组合就怎么组合,没有额外的配置负担。
方案二:用RecordInterceptor配合ByteArrayDeserializer实现解耦
如果你的项目中有多个Kafka消费者,都需要记录原始消息和元数据,那可以把日志逻辑抽成拦截器,实现业务逻辑和日志逻辑的解耦。
核心思路是先把消费者的value反序列化器配置为ByteArrayDeserializer,让拦截器能拿到原始字节数组,然后在拦截器里完成日志记录,再把消息传递给后续的业务处理逻辑。
步骤1:配置消费者工厂,指定原始字节反序列化
import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.common.serialization.ByteArrayDeserializer; import org.apache.kafka.common.serialization.StringDeserializer; 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; @Configuration public class KafkaConsumerConfig { @Bean public ConsumerFactory<String, byte[]> rawConsumerFactory() { Map<String, Object> configProps = new HashMap<>(); configProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "${kafka.bootstrap.servers}"); configProps.put(ConsumerConfig.GROUP_ID_CONFIG, "${kafka.consumer.group.id}"); // 配置key为字符串反序列化,value为原始字节反序列化 configProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); configProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class); return new DefaultKafkaConsumerFactory<>(configProps); } }
步骤2:自定义日志拦截器
import org.apache.kafka.clients.consumer.ConsumerRecord; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.kafka.listener.RecordInterceptor; import org.springframework.stereotype.Component; import java.nio.charset.StandardCharsets; @Component public class RawKafkaMessageLoggingInterceptor implements RecordInterceptor<String, byte[]> { private static final Logger log = LoggerFactory.getLogger(RawKafkaMessageLoggingInterceptor.class); @Override public ConsumerRecord<String, byte[]> intercept(ConsumerRecord<String, byte[]> record) { // 提取元数据 String topic = record.topic(); int partition = record.partition(); long offset = record.offset(); // 转换原始字节为JSON字符串 String rawMessage = new String(record.value(), StandardCharsets.UTF_8); // 统一记录日志 log.info("Kafka Message Log - Topic: {}, Partition: {}, Offset: {}, Raw Content: {}", topic, partition, offset, rawMessage); // 将消息传递给后续的消费者方法 return record; } }
步骤3:把拦截器配置到监听容器工厂
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.core.ConsumerFactory; @Configuration public class KafkaListenerConfig { @Bean public ConcurrentKafkaListenerContainerFactory<String, byte[]> rawKafkaListenerContainerFactory( ConsumerFactory<String, byte[]> rawConsumerFactory, RawKafkaMessageLoggingInterceptor loggingInterceptor) { ConcurrentKafkaListenerContainerFactory<String, byte[]> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(rawConsumerFactory); // 给容器添加自定义的日志拦截器 factory.setRecordInterceptor(loggingInterceptor); return factory; } }
步骤4:在消费者方法中指定使用这个容器工厂
@KafkaListener(topics = "${kafka.target.topic}", containerFactory = "rawKafkaListenerContainerFactory") public void consumeMessage(ConsumerRecord<byte[], byte[]> record) { // 这里处理你的业务逻辑,比如反序列化原始消息为业务对象 ObjectMapper objectMapper = new ObjectMapper(); try { YourBusinessObject businessObj = objectMapper.readValue(record.value(), YourBusinessObject.class); // 执行业务操作 } catch (Exception e) { log.error("Business logic processing failed", e); } }
这个方案的优势是日志逻辑复用性高,多个消费者可以共用同一个拦截器,不用在每个方法里重复写日志代码。
补充说明
你之前用RecordInterceptor没成功,大概率是因为默认配置了业务对象的反序列化器,导致拦截器拿到的是已经反序列化后的对象,而非原始字节。只要把value的反序列化器换成ByteArrayDeserializer,就能拿到原始消息了。
备注:内容来源于stack exchange,提问作者KC Baltz

