【问题标题】:Spring Kafka- When is exactly Consumer.poll() called behind the hood?Spring Kafka-何时在后台调用 Consumer.poll()?
【发布时间】:2018-10-19 01:30:02
【问题描述】:

我有一个 Spring Boot 应用程序,其中有一个 Kafka Consumer。

我正在使用具有默认消费者配置的 DefaultKafkaConsumerFactory。我有一个并发设置为 1 的 ConcurrentListenerContainerFactory,我有一个用 @KafkaListener 注释的方法。

我正在收听一个具有 3 个分区的主题,并且我有 3 个这样的消费者分别部署在不同的应用程序中。因此,每个消费者都在监听一个分区。

假设在后台调用消费者轮询并获取 40 条记录。然后是每条记录,依次提供给带有@KafkaListener注解的方法,即提供记录1,等待方法完成处理,提供记录2,等待方法完成处理等等。 是出现上述情况,还是每获取一条记录,都会创建一个单独的线程,方法调用发生在单独的线程上,这样主线程就不会阻塞,可以更快的轮询记录。

我还想更清楚地了解消息侦听器容器是什么以及最终的消息侦听器。

提前谢谢你。

【问题讨论】:

    标签: java apache-kafka spring-kafka


    【解决方案1】:

    嗯,这正是 Apache Kafka 的立场 - 保证订单处理来自同一线程中同一分区的记录。因此,当您在 3 个实例之间分配具有 3 个分区的主题时,每个实例都会获得自己的分区并在单个线程中进行轮询。

    KafkaMessageListenerContainer 是围绕KafkaConsumer 的事件驱动、自我控制的包装器。它确实在while (isRunning()) { 循环中调用poll(),该循环被安排在TaskExecutor 中:

    this.listenerConsumerFuture = containerProperties
                .getConsumerTaskExecutor()
                .submitListenable(this.listenerConsumer);
    

    它处理ConsumerRecords调用监听器:

    private void invokeListener(final ConsumerRecords<K, V> records) {
            if (this.isBatchListener) {
                invokeBatchListener(records);
            }
            else {
                invokeRecordListener(records);
            }
        }
    

    【讨论】:

      【解决方案2】:

      在 1.3 及更高版本中,每个消费者只有一个线程;下一个poll() 是在侦听器处理完上一个轮询的最后一条消息之后执行的。

      在早期版本中,有两个线程,并且在侦听器线程处理第一批时执行了第二个(可能是第三个)轮询。这是为了避免由于侦听器速度慢而导致的重新平衡。线程模型非常复杂,我们必须在必要时暂停/恢复消费者。 KIP-62 修复了重新平衡问题,因此我们能够使用当今使用的更简单的线程模型。

      【讨论】:

      • 感谢您的回复。这正是我的想法,但只是想确认一下。当 Kafka 获取的记录需要大量时间处理时,我们遇到了问题。正如你用 kip-62 指出的那样,心跳是在后台线程中发送的,但即使调用 poll() 也会有超时,如果你在默认情况下一段时间(300000 毫秒)不调用 poll ,消费者就会死去,而且我们正在使用自动提交,未提交偏移量,并且再次处理相同的记录。
      • 你应该调整max.poll.records和max.poll.interval.ms,使得在处理投票结果时监听器不会超过后者。
      猜你喜欢
      • 1970-01-01
      • 2019-09-20
      • 2019-07-26
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-11-16
      • 2019-09-16
      • 1970-01-01
      相关资源
      最近更新 更多