未设置墓碑标记时KTable记录何时过期?会继承源Topic保留策略吗
Kafka Streams 相关问题解答
1. KStream 是否遵守源Topic的2天过期规则?
KStream 是源Topic消息流的无状态抽象,本身不会持久化任何消息数据,因此不存在是否遵守过期规则的说法:
- 运行中的
KStream只会消费当前源Topic中还未被清理的消息,已经被retention.ms规则清理掉的消息不会被消费到 KStream不会存储已经处理完成的消息,因此不存在源Topic消息清理后需要从KStream中移除数据的逻辑
2. KTable 是否遵守源Topic的2天过期规则?
不会。KTable是有状态抽象,它的状态数据存储在Kafka Streams应用独立的状态存储(默认是RocksDB+内存)和对应的changelog Topic中,和源Topic T的清理规则完全隔离:
- 源Topic
T中消息因retention.ms到期被清理的动作,完全不会影响KTable已经生成的状态数据,旧数据不会自动从KTable中移除 - 非窗口型的
KTable状态默认永久保留,除非收到对应key的墓碑消息,或者主动配置了状态保留策略(仅窗口型KTable支持TTL自动过期配置)
3. 是否需要生成墓碑标记完成KTable数据清理?
是的,如果你需要自动删除KTable中指定key的状态数据,标准实现方式就是往源Topic T 发送一条同key、value为null的墓碑消息:
KStream消费到这条墓碑消息后,会同步传递给下游KTable,KTable会自动将该key从状态存储中删除,同时也会在changelog Topic中写入对应的墓碑标记,后续应用重启恢复状态时也不会再加载该key的旧数据- 如果不发送墓碑消息,哪怕源Topic中所有旧消息都被清理,
KTable中的对应状态会一直存在,只能通过手动重置Kafka Streams应用的状态存储来清理,属于专属运维操作,没有自动触发逻辑。
内容的提问来源于stack exchange,提问作者simonalexander2005
相关产品推荐
相关产品推荐

