Kafka Streams的num.stream.threads配置与输入topic分区数匹配问题咨询
问题1解答
num.stream.threads 是单个Kafka Streams实例内的并行处理线程数,单个线程同一时间最多承载1个流处理任务(StreamTask),而整个流应用的最大可并行任务数由流拓扑的输入分区规则决定:
- 如果你的拓扑仅消费单个输入topic,最大可并行任务数就等于该topic的分区数,此时单个实例的
num.stream.threads配置超过该topic分区数的部分会完全空闲,没有任务可分配,相当于浪费资源,这种场景下线程数确实不能超过单个输入topic的总分区数。 - 整个应用集群的总线程数(所有实例的
num.stream.threads之和)也不能超过最大可并行任务数,超出部分同样会空闲。
问题2解答
多输入topic场景的分配规则取决于你的流拓扑逻辑,你举的两个输入topic各12分区、num.stream.threads设为12的场景,分两种情况:
- 若拓扑中两个输入topic无关联操作(如分别过滤后统一写入输出topic,无join、分区对齐类聚合):总最大可并行任务数为两个topic分区数之和即24,12个线程每个会分配到2个任务,每个任务对应一个topic的单个分区,所以每个线程会消费2个分区,分别来自两个不同的topic。
- 若拓扑中两个输入topic存在join、窗口聚合等需要分区对齐的操作:Kafka Streams要求这类关联操作的输入topic分区数必须一致,此时总最大可并行任务数等于单个topic的分区数即12,每个任务会负责两个topic的同编号分区(如任务1负责topicA的分区1 + topicB的分区1),12个线程每个会分配到1个任务,每个线程同样会消费2个分区,分别来自两个topic的同号分区。
调优参考
配置num.stream.threads时优先计算流拓扑的总最大可并行任务数,单个实例的线程数不要超过该数值,多实例部署时所有实例的线程数总和也不要超过该数值,避免资源浪费。
内容的提问来源于stack exchange,提问作者AttitudeL
相关产品推荐
相关产品推荐

