Kafka反序列化报错:无类型信息头且未提供默认类型
Kafka自定义Payload反序列化问题排查与解决
报错信息
Error: No type information in headers and no default type provided. Error deserializing key/value for partition slack_command_event-0 at offset 0. If needed, please seek past the record to continue consumption.
问题背景
作为Kafka新手,尝试向消费者发送自定义SlashCommandPayload对象(含id、command、query、responseUrl四个字段)时触发上述反序列化错误,以下是相关配置与代码片段:
相关代码片段
生产者配置
@Configuration @RequiredArgsConstructor public class SlackCommandProducerConfig { private final KafkaProperties kafkaProperties; @Bean("slashCommandProducerFactory") public ProducerFactory<String, SlashCommandPayload> slackCommandProducerFactory() { var props = new HashMap<>(kafkaProperties.buildProducerProperties()); return new DefaultKafkaProducerFactory<>(props,() -> null,() -> new SlashCommandSerializer<>(SlashCommandPayload .class)); } @Bean("slackCommandKafkaTemplate") public KafkaTemplate<String, SlashCommandPayload> slackCommandKafkaTemplate(@Qualifier("slashCommandProducerFactory") ProducerFactory<String, SlashCommandPayload> slackCmdProducerFactory) { return new KafkaTemplate<>(slackCmdProducerFactory); } }
消费者配置
@Configuration @RequiredArgsConstructor public class SlackCommandConsumerConfig { private final KafkaProperties kafkaProperties; @Bean @ConditionalOnMissingBean(name = "slackCmdConsumerConfig") public ConsumerFactory<String, SlashCommandPayload> slackCmdConsumerConfig() { var kafkaProps = kafkaProperties.buildConsumerProperties(); //kafkaProps.put(JsonSerializer.ADD_TYPE_INFO_HEADERS, false); var props = new HashMap<>(kafkaProps); return new DefaultKafkaConsumerFactory<>(props,() -> null, () -> new JsonDeserializer<>(SlashCommandPayload.class)); } }
自定义序列化器
public class SlashCommandSerializer<T> implements Serializer<T> { private final Gson gson = GsonFactory.createSnakeCase(SlackConfig.DEFAULT); @Override public byte[] serialize(String data, T o) { if (data == null || o == null) { return null; } try { return gson.toJson(o).getBytes(StandardCharsets.UTF_8); } catch (Exception e) { throw new SerializationException("Error serializing JSON message", e); } } @Override public void configure(Map<String, ?> configs, boolean isKey) { Serializer.super.configure(configs, isKey); } @Override public void close() { } }
自定义反序列化器
public class SlashCommandDeserializer<T> implements Deserializer<T> { private final Gson gson = GsonFactory.createSnakeCase(SlackConfig.DEFAULT); private Class<T> destinationClass; public SlashCommandDeserializer(Class<T> destinationClass) { this.destinationClass = destinationClass; } @Override public void configure(Map<String, ?> configs, boolean isKey) { } @Override public void close() { } @Override public SlashCommandPayload deserialize(String s, byte[] bytes) { if (bytes == null) { return null; } try { return gson.fromJson(new String(bytes, StandardCharsets.UTF_8), destinationClass); } catch (Exception e) { throw new SerializationException("Error deserializing message", e); } } }
Payload类
@Data @AllArgsConstructor @NoArgsConstructor public class SlashCommandPayload { private String id; private String command; private String query; private String responseUrl; }
消息发送代码
private final KafkaTemplate<String, SlashCommandPayload> kafkaSlashCommandTemplate; public SlackEvent(@Qualifier("slackCommandKafkaTemplate") KafkaTemplate<String, SlashCommandPayload> kafkaSlashCommandTemplate){ this.kafkaSlashCommandTemplate = kafkaSlashCommandTemplate; } public void handleSlackSlashCommand(SlashCommandRequest slashCommandRequest) { var payload = new SlashCommandPayload(); payload.setId(slashCommandRequest.getPayload().getTeamId()); payload.setCommand(slashCommandRequest.getPayload().getCommand()); payload.setQuery(slashCommandRequest.getPayload().getText()); payload.setResponseUrl(slashCommandRequest.getResponseUrl()); kafkaSlashCommandTemplate.send("slack_command_event", slashCommandRequest.getContext().getTeamId(), payload); }
消息接收代码
@KafkaListener(id = "slack_command_event_consumer", topics = "slack_command_event", concurrency = "1") public void processSlackCommandEvent(SlashCommandPayload event) { var kafkaPayload = new SlackCommandEventPayload(); kafkaPayload.setId(event.getId()); kafkaPayload.setQuery(event.getQuery()); kafkaPayload.setResponseUrl(event.getResponseUrl()); handler.processCommandEvent(kafkaPayload); }
解决思路与步骤
1. 修复自定义反序列化器的类型错误
自定义反序列化器的deserialize方法硬编码返回SlashCommandPayload,与泛型设计冲突,会导致类型转换异常,修改为返回泛型类型:
@Override public T deserialize(String s, byte[] bytes) { if (bytes == null) { return null; } try { return gson.fromJson(new String(bytes, StandardCharsets.UTF_8), destinationClass); } catch (Exception e) { throw new SerializationException("Error deserializing message", e); } }
2. 统一序列化/反序列化方案(二选一)
方案一:使用自定义序列化器+反序列化器
当前消费者端错误使用了Spring自带的JsonDeserializer,需替换为自定义的SlashCommandDeserializer:
@Bean @ConditionalOnMissingBean(name = "slackCmdConsumerConfig") public ConsumerFactory<String, SlashCommandPayload> slackCmdConsumerConfig() { var kafkaProps = kafkaProperties.buildConsumerProperties(); var props = new HashMap<>(kafkaProps); return new DefaultKafkaConsumerFactory<>(props,() -> null, () -> new SlashCommandDeserializer<>(SlashCommandPayload.class)); }
方案二:改用Spring Kafka自带的JSON序列化组件(推荐)
放弃自定义序列化器,直接使用Spring提供的JsonSerializer和JsonDeserializer,自动处理类型头信息:
- 生产者配置修改:
@Bean("slashCommandProducerFactory") public ProducerFactory<String, SlashCommandPayload> slackCommandProducerFactory() { var props = new HashMap<>(kafkaProperties.buildProducerProperties()); props.put(JsonSerializer.ADD_TYPE_INFO_HEADERS, true); // 显式启用类型头(默认已开启) return new DefaultKafkaProducerFactory<>(props, new StringSerializer(), new JsonSerializer<>(SlashCommandPayload.class)); }
- 消费者配置修改:
@Bean @ConditionalOnMissingBean(name = "slackCmdConsumerConfig") public ConsumerFactory<String, SlashCommandPayload> slackCmdConsumerConfig() { var kafkaProps = kafkaProperties.buildConsumerProperties(); // 指定默认反序列化类型,确保无类型头时也能正常解析 kafkaProps.put(JsonDeserializer.VALUE_DEFAULT_TYPE, SlashCommandPayload.class.getName()); return new DefaultKafkaConsumerFactory<>(props, new StringDeserializer(), new JsonDeserializer<>(SlashCommandPayload.class)); }
3. 清理无效消息
之前发送的消息因序列化格式不匹配或缺少类型头,消费者无法解析,需手动跳过:
- 使用Kafka命令行工具重置消费者偏移量:
kafka-consumer-groups.sh --bootstrap-server <你的Kafka地址> --group slack_command_event_consumer --topic slack_command_event --reset-offsets --to-latest --execute
- 或在代码中配置错误处理器自动跳过(测试环境可用):
添加错误处理器Bean:
@Bean public SeekToCurrentErrorHandler seekToCurrentErrorHandler() { return new SeekToCurrentErrorHandler(); }
修改@KafkaListener注解:
@KafkaListener(id = "slack_command_event_consumer", topics = "slack_command_event", concurrency = "1", errorHandler = "seekToCurrentErrorHandler")
内容的提问来源于stack exchange,提问作者James K J
相关产品推荐
相关产品推荐

