【发布时间】:2017-10-15 15:22:46
【问题描述】:
我正在尝试让我的风暴拓扑在重新启动时从某个 kafka 偏移量中读取。
如果我理解正确的话,我可以通过ignoreZkOffsets 并设置startOffsetTime 来做到这一点,但到目前为止还没有奏效。
我尝试将 startOffsetTime 设置为 System.currentTimeMillis() - 60000L 以从一分钟前开始,并将其设置为当前偏移量。
【问题讨论】:
我正在尝试让我的风暴拓扑在重新启动时从某个 kafka 偏移量中读取。
如果我理解正确的话,我可以通过ignoreZkOffsets 并设置startOffsetTime 来做到这一点,但到目前为止还没有奏效。
我尝试将 startOffsetTime 设置为 System.currentTimeMillis() - 60000L 以从一分钟前开始,并将其设置为当前偏移量。
【问题讨论】:
来自kafka FAQ页面“Kafka允许按时间查询消息的偏移量,并且以段粒度进行。timestamp参数是unix时间戳,按时间戳查询偏移量返回消息的最新可能偏移量,即不迟于给定的时间戳附加。时间戳有 2 个特殊值 - 最新(从主题结束)和最早(从主题开始)。对于 unix 时间戳的任何其他值,Kafka 将获得不迟于给定时间戳创建的日志段的起始偏移量。因此,由于偏移量请求仅以段粒度提供服务,因此对于较大的段大小,偏移量获取请求返回的结果不太准确。" https://cwiki.apache.org/confluence/display/KAFKA/FAQ#FAQ-HowdoIaccuratelygetoffsetsofmessagesforacertaintimestampusingOffsetRequest?
如果您知道应用程序应该开始使用消息的偏移量,那么在 zookeeper 中设置它并将 ignoreZkOffsets 设置为 true。
仅供参考:zookeeper 的节点路径将是您在 spout 配置期间为 zkRoot 属性指定的值。
希望对你有所帮助。
【讨论】:
您对 ignoreZkOffsets 的理解部分正确,将此选项设置为 true 将缩短存储在 zookeeper 中的偏移量,但 startOffsetTime 不是任意的 Unix 时间戳。默认startOffsetTime的初始化如下:
public long startOffsetTime = kafka.api.OffsetRequest.EarliestTime();
Kafka api 只提供了EarliestTime 和LatestTime 2 种方法来设置初始偏移量,也就是说这种方法不起作用。
如果知道offset值,可以尝试修改storm-kafka在zookeeper中存储的offset值。这个值存储在${ZKRoot}/${ClientId}/${KafkaPartitionId}的ZKPath中,其中ClientId是你在SpoutConfig中指定的,如果你只有一个分区,KafkaPartitionId通常为0。
一旦你找到这个值,根据需要设置这个值并重新启动你的拓扑,它将从这个偏移量开始读取。如果此 ZKPath 不存在,您可以手动创建此路径。
此解决方案的一个缺陷是您应该了解自己的 clientId,这意味着您不能按照 Storm-starter 演示中的建议使用随机 UUID 作为您的 clientId。
【讨论】: