【问题标题】:Order of receiving messages if Kafka consumer subscribes to multiple topicsKafka 消费者订阅多个主题时接收消息的顺序
【发布时间】:2020-09-27 01:22:22
【问题描述】:

我有一个调查多个主题的消费者。对于这个问题,我限制了每个主题一个分区。假设当消费者开始轮询时,每个主题都有一些数据。 读取顺序是什么?

是循环赛吗?它是从第一个在下一个之前读取的吗?我使用consumer.poll(N) 进行投票。

【问题讨论】:

  • 它在给定的主题分区内是线性的,但我认为每个轮询循环都会轮询其他主题...当然您可以通过打印记录元数据轻松测试这一点?
  • 是的,第一个。关于第二个 - 是的,我可以这样做,但这可能是间接的。 Kafka 规范是否在任何地方都这么说(找不到)?
  • 如果不对更多的消息进行测试,将很难找到任何关于它的信息。常识可能会说它是循环/随机的,否则更活跃的主题会为自己吸引消费者,并且永远不会读取来自其他主题的消息。

标签: apache-kafka kafka-consumer-api


【解决方案1】:

排序相当复杂。以下是 Kafka 2.6 的工作原理:

  • 当您将主题分区分配给消费者时,这些分区将保存在哈希表中,因此顺序将是稳定的,但不一定是您使用的那个
  • 当您调用Consumer.poll(N) 时,它会返回所有排队的消息,但最多返回max.poll.records(见下文)
  • 当没有任何内容入队时,您分配的所有主题分区将按该主题分区的领导者所在的 Kafka 节点进行分区
  • 这些列表中的每一个都在提取请求中发送到每个相应的节点
  • 每个节点最多返回fetch.max.bytes(或至少一条消息,如果可用)
  • 节点将使用来自请求分区的消息填充这些字节,始终从第一个开始
  • 如果当前分区中没有消息了,但仍有字节要填充,则移动到下一个分区,直到没有消息或缓冲区已满
  • 节点还可以决定停止使用当前分区并继续使用下一个分区,即使当前分区中仍有可用消息
  • 客户端/消费者收到缓冲区后,会将其拆分为CompletedFetches,其中一个CompletedFetch正好包含缓冲区中一个主题分区的所有消息
  • 那些CompletedFetches 已入队(它们可能包含 0 条消息或 1000 条或更多消息)。每个请求的主题分区都会有一个CompletedFetch
  • 由于对节点的所有请求都是并行运行的,但只有一个队列,CompletedFetches/topic 分区可能会在最终结果中混淆,而不是原始分配顺序
  • 入队的 CompletedFetches 在逻辑上被扁平化为一个大队列
  • Consumer.poll(N) 最多会从扁平的大队列中读取和出列max.poll.records
  • 在记录返回给poll的调用者之前,开始了对所有节点的另一个fetch请求,但这一次,所有已经在扁平化队列中的主题分区都被排除在外
  • 这适用于所有未来的poll 电话

在实践中,这意味着您不会挨饿,但您可能会收到来自一个主题的大量消息,然后才能收到大量用于下一个主题的消息。

在消息大小为 10 字节的测试中,从一个主题读取了大约 58000 条消息,然后从下一个主题读取了大致相同的数量。 所有主题都预先填充了 100 万条消息。

因此你会有一种批量循环。

【讨论】:

    【解决方案2】:

    没有排序,因为底层协议允许在一个请求中发送多个分区的请求。

    当您调用 consumer.poll(N) 时,客户端确实将 FetchRequest 对象发送到托管分区领导者的代理(请参阅org.apache.kafka.clients.consumer.internals.Fetcher.createFetchRequests()) - 每个节点只有一个请求,而不是每个分区。

    重要的是客户端可以为多个分区发送一个 FetchRequest(参见protocol spec)。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2017-05-16
      • 1970-01-01
      • 2017-07-31
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-09-09
      • 2019-01-10
      相关资源
      最近更新 更多