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

Spring Kafka反序列化异常求助:配置ErrorHandlingDeserializer仍未解决

问题:Kafka消费者遇非预期消息时重复输出序列化错误,配置ErrorHandlingDeserializer无效

发送非预期类型的消息到Kafka主题后,消费者会持续重复输出相同的错误日志,报错堆栈信息如下:

java.lang.IllegalStateException: This error handler cannot process 'SerializationException's directly; please consider configuring an 'ErrorHandlingDeserializer' in the value and/or key deserializer
    at org.springframework.kafka.listener.DefaultErrorHandler.handleOtherException(DefaultErrorHandler.java:192) ~[spring-kafka-3.1.1.jar:3.1.1]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.handleConsumerException(KafkaMessageListenerContainer.java:1934) ~[spring-kafka-3.1.1.jar:3.1.1]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1365) ~[spring-kafka-3.1.1.jar:3.1.1]
    at java.base/java.util.concurrent.CompletableFuture$AsyncRun.run(CompletableFuture.java:1804) ~[na:na]
    at java.base/java.lang.Thread.run(Thread.java:833) ~[na:na]
Caused by: org.apache.kafka.common.errors.RecordDeserializationException: Error deserializing key/value for partition like-0 at offset 1. If needed, please seek past the record to continue consumption.
    at org.apache.kafka.clients.consumer.internals.CompletedFetch.parseRecord(CompletedFetch.java:309) ~[kafka-clients-3.6.1.jar:na]
    at org.apache.kafka.clients.consumer.internals.CompletedFetch.fetchRecords(CompletedFetch.java:263) ~[kafka-clients-3.6.1.jar:na]
    at org.apache.kafka.clients.consumer.internals.AbstractFetch.fetchRecords(AbstractFetch.java:340) ~[kafka-clients-3.6.1.jar:na]
    at org.apache.kafka.clients.consumer.internals.AbstractFetch.collectFetch(AbstractFetch.java:306) ~[kafka-clients-3.6.1.jar:na]
    at org.apache.kafka.clients.consumer.KafkaConsumer.pollForFetches(KafkaConsumer.java:1235) ~[kafka-clients-3.6.1.jar:na]
    at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1186) ~[kafka-clients-3.6.1.jar:na]
    at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1159) ~[kafka-clients-3.6.1.jar:na]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollConsumer(KafkaMessageListenerContainer.java:1649) ~[spring-kafka-3.1.1.jar:3.1.1]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doPoll(KafkaMessageListenerContainer.java:1624) ~[spring-kafka-3.1.1.jar:3.1.1]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:1421) ~[spring-kafka-3.1.1.jar:3.1.1]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1313) ~[spring-kafka-3.1.1.jar:3.1.1]
    ... 2 common frames omitted
Caused by: org.apache.kafka.common.errors.SerializationException: Can't deserialize data  from topic [like]
    at org.springframework.kafka.support.serializer.JsonDeserializer.deserialize(JsonDeserializer.java:588) ~[spring-kafka-3.1.1.jar:3.1.1]
    at org.apache.kafka.common.serialization.Deserializer.deserialize(Deserializer.java:73) ~[kafka-clients-3.6.1.jar:na]
    at org.apache.kafka.clients.consumer.internals.CompletedFetch.parseRecord(CompletedFetch.java:300) ~[kafka-clients-3.6.1.jar:na]
    ... 12 common frames omitted
