Spring Kafka重试消息时FailedDeserializationInfo序列化失败求助
Spring Kafka重试消息时处理FailedDeserializationInfo序列化问题
问题原因
当使用ErrorHandlingDeserializer处理反序列化错误时,springDeserializerExceptionValue头中存储的FailedDeserializationInfo对象包含DeserializationExceptionHeader类实例,该类没有可被Jackson序列化的属性,导致DefaultKafkaHeaderMapper在映射Header时触发序列化错误。
解决方案
方案一:直接传递原始Header字节(推荐)
ErrorHandlingDeserializer本身是将FailedDeserializationInfo序列化为字节后存入Header的,转发时直接提取原始字节重新添加到新消息的Header中,绕过HeaderMapper的序列化逻辑:
// 获取原始消息的Header Headers originalHeaders = failedRecord.headers(); // 提取springDeserializerExceptionValue的原始字节 Header exceptionHeader = originalHeaders.lastHeader("springDeserializerExceptionValue"); // 构建新消息的Headers RecordHeaders newHeaders = new RecordHeaders(); // 复制其他原始Header originalHeaders.forEach(header -> { if (!"springDeserializerExceptionValue".equals(header.key())) { newHeaders.add(header); } }); // 直接添加原始字节形式的异常Header if (exceptionHeader != null) { newHeaders.add(exceptionHeader.key(), exceptionHeader.value()); } // 发送到重试主题 kafkaTemplate.send("retry-topic", failedRecord.key(), failedRecord.value(), newHeaders);
方案二:自定义HeaderMapper转换异常信息
继承DefaultKafkaHeaderMapper,将FailedDeserializationInfo转换为自定义的可序列化对象,保留关键错误信息:
public class CustomErrorHeaderMapper extends DefaultKafkaHeaderMapper { @Override public void mapHeaders(Headers headers, Map<String, Object> target) { super.mapHeaders(headers, target); Object errorInfo = target.get("springDeserializerExceptionValue"); if (errorInfo instanceof FailedDeserializationInfo) { FailedDeserializationInfo failedInfo = (FailedDeserializationInfo) errorInfo; // 封装为可序列化的自定义对象 SerializableErrorDetail detail = new SerializableErrorDetail( failedInfo.getException().getMessage(), failedInfo.getTopic(), Base64.getEncoder().encodeToString(failedInfo.getRawValue()) ); target.put("springDeserializerExceptionValue", detail); } } // 自定义可序列化的错误详情类 public static class SerializableErrorDetail { private String exceptionMsg; private String topic; private String rawDataBase64; // 构造器、Getter、Setter public SerializableErrorDetail(String exceptionMsg, String topic, String rawDataBase64) { this.exceptionMsg = exceptionMsg; this.topic = topic; this.rawDataBase64 = rawDataBase64; } public String getExceptionMsg() { return exceptionMsg; } public void setExceptionMsg(String exceptionMsg) { this.exceptionMsg = exceptionMsg; } public String getTopic() { return topic; } public void setTopic(String topic) { this.topic = topic; } public String getRawDataBase64() { return rawDataBase64; } public void setRawDataBase64(String rawDataBase64) { this.rawDataBase64 = rawDataBase64; } } }
然后将自定义HeaderMapper配置到KafkaTemplate:
@Bean public KafkaTemplate<String, Object> retryKafkaTemplate(ProducerFactory<String, Object> producerFactory) { KafkaTemplate<String, Object> template = new KafkaTemplate<>(producerFactory); template.setHeaderMapper(new CustomErrorHeaderMapper()); return template; }
方案三:配置Jackson忽略空对象序列化错误
通过修改ObjectMapper配置,让Jackson忽略DeserializationExceptionHeader这类无属性类的序列化错误,不推荐全局配置,建议针对特定类型处理:
@Bean public ObjectMapper kafkaObjectMapper() { ObjectMapper mapper = new ObjectMapper(); // 为DeserializationExceptionHeader添加忽略空属性的混合类 mapper.addMixIn(DeserializationExceptionHeader.class, IgnoreEmptyPropsMixin.class); return mapper; } // 混合类,用于忽略空属性 @JsonIgnoreProperties(ignoreUnknown = true) interface IgnoreEmptyPropsMixin {} // 配置HeaderMapper使用自定义ObjectMapper @Bean public DefaultKafkaHeaderMapper kafkaHeaderMapper(ObjectMapper kafkaObjectMapper) { DefaultKafkaHeaderMapper mapper = new DefaultKafkaHeaderMapper(); mapper.setObjectMapper(kafkaObjectMapper); return mapper; }
内容的提问来源于stack exchange,提问作者dwb5013
相关产品推荐
相关产品推荐

