【问题标题】:Kafka CommitFailedException consumer exceptionKafka CommitFailedException 消费者异常
【发布时间】:2016-06-10 01:11:20
【问题描述】:

创建多个消费者(使用 Kafka 0.9 java API)并启动每个线程后,我得到以下异常

Consumer has failed with exception: org.apache.kafka.clients.consumer.CommitFailedException: Commit cannot be completed due to group rebalance
class com.messagehub.consumer.Consumer is shutting down.
org.apache.kafka.clients.consumer.CommitFailedException: Commit cannot be completed due to group rebalance
at org.apache.kafka.clients.consumer.internals.ConsumerCoordinator$OffsetCommitResponseHandler.handle(ConsumerCoordinator.java:546)
at org.apache.kafka.clients.consumer.internals.ConsumerCoordinator$OffsetCommitResponseHandler.handle(ConsumerCoordinator.java:487)
at org.apache.kafka.clients.consumer.internals.AbstractCoordinator$CoordinatorResponseHandler.onSuccess(AbstractCoordinator.java:681)
at org.apache.kafka.clients.consumer.internals.AbstractCoordinator$CoordinatorResponseHandler.onSuccess(AbstractCoordinator.java:654)
at org.apache.kafka.clients.consumer.internals.RequestFuture$1.onSuccess(RequestFuture.java:167)
at org.apache.kafka.clients.consumer.internals.RequestFuture.fireSuccess(RequestFuture.java:133)
at org.apache.kafka.clients.consumer.internals.RequestFuture.complete(RequestFuture.java:107)
at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient$RequestFutureCompletionHandler.onComplete(ConsumerNetworkClient.java:350)
at org.apache.kafka.clients.NetworkClient.poll(NetworkClient.java:288)
at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.clientPoll(ConsumerNetworkClient.java:303)
at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:197)
at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:187)
at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:157)
at org.apache.kafka.clients.consumer.internals.ConsumerCoordinator.commitOffsetsSync(ConsumerCoordinator.java:352)
at org.apache.kafka.clients.consumer.KafkaConsumer.commitSync(KafkaConsumer.java:936)
at org.apache.kafka.clients.consumer.KafkaConsumer.commitSync(KafkaConsumer.java:905)

然后开始正常消费消息,我想知道是什么导致了这个异常以便修复它。

【问题讨论】:

  • 雨果,你还在遇到这个问题吗?你能提供更多信息吗?
  • 是的@nautilus,我仍然有这个问题。我有 3 个消费者,都在同一个消费者组中,我有一个有 20 个分区的主题,应该从中读取数据。这个异常是随机发生的,但是消费者可以从主题/分区中读取数据,尽管这个异常被触发了。
  • 消费者只是在消费数据还是也在处理数据?我在您的堆栈跟踪中看到,当您尝试提交同步偏移量时会发生异常,您能否描述在消息的消耗和偏移量的提交之间发生了什么?我认为您的消费者可能会失去与协调员的心跳。
  • 消费消息后已经提交了偏移量,但仍然出现异常。换句话说,无论是否触发异常,都会消耗所有消息。

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


【解决方案1】:

两个可能的原因 -->

  1. 如果出现任何网络故障,消费者无法联系代理并会抛出此异常。但发生这些异常时并没有出现网络故障。
  2. 如错误跟踪中所述,如果处理消息花费的时间过多,ConsumerCoordinator 将失去连接,提交将失败。这是因为轮询。

此处给出的值是默认的 Kafka 消费者配置值。

request.timeout.ms=40000
heartbeat.interval.ms=3000
max.poll.interval.ms=300000
max.poll.records=500
session.timeout.ms=10000

解决方案-->

将 max.poll.records 减少到 100 条但仍然,该异常有时会发生。所以改变了如下配置;

request.timeout.ms=300000
heartbeat.interval.ms=1000 
max.poll.interval.ms=900000
max.poll.records=100
session.timeout.ms=600000

减少了心跳间隔,以便代理会频繁更新消费者处于活动状态。并且还增加了会话超时配置。

【讨论】:

    【解决方案2】:

    同时尝试调整以下参数:

    • heartbeat.interval.ms - 这告诉 Kafka 在考虑消费者将被视为“死亡”之前等待指定的毫秒数
    • ma​​x.partition.fetch.bytes - 这将限制消费者在轮询时收到的消息数量(最多)。

    我注意到如果消费者在心跳超时之前没有提交到 Kafka,就会发生重新平衡。如果在处理消息之后发生提交,则处理它们的时间将决定这些参数。因此,减少消息数量并增加心跳时间将有助于避免重新平衡。

    还要考虑使用更多的分区,这样就会有更多的线程处理您的数据,即使每次轮询的消息更少。

    我编写了这个小应用程序来进行测试。希望对您有所帮助。

    https://github.com/ajkret/kafka-sample

    更新

    Kafka 0.10.x 现在提供了一个新参数来控制接收到的消息数量: - ma​​x.poll.records - 在一次 poll() 调用中返回的最大记录数。

    更新

    Kafka 提供了一种暂停队列的方法。当队列暂停时,您可以在单独的线程中处理消息,允许您调用 KafkaConsumer.poll() 来发送心跳。然后在处理完成后调用KafkaConsumer.resume()。通过这种方式,您可以缓解由于不发送心跳而导致重新平衡的问题。以下是您可以做什么的概述:

    while(true) {
        ConsumerRecords records = consumer.poll(Integer.MAX_VALUE);
        consumer.commitSync();
    
        consumer.pause();
        for(ConsumerRecord record: records) {
    
            Future<Boolean> future = workers.submit(() -> {
                // Process
                return true;
            }); 
    
    
           while (true) {
                try {
                    if (future.get(1, TimeUnit.SECONDS) != null) {
                        break;
                    }
                } catch (java.util.concurrent.TimeoutException e) {
                    getConsumer().poll(0);
                }
            }
        }
    
        consumer.resume();
    }
    

    【讨论】:

    • 版本 0.10.x 现在有一个新参数 max.poll.records,用于代替 max.partition.fetch.bytes。
    • 我使用了与您的暂停和恢复相同的方法,但仍然出现相同的错误。唯一的区别是我在 pause() 之后和 resume() 之前调用 commitSync(),因为只有在处理记录时我才需要提交。知道我做错了什么吗?
    • @mav3n:我遇到了同样的问题。尝试增加 session.timeout.ms 和 max.poll.records 但没有成功。找到方法了吗?
    • @HoàngLong :不,我没有得到任何解决方案。最后,我将设计更改为单一消费者以使其正常工作。如果您对此有解决方案,请分享您的发现。谢谢
    • @mav3n:我对这个异常有一些“发现”。首先我在Kafka服务器上增加了session.timeout.ms,现在也不例外,但这并不是最优的。然后我们将“类似”的方式应用到ajkret 的方法中,只是我们不使用pause() 和resume(),而是每隔几秒定期轮询() Kafka 服务器。需要计算出间隔,但它适用于我的情况。
    猜你喜欢
    • 2020-07-10
    • 1970-01-01
    • 2017-08-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-09-19
    • 2017-11-16
    • 1970-01-01
    相关资源
    最近更新 更多