【问题标题】:Spring Kafka - Display Custom Error Message with @RetrySpring Kafka - 使用 @Retry 显示自定义错误消息
【发布时间】:2020-06-24 19:59:03
【问题描述】:

我目前正在使用 Spring Kafka 以及 Spring 的 @Retry 来使用来自主题的消息。所以基本上,如果出现错误,我会重试处理消费者消息。但是在这样做的同时,我想避免KafkaMessageListenerContainer 抛出的异常消息。相反,我想显示一条自定义消息。我尝试在ConcurrentKafkaListenerContainerFactory 中添加错误处理程序,但这样做时,我的重试不会被调用。 有人可以指导我如何显示自定义异常消息以及@Retry 场景吗?以下是我的代码 sn-ps:

ConcurrentKafkaListenerContainerFactory Bean 配置


@Bean
ConcurrentKafkaListenerContainerFactory << ? , ? > concurrentKafkaListenerContainerFactory(ConcurrentKafkaListenerContainerFactoryConfigurer configurer, ConsumerFactory < Object, Object > kafkaConsumerFactory) {
    ConcurrentKafkaListenerContainerFactory < Object, Object > kafkaListenerContainerFactory =
        new ConcurrentKafkaListenerContainerFactory < > ();
    configurer.configure(kafkaListenerContainerFactory, kafkaConsumerFactory);
    kafkaListenerContainerFactory.setConcurrency(1);

    // Set Error Handler
    /*kafkaListenerContainerFactory.setErrorHandler(((thrownException, data) -> {
        log.info("Retries exhausted);
    }));*/
    return kafkaListenerContainerFactory;
}

卡夫卡消费者

@KafkaListener(
    topics = "${spring.kafka.reprocess-topic}",
    groupId = "${spring.kafka.consumer.group-id}",
    containerFactory = "concurrentKafkaListenerContainerFactory"
)
@Retryable(
    include = RestClientException.class,
    maxAttemptsExpression = "${spring.kafka.consumer.max-attempts}",
    backoff = @Backoff(delayExpression = "${spring.kafka.consumer.backoff-delay}")
)
public void onMessage(ConsumerRecord < String, String > consumerRecord) throws Exception {
    // Consume the record
    log.info("Consumed Record from topic : {} ", consumerRecord.topic());


    // process the record
    messageHandler.handleMessage(consumerRecord.value());
}

以下是我得到的例外:

【问题讨论】:

    标签: spring-boot apache-kafka spring-kafka


    【解决方案1】:

    您不应该使用 @RetryableSeekToCurrentErrorHandler(现在是默认值,从 2.5 开始;所以我认为您正在使用该版本)。

    相反,配置一个自定义SeekToCurrentErrorHandler 具有最大尝试次数、回退和可重试异常。

    那个错误信息是正常的;它由容器记录;通过在SeekToCurrentErrorHandler 上设置logLevel 属性,可以将其日志记录级别从ERROR 降低到INFO 或DEBUG。您还可以向其中添加自定义恢复器,以便在重试完成后记录您的自定义消息。

    【讨论】:

    • 感谢您的意见!您是否可以添加代码 sn-p 来配置 SeekToCurrentErrorHandler 最大尝试次数、回退和可重试异常?使用@RetrySeekToCurrentErrorHandler 是否有任何问题
    • the reference manual。第一个示例有一个回退和自定义恢复器。向下滚动一点以查看如何添加不可重试的异常。
    • @GaryRussell 为什么该日志消息正常?为什么异常处理程序本身要抛出异常?
    【解决方案2】:

    我的事件重试模板,

    @Bean(name = "eventRetryTemplate")
    public RetryTemplate eventRetryTemplate() {
    RetryTemplate template = new RetryTemplate();
    
    ExceptionClassifierRetryPolicy retryPolicy = new ExceptionClassifierRetryPolicy();
    Map<Class<? extends Throwable>, RetryPolicy> policyMap = new HashMap<>();
    policyMap.put(NonRecoverableException.class, new NeverRetryPolicy());
    policyMap.put(RecoverableException.class, new AlwaysRetryPolicy());
    retryPolicy.setPolicyMap(policyMap);
    
    ExponentialBackOffPolicy backOffPolicy = new ExponentialBackOffPolicy();
    backOffPolicy.setInitialInterval(backoffInitialInterval);
    backOffPolicy.setMaxInterval(backoffMaxInterval);
    backOffPolicy.setMultiplier(backoffMultiplier);
    
    template.setRetryPolicy(retryPolicy);
    template.setBackOffPolicy(backOffPolicy);
    return template;
    }
    

    我的 kafka 监听器使用重试模板,

    @KafkaListener(
      groupId = "${kafka.consumer.group.id}",
      topics = "${kafka.consumer.topic}",
      containerFactory = "eventContainerFactory")
      public void eventListener(ConsumerRecord<String, String> 
      events,
      Acknowledgment acknowledgment) {
      eventRetryTemplate.execute(retryContext -> {
      retryContext.setAttribute(EVENT, "my-event");
      eventConsumer.consume(events, acknowledgment);
      return null;
      });
    }
    

    我的 kafka 消费者属性,

    private ConcurrentKafkaListenerContainerFactory<String, String> 
    getConcurrentKafkaListenerContainerFactory(
      KafkaProperties kafkaProperties) {
    ConcurrentKafkaListenerContainerFactory<String, String> factory =
        new ConcurrentKafkaListenerContainerFactory<>();
    factory.getContainerProperties().setAckMode(AckMode.MANUAL_IMMEDIATE);
    factory.getContainerProperties().setAckOnError(Boolean.TRUE);
    kafkaErrorEventHandler.setCommitRecovered(Boolean.TRUE);
    factory.setErrorHandler(kafkaErrorEventHandler);
    factory.setConcurrency(1); 
    factory.getContainerProperties().setAckMode(AckMode.MANUAL_IMMEDIATE);
    factory.setConsumerFactory(getEventConsumerFactory(kafkaProperties));
    return factory;
    }
    

    kafka 错误事件处理程序是我的自定义错误处理程序,它扩展了 SeekToCurrentErrorHandler 并实现了一些类似这样的处理错误方法.....

    @Override
    public void handle(Exception thrownException, List<ConsumerRecord<?, ?>> 
    records,
      Consumer<?, ?> consumer, MessageListenerContainer container) {
    
    log.info("Non recoverable exception. Publishing event to Database");
    super.handle(thrownException, records, consumer, container);
    
    ConsumerRecord<String, String> consumerRecord = (ConsumerRecord<String, 
    String>) records.get(0);
    
    FailedEvent event = createFailedEvent(thrownException, consumerRecord);
    
    failedEventService.insertFailedEvent(event);
    
    log.info("Successfully Published eventId {} to Database...", 
    event.getEventId());
    }
    

    这里失败的事件服务再次是我的自定义类,它将这些失败的事件放入可查询的关系数据库中(我选择它作为我的 DLQ)。

    【讨论】:

      猜你喜欢
      • 2021-06-06
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2010-11-07
      • 2017-07-03
      • 2016-09-19
      • 2021-12-01
      • 1970-01-01
      相关资源
      最近更新 更多