【问题标题】:Rebalance issue with spring kafka max.poll.interval.ms, max.poll.records and idleTimeBetweenPollsspring kafka max.poll.interval.ms、max.poll.records 和 idleTimeBetweenPolls 的重新平衡问题
【发布时间】:2021-09-08 06:03:18
【问题描述】:

我看到我的应用程序在不断地重新平衡。我的应用程序是在批处理模式下开发的,这里是已添加的配置属性。

myapp.consumer.group.id= cg-id-local
myapp.changefeed.topic= test_topic
myapp.auto.offset.reset=latest
myapp.enable.auto.commit=false
myapp.max.poll.interval.ms=300000
myapp.max.poll.records= 20000
myapp.idle.time.between.polls=240000
myapp.concurrency = 10

容器工厂:

ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();

        factory.setConsumerFactory(consumerFactory(poSummaryCGID));
        factory.setConcurrency(poSummNoOfConsumers);
        factory.setBatchListener(true);
        factory.setAckDiscarded(true);
        factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
        factory.getContainerProperties().setIdleBetweenPolls(idleTimeBetweenPolls);

我这里有几个问题:

  1. 我已设置每次轮询(4 分钟)的最大记录数为 20000,并且我们在一个 TOPIC 中有 10 个分区。由于我将并发设置为 10,因此 10 个消费者将启动并运行,每个消费者将监听 1 个分区。我的问题是,记录数是否会像每个消费者可以处理 2000 条记录一样分配给所有消费者?

  2. max.poll.interval.ms 已设置为 5 分钟。我确信消费者将在给定的轮询间隔(4 分钟)内处理 2000 条(如果我的上述理解是正确的)记录,该间隔小于具有上限的 max.poll.interval.ms。但不确定为什么会发生再平衡?我还需要设置其他配置属性吗?

我们将不胜感激!

Tried with these configurations: 

myapp.max.poll.interval.ms=600000 
myapp.max.poll.records= 2000 
myapp.idle.time.between.polls=360000 

myapp.max.poll.interval.ms=300000 
myapp.max.poll.records= 2000 
myapp.idle.time.between.polls=300000 

myapp.max.poll.interval.ms=300000 
myapp.max.poll.records= 2000 
myapp.idle.time.between.polls=180000

编辑修复: 我们应该永远 myapp.max.poll.interval.ms > (myapp.idle.time.between.polls + myapp.max.poll.records 处理时间)。

【问题讨论】:

    标签: spring-kafka


    【解决方案1】:

    没有。 max.poll.records 是每个消费者,而不是每个主题或容器。

    如果您有 concurrency=10 和 10 个分区,则应将 max.poll.records 减少到 2000,以便每个消费者每次轮询最多获得 2000 个。

    容器将自动减少轮询之间的空闲,以便不会超过 max.poll.interval.ms,但您应该对这些属性(max.poll.recordsmax.poll.interval.ms)保持保守,这样它就永远不可能超过间隔。

    【讨论】:

    • 感谢@Gary 的及时支持!!。我对第二个答案感到困惑。我明确将轮询之间的空闲时间设置为 4 分钟,以便每批事件将在 max.poll.interval(5 分钟)内处理。我可以将此 max.poll.interval 增加到 10 分钟,以便我的应用程序将处理更多事件,例如 10000 条记录/消费者,理想的轮询间隔为 6 分钟?
    • 我的意思是,假设您的消费者需要 2 分钟来处理其记录。如果我们睡了配置的 4 分钟,我们将超过 5 分钟并导致重新平衡;容器会将其睡眠时间减少到 3 分钟(- 5 秒以避免比赛)。如果你想让你的消费者慢下来,你只需要设置idleBetweenPolls。这很不寻常,大多数人都希望尽快处理事情。使用它的一个原因是,如果您对调用的某些服务有速率限制,这意味着您必须减慢消费者的速度。
    • 谢谢@Gary。设置理想时间并让消费者慢下来的原因是,源事件在 8 到 10 分钟的间隔时间内几乎是重复的事件。因此,我想在此持续时间开始时进行轮询,以便我的消费者将消除所有大多数重复项并在空闲时间限制内处理事件。所以我想问你在这种情况下你有什么建议?
    • 所以这就是原因,我想将此 max.poll.interval 增加到 10 分钟,以便我的应用程序将处理更多事件,例如 10000 条记录/消费者,理想的轮询间隔为 6 分钟,所以轮询间隔总是在 max.poll.interval 的 10 分钟以内,以避免重新平衡。这会是一个有效的配置吗?
    • 你可以在我描述的参数范围内做任何你想做的事情,以避免重新平衡。目前尚不清楚您为什么设置 idleBetweenPolls。
    猜你喜欢
    • 2020-03-13
    • 2022-10-05
    • 1970-01-01
    • 2017-06-19
    • 2014-04-29
    • 1970-01-01
    • 1970-01-01
    • 2021-10-19
    • 1970-01-01
    相关资源
    最近更新 更多