【发布时间】:2021-11-28 14:30:01
【问题描述】:
给定this answer:
所有主题的最大分区数决定了任务的数量。
和AbstractTask的代码如下:
final Set<TopicPartition> inputPartitions
我想知道什么时候(如果有的话)一个任务可以分配多个分区?
【问题讨论】:
给定this answer:
所有主题的最大分区数决定了任务的数量。
和AbstractTask的代码如下:
final Set<TopicPartition> inputPartitions
我想知道什么时候(如果有的话)一个任务可以分配多个分区?
【问题讨论】:
例如,如果您执行join()、merge() 或copartition(),则任务将具有多个输入分区。此外,如果您通过模式订阅一次阅读多个主题。
它是正交的
所有主题的最大分区数决定了任务的数量
引用是关于创建任务的数量,与每个任务的分区数量无关。
假设您有两个输入主题 A 有 2 个分区(A-0 和 A-1)和 B 只有一个分区(B-0)。你的程序是:
KStream a = builder.stream("A",...);
KStream b = builder.stream("B",...);
a.merge(b);
你的程序是合乎逻辑的:
topic-A ---+
+---> merge()
topic-B ---+
对于这种情况,你会得到两个任务:
Task 0_0:
A-0 ---+
+--- merge() -->
B-0 ---+
Task 0_1:
A-1 ---+
+--- merge() -->
注意第二个任务0_1只有一个输入分区,因为主题B只有1个分区。
任务基本上是您的(逻辑)程序的副本(物理实例化),它处理了名称为分区号的所有分区。因为 topic-A 有两个分区,所以需要创建两个任务。
【讨论】:
Processor 来演示这个每个任务多分区的案例?可以从 Kafka Streams 的哪一部分资源中学习?