【问题标题】:Avro Deserialization exception handling with Spring Cloud Stream使用 Spring Cloud Stream 处理 Avro 反序列化异常
【发布时间】:2021-05-21 14:02:57
【问题描述】:

我有一个使用 Spring Cloud Stream 和 Spring Kafka 的应用程序,它处理 Avro 消息。该应用程序运行良好,但现在我想添加一些错误处理。

目标:我想捕获反序列化异常,用异常详细信息 + 原始 Kafka 消息 + 自定义上下文信息构建一个新对象,并将此对象推送到专用 Kafka 主题。基本上是一个DLQ,但是原始消息会被截取和修饰。

问题:虽然我可以拦截异常,但我不知道如何从 Kafka 获取原始消息(TODO 1,下面)。我已经浏览过 ConsumerAwareErrorHandler.handle 中返回的数据对象,但我没有看到它。

下面是我的代码:

@EnableBinding(EventStream.class)
@SpringBootApplication
@Slf4j
public class SpringcloudApplication {

    public static void main(String[] args) {
        SpringApplication.run(SpringcloudApplication.class, args);
    }

    /* Configure custom exception handler */
    @Bean
    public ListenerContainerCustomizer<AbstractMessageListenerContainer<?, ?>> cust() {
        return (container, destination, group) -> {
            container.setErrorHandler(new ConsumerAwareErrorHandler() {
                @Override
                public void handle(Exception thrownException, ConsumerRecord<?, ?> data, Consumer<?, ?> consumer) {
                    log.info("Got error with data: {}", data);
                    // TODO 1 - How to get original message?
                    // TODO 2 - Send to dedicated (DLQ) topic
                }
            });
        };
    }

    @StreamListener(EventStream.INBOUND)
    public void consumeEvent(@Payload Message message) {
        log.info("Consuming event --> {}", message.toString());
        produceEvent(message);
    }

    @Autowired private EventStream eventStream;
    public Boolean produceEvent(Message message) {
        log.info("Producing event --> {}", message.toString());
        return eventStream
                .producer()
                .send(MessageBuilder.withPayload(message)
                    .setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.APPLICATION_JSON)
                    .build());
    }

}

还有属性文件:

spring:
  cloud:
    stream:
      default-binder: kafka
      default:
        consumer:
          useNativeEncoding: true
        producer:
          useNativeEncoding: true
      kafka:
        binder:
          brokers: localhost:9092
          producer-properties:
            key.serializer: org.apache.kafka.common.serialization.StringSerializer
            value.serializer: io.confluent.kafka.serializers.KafkaAvroSerializer
            schema.registry.url: "http://localhost:8081"
          consumer-properties:
            key.deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
            value.deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
            schema.registry.url: "http://localhost:8081"
            specific.avro.reader: true
            spring.deserializer.key.delegate.class: org.apache.kafka.common.serialization.StringDeserializer
            spring.deserializer.value.delegate.class: io.confluent.kafka.serializers.KafkaAvroDeserializer
      bindings:
        event-consumer:
          destination: data_stream_in # incoming topic
          contentType: application/**avro
          group: data_stream_consumer
        event-producer:
          destination: data_stream_out
          contentType: application/**avro

我正在使用以下版本:

  • Spring Boot 2.3.2.RELEASE
  • Spring Cloud:Hoxton.SR8
  • spring-cloud-stream-binder-kafka 3.0.8.RELEASE
  • spring-kafka 2.5.12

感谢任何帮助!

【问题讨论】:

    标签: avro spring-kafka spring-cloud-stream spring-cloud-stream-binder-kafka


    【解决方案1】:

    handle 方法中的第二个参数是 ConsumerRecord,它是原始 Kafka 记录,但如果您希望记录自动发送到 DLQ,您可以执行以下操作。

    @Bean
        public ListenerContainerCustomizer<AbstractMessageListenerContainer<byte[], byte[]>> customizer(SeekToCurrentErrorHandler errorHandler) {
            return (container, dest, group) -> {
                container.setErrorHandler(errorHandler);
            };
        }
    
    
    
        @Bean
        public SeekToCurrentErrorHandler errorHandler(DeadLetterPublishingRecoverer deadLetterPublishingRecoverer) {
            return new SeekToCurrentErrorHandler(deadLetterPublishingRecoverer);
        }
    
    
        @Bean
        public DeadLetterPublishingRecoverer publisher(KafkaOperations bytesTemplate) {
            return new DeadLetterPublishingRecoverer(bytesTemplate);
        }
    

    本质上,您正在设置一个SeekToCurrentErrorHandler,它能够将失败的记录发送到DLQ。有关 DeadLetterPublishingRecoverer 如何工作的更多详细信息,请参阅 Spring for Apache Kafka 的参考文档:https://docs.spring.io/spring-kafka/docs/current/reference/html/#dead-letters ; 更多信息SeekToCurrentErrorHandler:https://docs.spring.io/spring-kafka/docs/current/reference/html/#seek-to-current

    你还需要配置和ErrorHandlingDeserializer,

    spring.cloud.stream.kafka.binder.configuration.value.deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
    spring.cloud.stream.kafka.binder.configuration.spring.deserializer.key.delegate.class: org.apache.kafka.common.serialization.StringDeserializer
    ...
    similar for the value class. 
    

    更多信息ErrorHandlingDeserializer: https://docs.spring.io/spring-kafka/docs/current/reference/html/#error-handling-deserializer

    如果您想修改记录并向DLQ添加自定义消息,您可以通过覆盖handle方法然后访问ConsumerRecord然后调用超类方法来实现。

    【讨论】:

    • 我已经按照您的建议添加了配置并覆盖了句柄方法,但是 ConsumerRecord 中仍然没有数据。我有原始消息的标头,但没有有效负载。
    • 消费者记录值为null,因为无法反序列化。该值在DeserializationException 中(这是ListenerExecutionFailedException 的原因)。
    • 谢谢加里,这是有道理的。 @sobychacko,您能否详细说明您关于向 DLQ 添加自定义消息的最后一点。如果我调用超类方法,它仍然只会将 ConsumerRecord 作为第二个参数,因此我将如何发送自定义对象?
    猜你喜欢
    • 2020-07-04
    • 1970-01-01
    • 1970-01-01
    • 2017-11-24
    • 2018-01-31
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-06-23
    相关资源
    最近更新 更多