【问题标题】:How does one Kafka consumer read from more than one partition?一个 Kafka 消费者如何从多个分区中读取数据?
【发布时间】:2022-01-14 07:11:50
【问题描述】:

我想知道一个消费者如何从多个分区消费,具体来说,从不同分区读取消息的顺序是什么?

我看过源代码(ConsumerFetcher),但我无法真正了解所有内容。

这是我认为会发生的事情:

按顺序读取分区。即:一个分区中的所有消息都将被读取,然后再继续下一个分区。如果我们到达max.poll.records 而没有消耗整个分区,那么下一次提取将继续读取当前分区,直到它耗尽,然后再继续下一个。

我尝试将max.poll.records 设置为一个相对较低的数字,然后看看会发生什么。 如果我向一个主题发送消息然后启动一个消费者,那么在继续到下一个分区之前,所有消息都会从一个分区读取,即使该分区中的消息数高于max.poll.records

然后我尝试通过连续向该分区发送消息(使用 JMeter)来查看是否可以将消费者“锁定”在一个分区中。但我做不到:来自其他分区的消息也在被读取。

【问题讨论】:

标签: apache-kafka kafka-consumer-api


【解决方案1】:

消费者以贪婪的循环方式从其分配的分区轮询消息。 例如如果max.poll.records 设置为 100,并且分配了 2 个分区 A、B。消费者将尝试从 A 轮询 100 条消息。如果分区 A 没有 100 条可用消息,它将从分区 B 轮询剩下的要完成的 100 条消息。

虽然这样不太理想,但是这样就不会饿死分区了。

这也解释了为什么分区之间不能保证排序。

【讨论】:

    【解决方案2】:

    我已阅读 cmets 中链接的问题的答案中提到的KIP,我想我终于明白了消费者的工作方式。

    有两个主要的配置选项会影响数据的使用方式:

    • max.partition.fetch.bytes: 服务器将为给定分区返回的最大数据量

    • max.poll.records:消费者每次轮询时返回的最大记录数


    从每个分区获取的过程是贪婪的,并且以循环方式进行。贪婪意味着将从每个分区中检索尽可能多的记录;如果一个分区中的所有记录占用少于max.partition.fetch.bytes,则将全部取出;否则,只会获取max.partition.fetch.bytes

    现在,并非所有获取的记录都会在轮询调用中返回。只会返回max.poll.records

    将保留剩余的记录以供下一次轮询调用。 而且,如果保留的记录数少于max.poll.records,poll 方法会在返回前开始新一轮的抓取(pre-fetching)。这意味着,通常,消费者在获取新记录时正在处理记录。


    如果某些分区收到的消息比其他分区多得多,这可能会导致活动较少的分区长时间不被处理。

    这种方法的唯一缺点是,当分区各自的消息速率之间存在很大的不平衡时,它可能会导致某些分区长时间未使用。例如,假设 max messages 设置为 1 的消费者从分区 A 和 B 获取数据。如果返回的 fetch 包括来自 A 的 1000 条记录而没有来自 B 的记录,则消费者必须在获取之前处理来自 A 的所有 1000 条可用记录再次在分区 B 上。

    为了防止这种情况,我们可以减少max.partition.fetch.bytes

    【讨论】:

      猜你喜欢
      • 2022-06-13
      • 1970-01-01
      • 2022-11-19
      • 1970-01-01
      • 2015-08-19
      • 1970-01-01
      • 1970-01-01
      • 2021-11-28
      • 2016-06-04
      相关资源
      最近更新 更多