【问题标题】:Kafka consumer infinite retry | Rebalance issue on scalingKafka消费者无限重试 |缩放的重新平衡问题
【发布时间】:2021-01-14 19:20:21
【问题描述】:

我的 Kafka 消费者面临一些奇怪的再平衡问题。我已经使用SeekToCurrentErrorHandler 设置max.poll.interval.ms=1,200,000 即20 分钟并将重试延迟设置为900,000 15 分钟延迟,实现了一个无限重试策略。
如您所见,poll.interval > 延迟。我面临的问题是当一个新的消费者被添加时,它重新平衡并且旧的消费者离开了组,只有新的消费者接收并处理消息。在新的消费者中,我看到了日志Attempt to heartbeat failed since group is rebalancing,但它仍然在处理消息。老消费者收到0个数据。

我的消费者配置如下:

spring:
    kafka:
       ......
       ......
        consumer:
             max.poll.records: 150
             group.id: xxxxx
             properties:
                   enable.auto.commit: false
                   max.poll.interval.ms: 1200000 #20 minutes greater than retry interval

Kafka java 配置:

public ConcurrentkafkaListenerContainerFactory<String, byte[]> kafkaFactory() {
    
       ConcurrentkafkaListenerContainerFactory<String, byte[]> = new 
                                ConcurrentkafkaListenerContainerFactory();
        ......
        ......
   factory.setErrorHandler(kafkaErrorHandler);
   facory.setRetryTemplate(retrtTemplate());
   factory.setStatefulRetry(true)
   factory.getContainerProperties().setAckOnError(false);
   factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANIAL_IMMEDIATE);
   return factory;
  }

错误处理程序:

@Component
public class KafkaErrorHandler extends SeekToCurrentErrorHandler {
   KafkaErrorHanbdler(){super(-1)}
    @Override
        public void handle(Exception thrownException, List<ConsumerRecord<?, ?>> records, Consumer<?, 
          ?> consumer, MessageListenerContainer container) {

            LOG.info("handle");
            super.handle(thrownException, records, consumer, container);
        }
    
}

我的应用程序需要大约 10 秒来处理每个事件,并且只有在成功处理后才为每个事件发送确认。按照配置,它需要 150 条消息。最大轮询间隔设置为 20 分钟。 kafka 是否仅在处理 150 条消息后才进行轮询?

【问题讨论】:

    标签: kafka-consumer-api spring-kafka


    【解决方案1】:

    通过有状态重试,将异常抛出到容器中,我们将重新寻找未处理的记录并再次轮询。

    DEBUG 日志记录应该可以帮助您找出问题所在。

    对于较新的版本(自 2.3 起),您可以使用 BackOff 而不是重试模板配置 SeekToCurrentErrorHandler

    在当前消费者忙碌时添加新消费者意味着重新平衡将延迟到第一个消费者再次轮询。现有消费者在超时之前不会“离开组”。

    【讨论】:

      猜你喜欢
      • 2017-06-19
      • 2018-05-23
      • 2020-08-04
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-02-13
      • 2022-10-05
      • 2021-01-26
      相关资源
      最近更新 更多