【问题标题】:How to achieve Always retry policy/Custom Retry Policy for a Kafka Consumer without causing a rebalance in the group如何在不引起组重新平衡的情况下为 Kafka 消费者实现始终重试策略/自定义重试策略
【发布时间】:2020-02-13 19:48:55
【问题描述】:

我正在尝试编写一个有弹性的 Kafka 消费者。如果在侦听器方法中处理消息时出现异常,我想重试它。对于某些异常,我想重试几次,总是针对某些异常,从不针对其他异常。我已阅读有关 SeekToCurrentErrorHandler 的 Spring 文档,但不能 100% 确定如何实现它。

我对 ExceptionClassifierRetryPolicy 进行了子类化,并且正在根据 Listener 方法中发生的异常返回适当的重试策略。

我已经创建了 RetryTemplate,并在子类中使用我的自定义实现设置了它的 RetryPolicy。

我已经在 Kafka 容器上设置了 retryTemplate。我已将错误处理程序设置为新的 SeekToCurrentHandler,并将有状态重试属性设置为 true。

监听方法

@KafkaListener(topics = "topicName", containerFactory = "containerFactory")
   public void listenToKafkaTopic(@Payload Message<SomeAvroGeneratedClass> message, Acknowledgement ack){
      SomeAvroGeneratedClass obj = message.getPayLoad();
      processIncomingMessage(obj);
      ack.acknowledge();
   }

自定义重试策略类

 @Component
 public class MyRetryPolicy extends ExceptionClassifierRetryPolicy
   {
      @PostConstruct
       public void init(){
             final SimpleRetryPolicy simpleRetryPolicy = new SimpleRetryPolicy();
             simpleRetryPolicy.setMaxAttempts(8);

     this.setExceptionClassifier( classifiable ->
           {
        // Always Retry when instanceOf   TransientDataAccessException
       if( classifiable.getCause() instanceof TransientDataAccessException)
            {
               return new AlwaysRetryPolicy();                                         

            }
      else if(classifiable.getCause() instanceOf NonTransientDataAccessException)
           {
              return new NeverRetryPolicy();
           }
          else
            {
            return simpleRetryPolicy;
            }

      } );
 }}

重试模板和容器配置

@Configuration
public class RetryConfig{


 @Bean
 public RetryTemplate retryTemplate(@Autowired ConcurrentKafkaListenerContainerFactory factory){

   RetryTemplate retryTemplate = new RetryTemplate();
   retryTemplate.setRetryPolicy(new MyRetryPolicy());
   FixedBackOffPolicy fixedBackOffPolicy = new FixedBackOffPolicy()
   fixedBackOffPolicy.setBackOffPeriod(2000l);
   retryTemplate.setBackOffPolicy(fixedBackOffPolicy);

   factory.setRetryTemplate(retryTemplate);
   factory.setAckOnError(false);
   factory.setErrorHandler(new SeekToCurrentErrorHandler());
   factory.setStateFulRetry(true);
   factory.setRecoveryCallback(//configure recovery after retries are exhausted and commit offset
   );

  }
}

监听器属性:

  1. AckMode = 手动
  2. auto.offset.commit = false

问题:

  1. 使用我当前的代码,我能否在返回 AlwaysRetryPolicy 时实现我在 MyRetryPolicy 中定义的重试逻辑而不导致消费者重新平衡?如果没有,请指引我正确的道路。

  2. 我的方法在使用错误处理程序和重试时是否正确?

【问题讨论】:

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


    【解决方案1】:

    有状态重试与SeekToCurrentErrorHandler 结合使用是正确的方法。

    但是,如果您使用的是最新版本(2.2 或更高版本),错误处理程序将在 10 次尝试后放弃;您可以通过将 maxFailures 设置为 -1 (2.2) 或将 backOff 设置为 Long.MAX_VALUE(2.3 或更高版本)来实现无穷大。

    【讨论】:

    • 当使用有状态恢复时,恢复必须在监听器的重试恢复回调中进行,而不是STCEH——否则会出现内存泄漏(重试模板中的状态) )。即,recoveryCallback 应该在重试用尽后正常返回,以获得可重试的异常。侦听器正常退出时不会调用 STCEH,因此 maxAttempts 设置不适用。
    • 是的;那是正确的;当重试用尽时,会调用回调并执行一些操作(记录、发送到 DLQ 主题等)。如果正常退出,我们不会重新搜索记录。
    • SeekToCurrentErrorHandler 的 2.1 版本总是重试,所以一切都会如您所愿。
    • 是的,但是当前的spring-kafka 2.4.x版本是2.4.12。 spring.io/projects/spring-kafka#learn 我不建议使用 2.4.0。
    • 由于你在重试模板中也有一个回退,所以它将是两者的总和;最好只在一个地方定义回退并将另一个设置为零。
    猜你喜欢
    • 2015-10-06
    • 2019-12-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-06-21
    • 1970-01-01
    • 2020-08-04
    相关资源
    最近更新 更多