Java21+Spring Boot3.2下Kafka收发Map<Short,List<Customer>>失败求助
问题排查与解决
核心原因
生产者与消费者的序列化/反序列化器配置不匹配,导致消费者接收的消息为字符串格式,无法直接转换为Map<Short, List<Customer>>类型。
具体修复步骤
1. 统一序列化器配置
生产者和消费者必须使用相同的序列化/反序列化器,推荐使用Spring Kafka提供的JsonSerializer和JsonDeserializer:
生产者配置(KafkaProducerConfig)
@Bean public ProducerFactory<String, Map<Short, List<Customer>>> producerFactory() { Map<String, Object> configProps = new HashMap<>(); configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "你的Kafka地址:9092"); configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); // 配置JSON序列化器处理value configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class); // 映射类型别名,帮助消费者识别目标类型 configProps.put(JsonSerializer.TYPE_MAPPINGS, "customerMap:java.util.Map<java.lang.Short, java.util.List<你的包路径.Customer>>"); return new DefaultKafkaProducerFactory<>(configProps); } @Bean public KafkaTemplate<String, Map<Short, List<Customer>>> kafkaTemplate() { return new KafkaTemplate<>(producerFactory()); }
消费者配置(KafkaConsumerConfig)
@Bean public ConsumerFactory<String, Map<Short, List<Customer>>> consumerFactory() { Map<String, Object> configProps = new HashMap<>(); configProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "你的Kafka地址:9092"); configProps.put(ConsumerConfig.GROUP_ID_CONFIG, "你的消费者组ID"); configProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); // 配置JSON反序列化器处理value configProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class); // 指定默认反序列化类型 configProps.put(JsonDeserializer.VALUE_DEFAULT_TYPE, "java.util.Map<java.lang.Short, java.util.List<你的包路径.Customer>>"); // 与生产者保持一致的类型映射 configProps.put(JsonDeserializer.TYPE_MAPPINGS, "customerMap:java.util.Map<java.lang.Short, java.util.List<你的包路径.Customer>>"); // 指定信任的包(生产环境避免使用*) configProps.put(JsonDeserializer.TRUSTED_PACKAGES, "你的包路径,java.util,java.lang"); return new DefaultKafkaConsumerFactory<>(configProps); } @Bean public ConcurrentKafkaListenerContainerFactory<String, Map<Short, List<Customer>>> kafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, Map<Short, List<Customer>>> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); return factory; }
2. 检查Customer实体类
确保Customer满足JSON反序列化要求:
- 有无参构造函数(Jackson必填)
- 字段具备完整的getter/setter方法
- 若使用Lombok,可添加
@Data+@NoArgsConstructor注解
@Data @NoArgsConstructor public class Customer implements Serializable { private Long id; private String name; // 其他业务字段 }
3. 消费者监听器参数匹配
确保@KafkaListener方法的参数类型与配置的目标类型完全一致:
@KafkaListener(topics = "你的主题名", groupId = "你的消费者组ID") public void handleMessage(Map<Short, List<Customer>> message) { // 业务处理逻辑 System.out.println("接收消息: " + message); }
4. 验证生产者发送逻辑
确保KafkaTemplate的泛型与发送数据类型匹配:
@Service public class SendToDatabase { private final KafkaTemplate<String, Map<Short, List<Customer>>> kafkaTemplate; public SendToDatabase(KafkaTemplate<String, Map<Short, List<Customer>>> kafkaTemplate) { this.kafkaTemplate = kafkaTemplate; } public void sendData(Map<Short, List<Customer>> data) { kafkaTemplate.send("你的主题名", data); } }
常见误区排查
- 禁止混用String序列化器与JSON序列化器:若生产者用
StringSerializer,消费者只能收到字符串,无法直接转Map - 类型映射必须一致:生产者与消费者的
TYPE_MAPPINGS配置要完全匹配,否则Jackson无法识别目标类型 - 信任包配置:生产环境不要设置
TRUSTED_PACKAGES: "*",仅指定业务包和必要的JDK基础包
内容的提问来源于stack exchange,提问作者Al3x4ndru1
相关产品推荐
相关产品推荐

