【问题标题】:Can I retrieve the latest available offset for a Kafka partition without retrieving all the messages?我可以在不检索所有消息的情况下检索 Kafka 分区的最新可用偏移量吗?
【发布时间】:2018-02-23 00:31:17
【问题描述】:

看最新(v0.10)Kafka Consumer documentation

“消费者的位置给出了下一条记录的偏移量。它将比消费者在该分区中看到的最高偏移量大一个。它会自动前进每个消费者接收数据调用poll(long)并接收消息的时间。"

有没有办法查询服务器端分区可用的最大偏移量,而不检索所有消息?

我试图实现的逻辑如下:

  1. 每秒查询主题中待处理消息的数量 (A)
  2. 如果 A > 阈值,则唤醒将继续检索所有消息并处理它们的处理器
  3. 否则什么都不做(睡眠 1)

动机是我需要做一些批处理,但我希望处理器只有在有足够的数据时才唤醒(并且我不想两次检索所有数据)。

【问题讨论】:

    标签: apache-kafka kafka-consumer-api


    【解决方案1】:

    您可以使用Consumer.seekToEnd() 方法,运行Consumer.poll(0) 使其生效但立即返回,然后Consumer.position() 查找所有订阅(或分配)主题分区的位置。这些将是所有分区的当前最终偏移量。这也将开始从代理获取这些偏移量的一些数据,但如果您随后返回不同的位置,则任何返回的数据都将被忽略。

    正如 serejja 所提到的,目前的替代方法是使用旧的简单消费者,尽管该过程相当复杂,因为您需要手动找到每个分区的领导者。

    【讨论】:

    • 谢谢。我想知道是否可以避免两次读取所有数据(在我上面描述的场景中)。例如,我可以将 max.partition.fetch.bytes 减小到一个非常小的值,以消除 poll(0) 检索实际数据的“副作用”吗?
    • 你不需要调用 poll()。 seekToEnd() 是一个异步调用,您可以使用 poll() 或 position() 强制完成。使用 seek...() 和 position() 不会读取任何消息,只会读取少量元数据
    • @ChrisGerken 如果您正在与消费者团体合作并且还没有任务,则不确定这是否可靠(我必须更仔细地研究代码,但看起来它会抛出IllegalArgumentException)。对于手动分配的主题/主题分区,它似乎可以正常工作。
    【解决方案2】:

    遗憾的是,我不明白这对于 0.10 消费者来说是如何实现的。

    但是,如果您有任何较低级别的 Kafka 客户端,这是可行的(抱歉,我不确定是否存在用于 JVM 的客户端,但有很多用于其他语言的客户端)。

    因此,如果您有时间和灵感来实现这一点,这就是要走的路 - 每个 FetchResponse(这是每个“给我消息”请求的响应)都包含一个名为 HighwaterMarkOffset 的字段,本质上是分区末尾的偏移量 (https://cwiki.apache.org/confluence/display/KAFKA/A+Guide+To+The+Kafka+Protocol#AGuideToTheKafkaProtocol-FetchResponse)。这里的技巧是发送一个FetchRequest,它会立即返回(例如,不会阻止等待),除了 HighwaterMarkOffset 什么都没有。

    为此,您的FetchRequest 应该有:

    1. MaxWaitTime 设置为 0,这意味着“如果无法获取至少 MinBytes 字节,则立即返回”。
    2. MinBytes 设置为 0,意思是“如果你给我一个空的回复我就可以了”。
    3. FetchOffset 在这种情况下无关紧要,如果我没记错的话,它甚至可能是一个无效的偏移量,但最好是一个有效的偏移量。
    4. MaxBytes 设置为 0,意思是“给我不超过 0 个字节的数据”,例如什么都没有。

    这样,这个请求将立即返回,没有数据,但仍将高水位偏移设置为适当的值。获得高水位线偏移量后,您可以将其与当前偏移量进行比较,并计算出您落后了多少。

    希望这会有所帮助。

    【讨论】:

    • 谢谢@serejja!这无疑为进一步探索提供了方向。知道如何使用内部Fetcher class 中的想法/代码来实现这一目标吗? listOffset 或内部 sendListOffsetRequest 方法看起来很有希望..
    • @AlexGlikson 我们面临一个与您类似的问题。我们使用反射来调用具有最新时间戳的内部 listOffset。
    【解决方案3】:

    您可以使用下面 API 中的 public OffsetAndMetadata committed(TopicPartition partition) 方法来获取最后提交的偏移量

    https://kafka.apache.org/0100/javadoc/index.html?org/apache/kafka/clients/consumer/KafkaConsumer.html

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2013-07-24
      • 2022-01-18
      • 2015-04-03
      • 2016-04-22
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2022-10-25
      相关资源
      最近更新 更多