【发布时间】:2019-09-19 22:16:08
【问题描述】:
我正在使用 Alpakka kafka 连接器来使用来自 kafka 的数据包。我使用 Consumer 作为 CommittableSource。我想在一台机器上创建多个消费者线程并将它们用作单一来源。我怎样才能做到这一点?
目前,我使用 Consumer.CommittableSource 创建了多个源,并使用“合并”功能将所有源合并为一个源。但我不确定这是否是正确的方法,因为我没有创建线程。
请在下面找到我目前使用的源代码:
public Source<ConsumerMessage.CommittableMessage<String, String>, Consumer.Control> source() {
Source finalSource = Source.empty();
for (int index = 0; index < consumerConfig.getNoOfConsumers(); index++) {
finalSource = finalSource.merge(Consumer.committableSource(consumerSettings, subscription));
}
return finalSource;
}
【问题讨论】:
-
(没有任何经验)我认为对单个主题的并发处理只能通过为该主题指定多个分区来实现。每个分区都可以有自己的消费者(参见stackoverflow.com/questions/38024514/…)
-
@Conffusion :我已经知道每个分区都可以有自己的消费者。我在问如何通过 alpakka 连接器消耗多个数据包。
-
我只是分享我的一点理论知识,希望对您有所帮助。我将无法进一步帮助您。祝你好运。
标签: java apache-kafka akka-stream alpakka