如何使用Erlang Brod组订阅者获取Kafka各分区的提交偏移量与Kafka偏移量?
Brod v3.16.2组订阅者获取分区偏移量问题解答
在Brod v3.16.2版本中,使用Group Subscriber时完全可以获取每个Kafka分区的已提交偏移量(committed offsets)和Kafka端的偏移量,具体实现方式如下:
获取已提交偏移量:
调用Brod提供的brod_group_subscriber:get_committed_offsets(SubscriberRef)函数,传入组订阅者的引用(SubscriberRef),就能得到所有订阅分区对应的已提交偏移量映射,数据结构通常为#{Topic => #{Partition => Offset}},直接对应每个主题分区的已提交偏移量值。获取Kafka分区的偏移量(通常指分区最新高水位偏移量):
使用brod:fetch_watermarks(Client, Topic, Partition)函数,传入Brod客户端实例、目标主题和分区编号,该函数会返回该分区的低水位(low watermark)和高水位(high watermark)偏移量,其中高水位就是当前Kafka集群中该分区的最新偏移量。如果需要批量获取多个分区的偏移量,遍历分区列表调用此函数即可。
这些都是Brod v3.16.2原生支持的API,无需额外扩展就能直接使用。
内容的提问来源于stack exchange,提问作者ditas
相关产品推荐
相关产品推荐

