消费Kafka Topic消息遇DeserializationException,如何避免消费者阻塞?
Kafka消费者因消息结构变更反序列化失败阻塞的解决方案
问题现象与复现
- 消费者启动读取Topic时,遇到结构变更的消息会触发反序列化失败,导致消费者阻塞,无法执行
listenOnMessage方法内的业务逻辑 - 复现步骤:
- 生产初始结构的消息(
Message含String name字段) - 修改消息结构为
List<String> names后再生产消息 - 消费者读取到新结构消息时触发反序列化异常,进程停滞
- 生产初始结构的消息(
相关代码
初始消息结构:
// 初始消息结构 class Message { String name; }
变更后的消息结构:
// 修改后的消息结构 class Message { List<String> names; }
消费者消费方法:
// 反序列化异常时,该方法不会执行 @KafkaHandler public void listenOnMessage(Message message) { log.info("message: {}", message); // 业务逻辑... }
错误信息
反序列化时抛出类型不匹配异常(如Jackson的MismatchedInputException),提示无法将字符串类型转换为List<String>类型,导致消费进程阻塞。
问题根源
生产者在某一时间点切换了消息结构,消费者的反序列化逻辑未做兼容,无法处理新旧两种消息格式,触发异常后未正确处理,导致消费者阻塞。
解决方案
1. 配置错误处理与重试机制
通过Spring Kafka的错误处理器捕获反序列化异常,避免消费者阻塞,同时添加重试逻辑,最终将无法处理的消息转发到死信队列:
@Bean public SeekToCurrentErrorHandler kafkaErrorHandler(KafkaTemplate<String, Object> kafkaTemplate) { // 配置死信队列转发规则 DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(kafkaTemplate, (record, ex) -> new TopicPartition(record.topic() + ".dlq", record.partition())); // 重试3次,每次间隔1秒,失败后转发到死信队列 return new SeekToCurrentErrorHandler(recoverer, new FixedBackOff(1000L, 3L)); }
2. 兼容新旧消息结构
修改Message类,通过注解兼容新旧字段,确保反序列化器能同时处理两种格式:
import com.fasterxml.jackson.annotation.JsonAlias; import java.util.Collections; import java.util.List; class Message { @JsonAlias("name") // 兼容旧字段"name" private List<String> names; // 可选:自动将旧字段值转换为新结构 public void setName(String name) { this.names = Collections.singletonList(name); } // getter、setter方法 }
3. 自定义反序列化器
实现自定义反序列化器,优先尝试新结构反序列化,失败则降级处理旧结构并转换为新格式:
import com.fasterxml.jackson.databind.ObjectMapper; import org.apache.kafka.common.serialization.Deserializer; import java.io.IOException; import java.util.Collections; public class MessageDeserializer implements Deserializer<Message> { private final ObjectMapper objectMapper = new ObjectMapper(); @Override public Message deserialize(String topic, byte[] data) { try { // 优先按新结构反序列化 return objectMapper.readValue(data, Message.class); } catch (Exception e) { try { // 失败则按旧结构解析,再转换为新结构 OldMessage oldMsg = objectMapper.readValue(data, OldMessage.class); Message newMsg = new Message(); newMsg.setNames(Collections.singletonList(oldMsg.getName())); return newMsg; } catch (IOException ex) { throw new RuntimeException("消息反序列化失败", ex); } } } // 临时旧结构类 private static class OldMessage { private String name; public String getName() { return name; } public void setName(String name) { this.name = name; } } }
4. 直接跳过异常消息(谨慎使用)
如果不需要保留异常消息,可以直接跳过当前异常消息,继续消费后续内容:
@Bean public ErrorHandler skipErrorHandler() { return (record, exception) -> { log.error("跳过异常消息,offset: {}", record.offset(), exception); // 手动移动到下一个offset Consumer<?, ?> consumer = ((ConsumerRecord<?, ?>) record).consumer(); consumer.seek(new TopicPartition(record.topic(), record.partition()), record.offset() + 1); }; }
内容的提问来源于stack exchange,提问作者Sweet Pipe
相关产品推荐
相关产品推荐

