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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 04:22:02