【发布时间】:2017-12-12 12:01:50
【问题描述】:
我正在使用 Confluent kafka C# 客户端。如何获取此主题中消耗的最新偏移量?
【问题讨论】:
标签: c# apache-kafka kafka-consumer-api confluent-platform
我正在使用 Confluent kafka C# 客户端。如何获取此主题中消耗的最新偏移量?
【问题讨论】:
标签: c# apache-kafka kafka-consumer-api confluent-platform
我可以从生产者那里读取主题偏移量(high 和 low),而不是从消费者那里检索偏移信息(我不想先消费消息):
var partitionOffset = _producer.QueryWatermarkOffsets(new TopicPartition("myTopic", myPartition), TimeSpan.FromSeconds(10));
【讨论】:
当您收到一条消息时,它应该包括主题、分区和来自它的来源的偏移量(除了键和值)。
来自example here:
consumer.OnMessage += (_, msg)
=> Console.WriteLine($"Topic: {msg.Topic} Partition: {msg.Partition} " +
$"Offset: {msg.Offset} {msg.Value}");
当它到达每个主题分区的末尾时,您还会收到一个事件
consumer.OnPartitionEOF += (_, end)
=> Console.WriteLine($"Reached end of topic {end.Topic} partition {end.Partition}" +
$" , next message will be at offset {end.Offset}");
【讨论】:
除了上一个答案,你可以使用
List<TopicPartitionOffsetError> Position(IEnumerable<TopicPartition> partitions)
它将返回从 librdkafka 轮询给定主题/分区的最后偏移量
您有一个类似的Committed 方法,用于获取消费者最新提交的偏移量
还可以查询最新的已知偏移量
WatermarkOffsets QueryWatermarkOffsets(TopicPartition topicPartition, TimeSpan timeout)
它将向 kafka 集群发送请求。呼叫被阻塞,设置适当的超时。目前,您不能一次在多个分区上发送请求。 您可以使用它来获取最后已知的偏移量,或者计算滞后
还有
WatermarkOffsets GetWatermarkOffsets(TopicPartition topicPartition)
这将查询 librdkafka 中的内部状态,并可能返回 INVALID_OFFSET (-1001)。您可以使用它来检测由于处理数据而导致的一些滞后。 (这个方法的位置和结果的区别)
【讨论】: