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

如何在Java Spring中记录包含分区与偏移量的原始Kafka消息?

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.15 09:30:29