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

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将反序列化失败的原始字节封装为自定义对象,再为该对象绑定字节序列化模板:

  1. 定义原始字节封装类:
public static class RawKafkaMessage {
    private final byte[] rawData;

    public RawKafkaMessage(byte[] rawData) {
        this.rawData = rawData;
    }

    public byte[] getRawData() {
        return rawData;
    }
}
  1. 修改消费者工厂配置,设置反序列化失败时返回封装对象:
@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
    );
}
  1. 调整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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 08:07:03