Spring-Kafka 3.0.4中DeserializationException的DLT配置处理异常问题
Spring Kafka 3.0.4 + AWS Glue 反序列化异常下的DLT配置问题
问题描述
使用Spring Kafka 3.0.4和AWS基础设施配置重试与死信队列(DLT)时,处理DeserializationException遇到异常:当消息反序列化失败时,DLT发送过程中错误使用了GlueSchemaRegistryKafkaSerializer,导致发送失败并阻塞消费者。
现有核心配置
@Bean public ConsumerFactory<String, Object> consumerFactory() { return new DefaultKafkaConsumerFactory<>( kafkaConsumerProps, StringDeserializer::new, () -> new ErrorHandlingDeserializer<>(kafkaAwsGlueValueDeserializer()), false ); } @Bean public DeadLetterPublishingRecoverer deadLetterPublishingRecoverer(KafkaTemplate<String, Object> kafkaTemplate) { DefaultKafkaProducerFactory<String, byte[]> defaultKafkaProducerFactory = new DefaultKafkaProducerFactory<>(kafkaProperties.buildProducerProperties(), new StringSerializer(), new ByteArraySerializer()); KafkaTemplate<String, byte[]> bytesKafkaTemplate = new KafkaTemplate<>(defaultKafkaProducerFactory); Map<Class<?>, KafkaOperations<?, ?>> templates = new LinkedHashMap<>(); templates.put(Object.class, kafkaTemplate); templates.put(byte[].class, bytesKafkaTemplate); DeadLetterPublishingRecoverer deadLetterPublishingRecoverer = new DeadLetterPublishingRecoverer(templates); return deadLetterPublishingRecoverer; }
错误日志
2023-08-31T09:06:47.952Z ERROR 1 --- [ntainer#6-0-C-1] k.r.DeadLetterPublishingRecovererFactory : Record: topic = test.TestData, partition = 1, offset = 0, main topic = test.TestData threw an error at topic test.TestData and won't be retried. Sending to DLT with name test.TestData-test-dlt. org.springframework.kafka.listener.ListenerExecutionFailedException: Listener failed ... Caused by: org.springframework.kafka.support.serializer.DeserializationException: failed to deserialize ... Caused by: com.amazonaws.services.schemaregistry.exception.AWSSchemaRegistryException: Exception occurred while de-serializing Avro message ... at com.amazonaws.services.schemaregistry.deserializers.GlueSchemaRegistryKafkaDeserializer.deserialize(GlueSchemaRegistryKafkaDeserializer.java:116) ~[schema-registry-serde-1.1.15.jar!/:na] at org.apache.kafka.common.serialization.Deserializer.deserialize(Deserializer.java:60) ~[kafka-clients-3.3.2.jar!/:na] at org.springframework.kafka.support.serializer.ErrorHandlingDeserializer.deserialize(ErrorHandlingDeserializer.java:179) ~[spring-kafka-3.0.4.jar!/:3.0.4] ... 15 common frames omitted Caused by: org.apache.avro.AvroTypeException: Found test.avro.TestData, expecting test.avro.TestData, missing required field testId ... 2023-08-31T09:06:47.955Z WARN 1 --- [ntainer#6-0-C-1] r.DeadLetterPublishingRecovererFactory$1 : Destination resolver returned non-existent partition test.TestData-test-dlt-1, KafkaProducer will determine partition to use for this topic 2023-08-31T09:06:47.955Z ERROR 1 --- [ntainer#6-0-C-1] c.a.s.schemaregistry.utils.AVROUtils : Unsupported Avro Data Formats 2023-08-31T09:06:47.956Z ERROR 1 --- [ntainer#6-0-C-1] c.a.s.schemaregistry.utils.AVROUtils : Unsupported Type of Record received 2023-08-31T09:06:47.956Z ERROR 1 --- [ntainer#6-0-C-1] r.DeadLetterPublishingRecovererFactory$1 : Dead-letter publication to test.TestData-test-dlt failed for: test.TestData-1@0 com.amazonaws.services.schemaregistry.exception.AWSSchemaRegistryException: Unsupported Type of Record received ... at com.amazonaws.services.schemaregistry.serializers.GlueSchemaRegistryKafkaSerializer.prepareInput(GlueSchemaRegistryKafkaSerializer.java:157) 2023-08-31T09:06:47.957Z ERROR 1 --- [ntainer#6-0-C-1] o.s.kafka.listener.DefaultErrorHandler : Failed to determine if this record (test.TestData-1@0) should be recovererd, including in seeks org.springframework.kafka.KafkaException: Dead-letter publication to test.TestData-test-dlt failed for:test.TestData-1@0 at org.springframework.kafka.listener.DeadLetterPublishingRecoverer.verifySendResult(DeadLetterPublishingRecoverer.java:669) ~[spring-kafka-3.0.4.jar!/:3.0.4] ...
问题根源
当发生DeserializationException时,ErrorHandlingDeserializer会返回null,但DeadLetterPublishingRecoverer的模板匹配逻辑是基于原始消息的value类型(而非反序列化后的null),同时null会匹配映射表中第一个兼容的类型(即Object.class对应的Glue序列化模板),导致尝试用Glue序列化器处理无法识别的原始字节,最终发送失败。
修复方案
方案1:基于异常类型选择DLT发送模板
自定义DeadLetterPublishingRecoverer的模板选择逻辑,识别反序列化异常场景,强制使用ByteArraySerializer的KafkaTemplate:
@Bean public DeadLetterPublishingRecoverer deadLetterPublishingRecoverer(KafkaTemplate<String, Object> kafkaTemplate) { DefaultKafkaProducerFactory<String, byte[]> defaultKafkaProducerFactory = new DefaultKafkaProducerFactory<>(kafkaProperties.buildProducerProperties(), new StringSerializer(), new ByteArraySerializer()); KafkaTemplate<String, byte[]> bytesKafkaTemplate = new KafkaTemplate<>(defaultKafkaProducerFactory); // 根据异常类型选择模板:反序列化异常使用字节模板,其余用Glue模板 return new DeadLetterPublishingRecoverer((record, exception) -> { Throwable rootCause = exception; while (rootCause.getCause() != null) { rootCause = rootCause.getCause(); } if (rootCause instanceof DeserializationException) { return bytesKafkaTemplate; } return kafkaTemplate; }); }
方案2:封装原始字节并匹配对应模板
通过ErrorHandlingDeserializer将反序列化失败的原始字节封装为自定义对象,再为该对象绑定字节序列化模板:
- 定义原始字节封装类:
public static class RawKafkaMessage { private final byte[] rawData; public RawKafkaMessage(byte[] rawData) { this.rawData = rawData; } public byte[] getRawData() { return rawData; } }
- 修改消费者工厂配置,设置反序列化失败时返回封装对象:
@Bean public ConsumerFactory<String, Object> consumerFactory() { ErrorHandlingDeserializer<Object> errorHandlingDeserializer = new ErrorHandlingDeserializer<>(kafkaAwsGlueValueDeserializer()); // 反序列化失败时,将原始字节封装为RawKafkaMessage errorHandlingDeserializer.setFailedDeserializationFunction((topic, data, ex) -> new RawKafkaMessage(data)); return new DefaultKafkaConsumerFactory<>( kafkaConsumerProps, new StringDeserializer(), errorHandlingDeserializer, false ); }
- 调整DLT恢复器的模板映射,优先匹配
RawKafkaMessage类型:
@Bean public DeadLetterPublishingRecoverer deadLetterPublishingRecoverer(KafkaTemplate<String, Object> kafkaTemplate) { DefaultKafkaProducerFactory<String, byte[]> defaultKafkaProducerFactory = new DefaultKafkaProducerFactory<>(kafkaProperties.buildProducerProperties(), new StringSerializer(), new ByteArraySerializer()); KafkaTemplate<String, byte[]> bytesKafkaTemplate = new KafkaTemplate<>(defaultKafkaProducerFactory); Map<Class<?>, KafkaOperations<?, ?>> templates = new LinkedHashMap<>(); templates.put(RawKafkaMessage.class, bytesKafkaTemplate); // 优先匹配封装类 templates.put(Object.class, kafkaTemplate); return new DeadLetterPublishingRecoverer(templates); }
额外优化
针对反序列化异常这类不可恢复错误,建议在RetryTopicConfiguration中跳过重试,直接发送到DLT:
@Bean public RetryTopicConfiguration retryTopicConfiguration(KafkaTemplate<String, Object> template) { return RetryTopicConfigurationBuilder .newInstance() .exponentialBackoff(1000, 2, 10000) .notRetryOn(DeserializationException.class) // 反序列化异常不重试 .create(template); }
内容的提问来源于stack exchange,提问作者CzarSkar
相关产品推荐
相关产品推荐

