我正在使用 spring kafka 2.7.8 来消费来自 kafka 的消息。消费者监听如下
@KafkaListener(topics = "topicName",
groupId = "groupId",
containerFactory = "kafkaListenerFactory")
public void onMessage(ConsumerRecord record) {
}
Run Code Online (Sandbox Code Playgroud)
上面的 onMessage 方法一次接收一条消息。
这是否意味着 max.poll.records 被 spring 库设置为 1 或者它一次轮询 500 个(默认值)并且该方法一个接一个地接收。
提出这个问题的原因是,我们经常在产品中看到以下错误。在不到一分钟的时间内收到多个消费者的以下所有 4 个错误。试图了解这是由于间歇性的 kafka 代理连接问题还是由于负载。请指教。
Seek to current after exception; nested exception is org.apache.kafka.clients.consumer.CommitFailedException: Offset commit cannot be completed since the consumer is not part of an active group for auto partition assignment; it is likely that the consumer was kicked out of the group.
seek to current after exception; nested exception is org.apache.kafka.common.errors.TimeoutException: Timeout of 60000ms expired before successfully committing offsets {topic-9=OffsetAndMetadata{offset=2729058, leaderEpoch=null, metadata=''}}
Consumer clientId=consumer-groupName-5, groupId=consumer] Offset commit failed on partition topic-33 at offset 2729191: The coordinator is not aware of this member.
Seek to current after exception; nested exception is org.apache.kafka.clients.consumer.CommitFailedException: Commit cannot be completed since the group has already rebalanced and assigned the partitions to another member. This means that the time between subsequent calls to poll() was longer than the configured max.poll.interval.ms, which typically implies that the poll loop is spending too much time message processing. You can address this either by increasing max.poll.interval.ms or by reducing the maximum size of batches returned in poll() with max.poll.records
max.poll.records不被 Spring 改变;它将采用默认值(或您设置的任何值)。在下一次轮询之前,一次将一条记录交给侦听器。
这意味着您的侦听器必须能够max.poll.records在max.poll.interval.ms.
您需要减少max.poll.records和/或增加,max.poll.interval.ms以便您可以在这段时间内处理记录,并留出足够的余量,以避免这些重新平衡。
| 归档时间: |
|
| 查看次数: |
9212 次 |
| 最近记录: |