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

Spring Boot Kafka消费者无法将JSON消息转为自定义Message类报错

Kafka JSON消息转换异常排查与解决

异常原因分析

  • 生产者泛型配置不匹配:ProducerFactory定义为<String, String>,但实际要发送Message对象且配置了JsonSerializer,泛型类型不匹配会导致序列化逻辑错误,最终发送的payload被错误处理为字符串类型,而非正确的JSON字节流。
  • KafkaTemplate冗余配置:调用kafkaTemplate.setConsumerFactory(messageConsumerFactory())完全多余,KafkaTemplate是生产者组件,无需关联消费者工厂,该配置可能干扰消息序列化/反序列化流程。
  • 类型头解析异常:消息头中的__TypeId__为字节数组形式,与消费者JsonDeserializer的类型解析逻辑不匹配,导致无法正确识别目标转换类型。

解决方案

1. 修正生产者泛型与序列化配置

将ProducerFactory泛型改为<String, Message>,匹配实际发送的对象类型,确保JsonSerializer正确序列化Message:

@Bean
public ProducerFactory<String, Message> producerFactory() {
    Map<String, Object> configProps = new HashMap<>();
    configProps.put(
            ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,
            bootstrapAddress);
    configProps.put(
            ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
            StringSerializer.class);
    configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
    // 配置类型映射,让消费者能识别消息类型
    configProps.put(JsonSerializer.TYPE_MAPPINGS, "message:com.xxx.xxx.kafka.messages.Message");
    return new DefaultKafkaProducerFactory<>(configProps);
}

2. 清理KafkaTemplate冗余配置

移除不必要的setConsumerFactory调用,KafkaTemplate仅需关联生产者工厂:

@Bean
public KafkaTemplate<String, Message> kafkaTemplate() {
    return new KafkaTemplate<>(producerFactory());
}

3. 优化消费者反序列化配置

为JsonDeserializer添加类型映射与信任包配置,确保正确解析消息:

public ConsumerFactory<String, Message> messageConsumerFactory() {
    Map<String, Object> props = new HashMap<>();
    props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress);
    props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
    props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
    props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class);
    props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
    // 与生产者配置一致的类型映射
    props.put(JsonDeserializer.TYPE_MAPPINGS, "message:com.xxx.xxx.kafka.messages.Message");
    // 信任目标类所在包,避免类型解析权限问题
    props.put(JsonDeserializer.TRUSTED_PACKAGES, "com.xxx.xxx.kafka.messages");
    
    return new DefaultKafkaConsumerFactory<>(
            props,
            new StringDeserializer(),
            new JsonDeserializer<>(Message.class));
}

4. 配置容器工厂的消息转换器(可选)

为消费者容器工厂设置MappingJackson2MessageConverter,增强JSON转换兼容性:

@Bean
public ConcurrentKafkaListenerContainerFactory<String, Message> messageKafkaListenerContainerFactory() {
    ConcurrentKafkaListenerContainerFactory<String, Message> factory =
            new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(messageConsumerFactory());
    
    MappingJackson2MessageConverter converter = new MappingJackson2MessageConverter();
    converter.setObjectMapper(new ObjectMapper());
    converter.setTypePrecedence(Jackson2JavaTypeMapper.TypePrecedence.TYPE_ID);
    converter.addTypeMapping("message", Message.class);
    factory.setMessageConverter(converter);
    
    return factory;
}

5. 确认发送逻辑

确保发送消息时直接传入Message对象,而非手动转为JSON字符串:

// 正确发送方式
kafkaTemplate.send("指定主题", message实例);

内容的提问来源于stack exchange,提问作者Hasan Isleem

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 09:00:39