KafkaStreams#store方法是否线程安全?跨线程调用场景咨询
Kafka Streams跨线程调用store方法的线程安全性问题
问题背景
我的Kafka Streams应用中存在两个线程:
- Thread A:创建包含状态存储等内容的Topology对象,随后调用KafkaStreams类的构造方法及
start()方法。 - Thread B:持有Thread A创建的KafkaStreams对象引用,定期调用
KafkaStreams#store方法获取ReadOnlyWindowStore实例,读取存储数据用于监控。
我担忧此实现的线程安全性:
ReadOnlyWindowStore的官方文档说明“实现应保证线程安全,因为会存在并发读写场景”,因此对此类的线程安全问题并不担心。- 但对于
KafkaStreams#store方法,我不确定跨线程调用是否安全——该方法涉及非线程安全的HashMap(对应Kafka源码3.6.1版本QueryableStoreProvider.java第39行),而KafkaStreams#removeStreamThread()方法会修改此HashMap。
我的问题:
- 跨线程调用
KafkaStreams#store是否可行? - 是否应在创建KafkaStreams的线程中调用
store()并仅共享ReadOnlyWindowStore实例? - KafkaStreams类整体是否设计为线程安全?
解答
- 跨线程调用
KafkaStreams#store完全可行:虽然底层依赖的HashMap本身非线程安全,但Kafka Streams内部对这个HashMap的访问做了同步控制,官方设计上允许外部线程调用store方法。实际生产里很多监控场景都是这么实现的,只要KafkaStreams实例处于正常运行状态(比如RUNNING或REBALANCING状态),调用store不会有线程安全问题。 - 不需要局限在创建线程调用
store:ReadOnlyWindowStore本身是线程安全的,但如果提前在Thread A获取实例并共享,当发生Rebalance这类场景时,旧实例可能无法访问到最新的存储分片。反而定期在Thread B调用store获取最新实例,才能保证访问到的是当前有效的状态存储,这才是更稳妥的做法。 - KafkaStreams类整体是线程安全的:官方明确标注KafkaStreams的所有公共方法都具备线程安全性,包括
start()、store()、close()等。内部的状态变更操作(比如removeStreamThread()修改HashMap)都有对应的同步机制,不会引发并发访问的异常。
内容的提问来源于stack exchange,提问作者Kohei Nozaki
相关产品推荐
相关产品推荐

