如何使用ErrorHandlingDeserializer配合Spring Cloud Stream Kafka binder发送错误记录到默认DLQ主题
问题根因
你当前配置的ErrorHandlingDeserializer默认会在反序列化失败后直接返回null并丢弃异常,Spring Cloud Stream Kafka binder的原生DLQ逻辑感知不到这类异常,自然不会触发DLQ投递。你不需要自行定义DeadLetterPublishingRecoverer,只要调整少量配置就能复用binder已有的全部Kafka配置,将反序列化错误投递到默认DLQ主题。
配置修改方案
仅需调整my-in-0消费者的对应配置即可,不需要修改全局binder配置,也不需要新增冗余的Kafka参数:
- 开启binder对反序列化异常的自动识别能力
- 配置
ErrorHandlingDeserializer将反序列化异常写入消息头,而非直接吞掉
修改后的对应配置片段如下:
spring: cloud: stream: kafka: bindings: my-in-0: consumer: enableDlq: true autoCommitOnError: true # 新增:开启反序列化异常处理开关 handle-deserialization-exceptions: true configuration: key.deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer value.deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer spring.deserializer.key.delegate.class: io.confluent.kafka.serializers.KafkaAvroDeserializer spring.deserializer.value.delegate.class: io.confluent.kafka.serializers.KafkaAvroDeserializer # 新增:反序列化失败不直接抛异常,将异常信息存入消息头传递给binder spring.deserializer.key.fail-on-error: false spring.deserializer.value.fail-on-error: false dlqProducerProperties: configuration: key.serializer: org.apache.kafka.common.serialization.ByteArraySerializer value.serializer: org.apache.kafka.common.serialization.ByteArraySerializer
效果说明
以上配置适用于Spring Cloud Stream 3.1及以上版本,是官方原生支持的实现方式。修改后,反序列化失败的消息会自动被binder捕获,完全复用你在binder层级配置的所有Kafka连接、安全、序列化参数,直接投递到你期望的默认DLQ主题error.my-in-topic.myGroup。
内容的提问来源于stack exchange,提问作者armkillbill
相关产品推荐
相关产品推荐

