Kafka Streams中RocksDBTimestampedStore频繁日志问题咨询
问题描述
我们环境中有两个应用实例,拓扑定义如下:
KStream<String, ObjectMessage> stream = kStreamBuilder.stream(inputTopic); stream.mapValues(new ProtobufObjectConverter()) .groupByKey() .windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMillis(100))) .aggregate(AggregatedObject::new, new ObjectAggregator(), buildStateStore(storeName)) .suppress(Suppressed.untilWindowCloses(Suppressed.BufferConfig.unbounded().withMaxRecords(config.suppressionBufferSize()))) .mapValues(new AggregatedObjectProtobufConverter()) .toStream((key, value) -> key.key()) .to(outputTopic); private Materialized<String, AggregatedObject, WindowStore<Bytes, byte[]>> buildStateStore(String storeName) { return Materialized.<String, AggregatedObject, WindowStore<Bytes, byte[]>>as(storeName) .withKeySerde(Serdes.String()) .withValueSerde(new JsonSerde<>(AggregatedObject.class)); }
通过for循环为多个输入主题创建上述拓扑,单个应用实例包含多个拓扑,每个拓扑的状态存储命名格式为KSTREAM-AGGREGATE-%s-STATE-STORE-0000000001。
此前未配置state-dir目录,基于K8S有状态集部署,状态无法在重启后持久化,应用需重建状态。日志中频繁出现如下内容(仅后缀时间戳变化):
INFO 1 --- [-StreamThread-1] o.a.k.s.s.i.RocksDBTimestampedStore Opening store KSTREAM-AGGREGATE-my.topic.name-STATE-STORE-0000000001.1675576920000 in regular mode
但日志中的时间戳(如1675576920000)为1天前甚至3天前的时间。今日配置state-dir后,该日志仍持续打印。请问是只需等待处理完成,还是配置存在问题?
另外,为何RocksDBTimestampedStore会频繁打印日志?且日志中的时间戳与窗口操作定义的100ms不符?
问题解答
配置state-dir后日志仍打印的处理方式
配置state-dir后无需调整配置,只需等待旧状态处理完毕即可:
- 之前未开启状态持久化时,每次重启都会从头消费消息重建状态,过程中生成了大量对应历史窗口的状态存储文件(日志中的旧时间戳就是这些窗口的时间标记)。
- 配置
state-dir后,Kafka Streams会先加载已存在的状态文件,同时后台异步清理过期窗口状态。当前的日志正是加载这些旧状态文件的过程,等所有旧状态文件完成加载和清理,这类日志就会停止。 - 若要确认配置有效性,可检查K8S有状态集的
state-dir挂载路径权限(容器需具备读写权限)以及PVC存储是否正常。
RocksDBTimestampedStore频繁打日志及时间戳不符的原因
- 频繁打日志的原因:Kafka Streams的窗口状态存储按窗口结束时间分目录存储,每个窗口对应一个独立的RocksDB实例。重建状态时需要逐个加载历史窗口的存储文件,每个文件加载都会触发这条日志,因此看起来频繁。
- 时间戳与窗口大小不符的原因:日志中的时间戳是历史窗口的结束时间戳,而非当前窗口的时间。你定义的100ms是窗口大小,但重建状态时消费的是几天前的历史消息,这些消息会被分配到当时对应的窗口,加载该窗口的状态文件时就会打印那个时间戳,和当前窗口大小无关。
内容的提问来源于stack exchange,提问作者ashur
相关产品推荐
相关产品推荐

