Kafka Streams中KTable输出全事件日志而非最新更新的问题
Kafka Streams Table输出全量事件而非仅最新值的问题
我有一个包含4条事件的Kafka Topic,Topic中的事件顺序为:
- CCC: {"stockVal":10,"times":1234}
- BBB: {"stockVal":10,"times":1234}
- AAA: {"stockVal":10,"times":1234}
- AAA: {"stockVal":20,"times":1240}
我编写的拓扑代码本应仅打印每只股票的最新更新,但实际却输出了完整的事件流。代码如下:
val builder = new streams.StreamsBuilder val tableData = builder.table[String, StockData](inputTopic) tableData.toStream().print(Printed.toSysOut[String, StockData].withLabel("table-form")) builder.build()
实际输出日志:
[table-form]: CCC, {"stockVal":10,"times":1234} [table-form]: BBB, {"stockVal":10,"times":1234} [table-form]: AAA, {"stockVal":10,"times":1234} [table-form]: AAA, {"stockVal":20,"times":1240}
为何股票AAA的记录会被打印两次?启动应用时Topic中已有所有消息,按理解应仅输出AAA的最后一条值,恳请解答!
原因分析
Kafka Streams的KTable本质是基于状态存储的增量更新流,问题核心在于对KTable.toStream()的行为理解偏差:
toStream()输出状态变更记录:每处理一条输入Topic的消息,KTable都会更新对应key的状态存储,toStream()会把每一次状态变更的结果输出。对于AAA的第一条消息,状态从无到有,会输出这条记录;处理第二条AAA消息时,状态被更新为新值,因此会再次输出更新后的记录。- 从头消费触发全量更新:启动应用时Topic已有所有消息,且消费者组为首次消费该Topic(或配置了
auto.offset.reset=earliest),Kafka Streams会从头消费所有历史消息,每一条消息都会触发一次状态变更,进而被toStream()打印。
解决方案
如果需要在应用启动后仅输出每个key的最终最新值,可以监听Kafka Streams的运行状态,待状态存储完全初始化后直接读取快照:
val topology = builder.build() val streams = new KafkaStreams(topology, streamsConfig) // 监听应用状态,进入RUNNING状态后读取最终状态 streams.setStateListener((newState, oldState) => { if (newState == KafkaStreams.State.RUNNING) { // 获取KTable对应的状态存储 val keyValueStore = streams.store( tableData.queryableStoreName, QueryableStoreTypes.keyValueStore[String, StockData]() ) // 遍历所有key的最新值 val entries = keyValueStore.all() while (entries.hasNext) { val entry = entries.next() println(s"[final-latest]: ${entry.key}, ${entry.value}") } entries.close() } }) streams.start()
如果仅需实时处理启动后新产生的消息并输出最新值,可配置auto.offset.reset=latest,让应用从Topic末尾开始消费,这样每条输出的记录都是对应key的最新状态。
内容的提问来源于stack exchange,提问作者user19104652
相关产品推荐
相关产品推荐

