必须在消费者配置层面启用errorhandlingdeserializer2才能跳过反序列化异常继续消费,核心是配置value-deserializer为errorhandlingdeserializer2.class、指定delegate class和兜底函数,并为key同样配置以避免中断。

Java中Kafka Consumer消费时抛反序列化异常(比如SerializationException或DeserializationException),不能靠业务层try-catch捕获,因为异常发生在消息拉取后、进入监听器前的反序列化阶段。必须在消费者配置层面启用ErrorHandlingDeserializer,才能让消费者跳过毒丸、继续消费后续消息。
核心配置:用ErrorHandlingDeserializer2替代原Deserializer
Spring Kafka提供了ErrorHandlingDeserializer2(推荐)或ErrorHandlingDeserializer(旧版),它本身不真正反序列化,而是包装一个“后备”反序列化器,并在失败时调用你指定的函数生成默认对象或兜底值。
- 把
value-deserializer设为ErrorHandlingDeserializer2.class - 通过
spring.deserializer.value.delegate.class(或配置项ErrorHandlingDeserializer2.VALUE_DESERIALIZER_CLASS)指定原始反序列化器,如JsonDeserializer.class - 通过
ErrorHandlingDeserializer2.VALUE_FUNCTION指定失败时的兜底逻辑(BiFunction<byte headers t></byte>)
application.yml示例(Spring Boot)
比硬编码更清晰、易维护:
spring:
kafka:
consumer:
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer2
properties:
spring.deserializer.key.delegate.class: org.apache.kafka.common.serialization.StringDeserializer
spring.deserializer.value.delegate.class: org.springframework.kafka.support.serializer.JsonDeserializer
spring.json.value.default.type: com.example.NTCMessageBody
# 失败时返回自定义兜底对象
spring.deserializer.value.function: com.example.FailedNTCMessageBodyProvider
实现兜底函数:FailedNTCMessageBodyProvider
这个类决定反序列化失败后返回什么对象,通常用于标记异常消息、保留原始字节以便排查:
- 实现
BiFunction<byte headers t></byte>接口 - 参数
byte[]是原始消息体,Headers含元数据(如topic、partition、offset) - 返回一个可识别的“坏消息”对象,比如继承自正常类型、带
failedDecode字段
例如:
public class FailedNTCMessageBodyProvider implements BiFunction<byte headers ntcmessagebody> {
@Override
public NTCMessageBody apply(byte[] data, Headers headers) {
return new NTCBadMessageBody(data); // 自定义兜底对象
}
}
</byte>
注意键(key)也要配,否则key反序列化失败同样卡住
如果消息key也用了自定义序列化(比如LongSerializer),而消费者配置了StringDeserializer,key反序列化也会失败并中断消费。所以建议:
- 对key同样启用
ErrorHandlingDeserializer2 - 设置
spring.deserializer.key.delegate.class指向真实key反序列化器 - 若key固定为字符串,可直接用
StringDeserializer,无需兜底
补充:配合DefaultErrorHandler做业务异常兜底
ErrorHandlingDeserializer只解决反序列化阶段失败;业务逻辑中抛出的异常(如空指针、数据库错误)需由DefaultErrorHandler处理:
- 配置重试次数、退避策略
- 设置死信主题(DLQ)转发不可恢复消息
- 显式声明
DeserializationException为不可重试(它默认已是fatal)
这样,反序列化失败走兜底对象,业务异常走重试/DLQ,职责分离,系统更健壮。
Java免费学习笔记:立即使用
解锁 Java 大师之旅:从入门到精通的终极指南











