在Apache Flux YAML拓扑中配置Storm-Kafka获取Kafka最新消息
解决Apache Flux Kafka Spout仅拉取最新消息的配置问题
我来帮你搞定这个问题!你遇到的报错是因为startOffsetTime用了ref引用,但LatestTime是KafkaConfig类里的静态常量,不是Flux容器中注册的bean,所以找不到导致赋值为null,触发了类型错误。
正确的配置方式是直接给startOffsetTime设置对应的常量数值,Storm的KafkaConfig里定义了:
LatestTime = -1L(从最新消息开始消费)EarliestTime = -2L(从最早消息开始消费)
结合你现有的配置,只需要把startOffsetTime的ref改成value并设为-1即可,同时你已经设置了ignoreZkOffsets: true,这个配置是必须的——它告诉Spout不要从ZooKeeper读取历史偏移量,而是使用我们指定的起始偏移时间。
修改后的完整spoutConfig配置如下:
- id: "stringScheme" className: "org.apache.storm.kafka.StringScheme" - id: "stringMultiScheme" className: "org.apache.storm.spout.SchemeAsMultiScheme" constructorArgs: - ref: "stringScheme" - id: "zkHosts" className: "org.apache.storm.kafka.ZkHosts" constructorArgs: - "172.25.33.191:2181" - id: "spoutConfig" className: "org.apache.storm.kafka.SpoutConfig" constructorArgs: - ref: "zkHosts" - "blockdata" - "" - "myId" properties: - name: "scheme" ref: "stringMultiScheme" - name: "ignoreZkOffsets" value: true - name: "startOffsetTime" value: -1 # 对应KafkaConfig.LatestTime,拉取最新消息
这样配置后,Storm Spout就会忽略ZooKeeper中的偏移量记录,直接从Kafka主题的最新位置开始消费消息了。
内容的提问来源于stack exchange,提问作者Obiii
相关产品推荐
相关产品推荐

