【问题标题】:Kafka Rejoining Group Quite Frequently卡夫卡经常重新加入团队
【发布时间】:2022-02-08 18:26:10
【问题描述】:

我有一个 Spring Boot 应用程序,它由 10 个不同的使用者组成,试图使用来自 10 个不同主题的消息。该应用程序使用相同的消费者组,因为所有 10 个不同的主题都列在同一个消费者组下。现在在运行应用程序一段时间后(5 小时后),我看到消费者正在尝试重新平衡,有时重新平衡永远不会完成,因此一段时间后无法接收消息。

这是我的消费者配置。

public Map<String, Object> consumerConfigs(String clientID) {
    Map<String, Object> props = new HashMap<>();
    props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);

    props.put("key.deserializer", "org.apache.kafka.common.serialization.ByteArrayDeserializer");
    props.put("value.deserializer", "org.apache.kafka.common.serialization.ByteArrayDeserializer");
    props.put(ConsumerConfig.GROUP_ID_CONFIG, consumerGroup);
    props.put(SECURITY_PROTOCOL, securityProtocol);
    props.put(SASL_MECHANISM, saslMechanism);
    props.put(SASL_JAAS_CONFIG,
            String.format("%s required username=\"%s\" password=\"%s\" ;", loginModule, username, kafkaSecretPass));

    props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
    props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
    props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 5);
    props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 10000);
    props.put(ConsumerConfig.CLIENT_ID_CONFIG, clientID);
    props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG,
            "org.apache.kafka.clients.consumer.RoundRobinAssignor");
   //props.put(ConsumerConfig.REQUEST_TIMEOUT_MS_CONFIG, 120000);
  //props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, 60000);
  //props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 600000);
 //props.put(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG, 2097164);
 //props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 30000);
     return props;
}

这是我在运行消费者一段时间后看到的错误。

[Consumer clientId=consumer-one, groupId=group-1] (Re-)joining group

这种重新加入群组的尝试持续了很长时间,最终它永远无法加入群组。

以下是其中一位消费者使用主题并将其发送以进行处理的方式。

  @EventListener(ApplicationStartedEvent.class)
public Disposable consume() {
    return reactiveKafkaConsumerTemplate
            .receive()
            .concatMap(consumerRecord -> {
                log.info("Received event:customer offset {} and value {}",
                        consumerRecord.offset(),new String(consumerRecord.value(), StandardCharsets.UTF_8));
                return updateCustomer(consumerRecord).flatMap(success -> {
                    consumerRecord.receiverOffset().acknowledge();
                    return Mono.just("Received");
                });
            })
            .retryWhen(Retry.backoff(3, Duration.ofSeconds(2)).transientErrors(true))
            .onErrorResume(e -> {
                log.info("Customer - Receiver On Error Resume");
                return Mono.empty();
            })
            .repeat()
            .subscribe(suc -> log.info("Subscribe successfully in Customer"), err -> log.error("Error occurred during subscribe Customer, " + err));
}

【问题讨论】:

  • The application uses same consumer group as all the 10 different topics are listed under the same consumer group - ```所有 10 个不同的主题都列在同一个消费者组下是什么意思``` - 你怎么说主题在一个组下。您指的是主题中的分区吗?
  • 还添加您如何订阅主题/将主题分配给消费者。这可能会增加清晰度
  • @Umeshwaran :添加了我们订阅主题的方式。 ReactitveKafkaConsumer 模板是根据上述消费者配置创建的。每个topic只有一个partition,并且都属于一个consumer group,用props.put(ConsumerConfig.GROUP_ID_CONFIG, consumerGroup);描述。
  • 消费者再平衡的原因不止一个。我给出了最常见的一种。尝试检查消费者是否花费了太长时间来处理消息。如果不是这种情况,请尝试为每个消费者提供唯一的消费者组 ID,以便其他主题的其他消费者可以毫无问题地消费消息
  • 在不相关的主题中使用相同的组 id 不是一个好习惯——正是因为这个原因;一个主题的重新平衡会导致所有主题的重新平衡。

标签: apache-kafka kafka-consumer-api spring-kafka


【解决方案1】:

频繁的重新平衡通常是因为消费者处理批次花费的时间太长。发生这种情况是因为消费者长时间处理批处理(并且没有发送心跳),因此代理认为消费者丢失并开始重新平衡。

我建议通过减少 max.partition.fetch.bytes 的值来创建更小的批次,或者通过增加 heartbeat.interval.ms 的值来延长/增加心跳间隔。

欲了解更多信息:More info

【讨论】:

    猜你喜欢
    • 2016-12-17
    • 2018-03-06
    • 1970-01-01
    • 1970-01-01
    • 2018-12-04
    • 1970-01-01
    • 2021-03-17
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多