【问题标题】:Spring Kafka don't respect max.poll.records with strange behaviorSpring Kafka 不尊重 max.poll.records 的奇怪行为
【发布时间】:2020-06-30 17:33:34
【问题描述】:

好吧,我正在尝试以下场景:

  1. 在 application.properties 中将 max.poll.records 设置为 50。
  2. 在 application.properties 中将 enable-auto-commit=false 和 ack-mode 设置为手动。
  3. 在我的方法中添加了@KafkaListener,但不提交任何消息,只阅读、记录但不发送ACK。

实际上,在我的 Kafka 主题中,我有 500 条消息要消耗,所以我期待以下行为:

  1. Spring Kafka poll() 50 条消息(偏移量 0 到 50)。
  2. 正如我所说,我没有提交任何内容,只是记录了 50 条消息。
  3. 在下一次 Spring Kafka poll() 调用中,获取与步骤 1 相同的 50 条消息(偏移量 0 到 50)。据我了解,Spring Kafka 应该继续此循环(步骤 1-3),始终阅读相同的消息。

但是会发生以下情况:

  1. Spring Kafka poll() 50 条消息(偏移量 0 到 50)。
  2. 正如我所说,我没有提交任何内容,只是记录了 50 条消息。
  3. 在下一次 Spring Kafka poll() 调用中,获取与步骤 1 不同的 NEXT 50 条消息(偏移量 50 到 100)。

Spring Kafka 以 50 条消息为单位读取 500 条消息,但不提交任何内容。如果我关闭应用程序并重新启动,则再次收到 500 条消息。

所以,我的疑惑:

  1. 如果我将 max.poll.recors 配置为 50,如果我没有提交任何内容,spring Kafka 如何获取接下来的 50 条记录?我知道 poll() 方法应该返回相同的记录。
  2. Spring Kafka 有缓存吗?如果是,如果我在缓存中获得 100 万条记录而没有提交,这可能是个问题。

【问题讨论】:

    标签: spring apache-kafka spring-kafka


    【解决方案1】:

    您的第一个问题:

    如果我将 max.poll.recors 配置为 50,spring Kafka 如何获取 如果我没有提交任何内容,接下来的 50 条记录?我了解 poll() 方法应该返回相同的记录。

    首先,为了确保你没有提交任何东西,你必须确保你理解以下3个参数,我相信你已经理解了。

    • ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG,将其设置为 false(这也是推荐的默认值)。如果设置为 false,请注意 auto.commit.interval.ms 变得无关紧要。查看this 文档:

    因为监听器容器有它自己的提交机制 偏移量,它更喜欢 Kafka ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG 是假的。从 2.3 版开始,它无条件地将其设置为 除非在消费者工厂或 容器的消费者属性覆盖。

    • factory.getContainerProperties().setAckMode(AckMode.MANUAL);您有责任承认。 (在使用事务时忽略)和ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG 不能是true

    • factory.getContainerProperties().setSyncCommits(true/false); 设置当容器负责提交时是否调用consumer.commitSync()commitAsync()。 默认为true。这负责与 Kafka 同步,没有别的,如果设置为 true,则该调用将阻塞,直到 Kafka 响应。

    其次,消费者 poll() 不会返回相同的记录。对于当前正在运行的消费者,它使用一些内部索引跟踪其在内存中的偏移量,我们不必关心提交偏移量。另请参阅@GaryRussell 的解释here

    简而言之,他解释说:

    一旦轮询返回了记录(并且偏移量不 已提交),除非您重新启动 消费者或对消费者执行 seek() 操作以重置 偏移到未处理的。


    您的第二个问题:

    Spring Kafka 有缓存吗?如果是,这可能是一个问题,如果我 无需提交即可在缓存中获取 100 万条记录。

    没有“缓存”,都是关于偏移量和提交的,解释同上。



    现在要实现您想要做的事情,您可以考虑在获取前 50 条记录后做 2 件事,即为下一个 poll():

    • 或者,以编程方式重新启动容器
    • 或致电consumer.seek(partition, offset);


    奖励:
    无论您选择什么配置,您都可以通过查看此输出的LAG 列来查看结果

    kafka-consumer-groups.bat --bootstrap-server localhost:9091 --describe --group your_group_name
    

    【讨论】:

      【解决方案2】:

      消费者不提交偏移量只会在以下情况下产生影响:

      • 您的消费者在读取 200 条消息后崩溃,当您重新启动它时,它会从 0 重新开始。
      • 您的消费者不再被分配分区。

      所以在一个完美的世界里,你根本不需要提交,它会消耗所有的消息,因为消费者首先要求 1-50,然后是 51-100。

      但是如果消费者崩溃了,没有人知道消费者读取的偏移量是多少。如果消费者已经提交了偏移量,当它重新启动时,它可以检查偏移量主题以查看崩溃的消费者离开的位置并从那里开始。

      max.poll.records 定义了一次要获取多少条记录,但它没有定义要获取哪些记录。

      【讨论】:

      • 这种我不需要提交的行为(在完美的世界中,认为我的应用程序永远不会关闭)是 Apache Kafka 或 Spring Kafka 功能?如果我使用另一个库,使用另一种语言,行为是否相同?
      • Apache kafka 功能。是的。如果使用普通的 Kafka 消费者或其他库,这将是相同的
      • 感谢您的回复。
      猜你喜欢
      • 2021-09-08
      • 1970-01-01
      • 2012-05-31
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-03-11
      • 2021-09-13
      相关资源
      最近更新 更多