【发布时间】:2019-06-05 05:03:31
【问题描述】:
我目前正在开发一个部分依赖于 Apache Kafka(2.2.0 版)的应用程序。我必须做的一件事是跟踪其他消费者提交其当前偏移量的内容(更重要的是何时)。据我所知,仅使用 Java 客户端,无法获取已提交偏移量的相关时间戳,因为 AdminClient 的 listConsumerGroupOffsets 方法最终会导致 OffsetAndMetadata 对象,其中不包括时间戳。因此,我只是开始阅读来自__consumer_offsets 主题的消息。如果有更好的方法,请告诉我。
现在,如果一个人直接读取__consumer_offsets 中的消息,那么一个人一下子就有了两个时间戳。一个是附加到实际提交消息的时间戳,另一个是commit_timestamp,它是消息内容的一部分。我的第一个想法是其中一个可能由代理设置,另一个可能由提交它的客户端设置(另外,如果您查看 ZooKeeper 中的/config/topics/__consumer_offsets,它没有指定LogAppendTime 消息时间戳,因此可以假设它只使用默认值)。唉,手动移动系统时间的快速实验表明,两者实际上都是由代理设置的。更重要的是,他们并不总是同意(消息的时间戳有时略早于commit_timestamp)。我试图深入研究 Kafka 代码以准确了解发生了什么,但它相当复杂,而且我对它还不够熟悉,无法快速掌握。所以这是我的问题:
- 为什么
__consumer_offsets中的消息时间戳会自动LogAppendTime,即使没有明确指定?只是用于发送提交消息的生产者将时间戳留空吗? - 为什么消息时间戳和消息中包含的
commit_timestamp不一致?我似乎记得曾经在某处读到过,以前可以显式设置commit_timestamp,从而手动控制提交偏移量的保留。 - 更重要的是:是否有任何理由使用其中一种?例如,如果仍然可以手动设置
commit_timestamp,则使用附加到消息的时间戳会更有意义。
我知道这是一个非常具体的问题,对大多数人来说可能并不重要。但直到现在,我总是能够通过使用 Google 并查看 Kafka 的源代码来了解后台发生的事情;然而,这个让我有点难过。因此,非常感谢任何见解。
【问题讨论】:
-
OffsetAndMetadata 的元数据可以包含您想要的任何内容。事实上,这是 Confluent Replicator 用于确保灾难场景中的消费者故障转移的一部分
-
@cricket_007 啊,好点子。但是,在元数据中,我只能包含客户端时间戳。我在这里感兴趣的(以及我提到的两个时间戳提供的)是代理端时间戳。
-
我没有花太多时间查看消费者偏移消息的生成,但我个人并不知道能够显式设置时间。不过,消息仍然是通过 ProducerRecord 形成的,并且是由客户端生成的(我相信是消费者组协调员),因此它不完全是“代理时间戳”
-
@cricket_007 非常感谢您的见解。首先,消费者组协调员不是经纪人(与组长相反)吗?消息当然是由客户端创建的,但时间戳似乎是代理生成的。为了测试这一点,我运行了一个系统时间偏移 5 秒的 Kafka 代理,这也是我在获得的时间戳中看到的。手动设置时间戳似乎是 Kafka 协议的一部分,如 here 所示。
-
嗯。所以你说你使用的是 Kafka 2.2,但那些文档说偏移提交请求的时间戳自 0.9 以来已被删除......我没看错吗?
标签: apache-kafka kafka-consumer-api