【问题标题】:How to know when consumer.poll() is called while using @KafkaListener annotation?如何知道使用 @KafkaListener 注释时何时调用 consumer.poll()?
【发布时间】:2020-07-06 22:49:21
【问题描述】:

我知道如果我使用@KafkaListener,我无法控制何时进行轮询,我从this answer 读到

下一个poll() 是在侦听器处理完上一个轮询的最后一条消息之后执行的。

所以我想知道如何知道每个poll() 何时执行?或者等效地,处理每个poll() 调用中收到的所有消息需要多长时间?

我问是因为我的程序出现“偏移提交失败......请求超时”异常,我想调整我的消费者配置,即max.poll.interval.ms 和max.poll.records,但我需要知道当前性能第一。

如果有帮助,这是我的@KafkaListener 方法的一部分:

@KafkaListener(id = "dataListener", topics ="${spring.kafka.topic}", containerFactory = "kafkaListenerContainerFactory")
public void listen(@Payload(required = false)  ConsumerRecord payload, @Header(KafkaHeaders.RECEIVED_PARTITION_ID) String partition,
                        @Header(KafkaHeaders.OFFSET)Long offset, @Header(KafkaHeaders.RECEIVED_MESSAGE_KEY)String messageKey){
    // processing messages
}

【问题讨论】:

  • @Gary Russell 你能解释一下吗?非常感谢!

标签: java apache-kafka spring-kafka


【解决方案1】:

您可以通过打开调试日志来查看轮询活动。

this.logger.debug(() -> "Received: " + records.count() + " records");

在每个poll() 之后记录。

【讨论】:

  • 谢谢加里!通过这一行,我现在可以找到source code,对于像我一样使用log4j 的人,您可以设置log4j.logger.org.springframework.kafka.listener=debug 但保留log4j.logger.org.springframework.kafka.listener.adapter=info,否则它将记录收到的每条原始消息。
猜你喜欢
  • 2019-09-01
  • 2018-10-19
  • 1970-01-01
  • 2020-11-25
  • 2010-11-05
  • 2013-10-10
  • 2016-04-17
  • 2019-03-19
  • 2011-11-10
相关资源
最近更新 更多