【发布时间】: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