Caused by: com.fasterxml.jackson.core.JsonParseException: Unrecognized token 'asd': was expecting (JSON String, Number, Array, Object or token 'null', 'true' or 'false')
 at [Source: (byte[])"asd"; line: 1, column: 4]
    at com.fasterxml.jackson.core.JsonParser._constructError(JsonParser.java:2477) ~[jackson-core-2.15.3.jar:2.15.3]
    at com.fasterxml.jackson.core.base.ParserMinimalBase._reportError(ParserMinimalBase.java:760) ~[jackson-core-2.15.3.jar:2.15.3]
    at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._reportInvalidToken(UTF8StreamJsonParser.java:3699) ~[jackson-core-2.15.3.jar:2.15.3]
    at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._handleUnexpectedValue(UTF8StreamJsonParser.java:2787) ~[jackson-core-2.15.3.jar:2.15.3]
    at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._nextTokenNotInObject(UTF8StreamJsonParser.java:908) ~[jackson-core-2.15.3.jar:2.15.3]
    at com.fasterxml.jackson.core.json.UTF8StreamJsonParser.nextToken(UTF8StreamJsonParser.java:794) ~[jackson-core-2.15.3.jar:2.15.3]
    at com.fasterxml.jackson.databind.ObjectReader._initForReading(ObjectReader.java:357) ~[jackson-databind-2.15.3.jar:2.15.3]
    at com.fasterxml.jackson.databind.ObjectReader._bindAndClose(ObjectReader.java:2095) ~[jackson-databind-2.15.3.jar:2.15.3]
    at com.fasterxml.jackson.databind.ObjectReader.readValue(ObjectReader.java:1583) ~[jackson-databind-2.15.3.jar:2.15.3]
    at org.springframework.kafka.support.serializer.JsonDeserializer.deserialize(JsonDeserializer.java:585) ~[spring-kafka-3.1.1.jar:3.1.1]
    ... 14 common frames omitted

已按照Spring官方文档建议配置ErrorHandlingDeserializer,但问题仍未解决,当前消费者配置代码如下:

@Bean
public Map<String, Object> consumerConfig() {
    Map<String, Object> props = new HashMap<>();
    props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP_SERVER);
    props.put(ConsumerConfig.GROUP_ID_CONFIG, "group-01");

    props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
    props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);

    props.put(ErrorHandlingDeserializer.KEY_DESERIALIZER_CLASS, StringDeserializer.class);
    props.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS, JsonDeserializer.class.getName());

    props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "true");

    props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
    return props;
}

(也尝试过YAML配置方式,问题依旧)


排查方向建议

  • 检查配置是否生效:确认没有其他消费者配置(如多个ConsumerFactory实例、@KafkaListener注解指定的自定义反序列化器)覆盖当前设置,导致ErrorHandlingDeserializer未实际生效。
  • 补充JsonDeserializer必要配置:当前仅指定了JsonDeserializer类名,需添加目标反序列化类型配置,例如:
    // 指定默认反序列化的Java类
    props.put(JsonDeserializer.VALUE_DEFAULT_TYPE, "com.yourpackage.YourTargetDto");
    // 如果消息没有携带类型头,关闭类型头校验
    props.put(JsonDeserializer.USE_TYPE_INFO_HEADERS, false);
    
  • 调整偏移量提交策略:当前开启了自动提交,但反序列化失败时,消费未完成导致偏移量不提交,消费者会持续重试该消息。建议:
    1. 关闭自动提交,改用手动提交;
    2. 配置DefaultErrorHandler添加跳过或死信队列策略,例如:
      @Bean
      public DefaultErrorHandler errorHandler(KafkaTemplate<?, ?> kafkaTemplate) {
          // 失败后重试2次,然后发送到死信队列
          return new DefaultErrorHandler(
              new DeadLetterPublishingRecoverer(kafkaTemplate),
              new FixedBackOff(1000L, 2L)
          );
          // 或者直接跳过无法反序列化的消息
          // return new DefaultErrorHandler((record, ex) -> {}, new FixedBackOff(0L, 0L));
      }
      
  • 统一反序列化器配置格式:将ErrorHandlingDeserializer.KEY_DESERIALIZER_CLASS的值改为类名,与VALUE配置保持一致,避免Spring处理类对象时出现异常:
    props.put(ErrorHandlingDeserializer.KEY_DESERIALIZER_CLASS, StringDeserializer.class.getName());
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 16:14:53