如何配置Kafka自定义反序列化器通过@KafkaListener获取消费的JSON数据
你现有代码存在3个核心问题,按以下步骤修改即可正常运行:
1. 修复无效的自定义反序列化器
你写的反序列化器调用了SerializationUtils.deserialize,这是处理Java原生序列化字节流的方法,完全无法解析JSON格式消息,且没有处理泛型目标类型的读取,修正后代码如下:
public class CustomDeserializer<T extends Serializable> implements Deserializer<T> { private final ObjectMapper objectMapper = new ObjectMapper(); public static final String VALUE_CLASS_NAME_CONFIG = "value.class.name"; private Class<T> targetClass; @Override public void configure(Map<String, ?> configs, boolean isKey) { // 从配置中读取要反序列化的目标实体类型 targetClass = (Class<T>) configs.get(VALUE_CLASS_NAME_CONFIG); } @Override public T deserialize(String topic, byte[] objectData) { if (objectData == null) { return null; } try { return objectMapper.readValue(objectData, targetClass); } catch (IOException e) { throw new RuntimeException("JSON消息反序列化失败", e); } } @Override public void close() {} }
注意:提前给KafkaPayload、EventHeader、NewFields、OldFields四个实体类加无参构造、getter/setter,或者直接加Lombok的@Data注解,ObjectMapper才能正常反射赋值;如果JSON字段名和实体类属性名不匹配(比如JSON里是大驼峰EventHeader,实体属性是小驼峰eventHeader),要在属性上加@JsonProperty("EventHeader")做映射。
2. 修正Kafka配置类
你原来的consumerFactory方法返回类型错误:Spring Kafka的监听容器需要的是ConsumerFactory(消费者实例工厂),而非直接new出来的原生KafkaConsumer实例——原生KafkaConsumer是非线程安全的,容器需要根据并发配置自行创建多个消费者实例,不能直接注入单例的消费者对象。修正后配置类如下:
@Configuration // 启动类记得加@EnableKafka注解,否则@KafkaListener不会生效 public class KafkaConfig { @Bean public ConsumerFactory<String, KafkaPayload> consumerFactory(){ Map<String, Object> config = new HashMap<>(); config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, CustomDeserializer.class); // 给自定义反序列化器传入目标类型 config.put(CustomDeserializer.VALUE_CLASS_NAME_CONFIG, KafkaPayload.class); return new DefaultKafkaConsumerFactory<>(config); } @Bean public ConcurrentKafkaListenerContainerFactory<String, KafkaPayload> kafkaListener(){ ConcurrentKafkaListenerContainerFactory<String, KafkaPayload> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); return factory; } }
提效提示:如果没有特殊自定义需求,完全可以不用自己写反序列化器,直接用Spring Kafka自带的JsonDeserializer,只需要把反序列化类配置替换为JsonDeserializer.class,再加一行配置config.put(JsonDeserializer.VALUE_DEFAULT_TYPE, KafkaPayload.class)即可,稳定性更高。
3. 修正消费监听类
给消费类加@Service注解让Spring上下文托管,否则监听器不会被扫描加载:
@Service public class ConsumerService { @KafkaListener(topics = "Topic", groupId = "groupId", containerFactory = "kafkaListener") public void consume(KafkaPayload kafkaPayload){ System.out.println("Consumed Message :" + kafkaPayload); } }
启动前检查项
- 确认Kafka服务端正常运行,监听的
Topic已提前创建 - 确认生产端发送的消息是标准JSON格式,无多余转义、不可见字符
- 确认Spring Boot启动类上添加了
@EnableKafka注解
内容的提问来源于stack exchange,提问作者NNFdo

