【问题标题】:Multiple Consumer threads using Alpakka connector使用 Alpakka 连接器的多个消费者线程
【发布时间】: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


【解决方案1】:

是什么让您相信您需要更多线程?更多时候,您希望跨多个流共享单个 Kafka 消费者客户端实例。

您不应将多个Consumer.committableSources 中的元素合并到一个流中,它不适用于批量提交。

多次运行相同的流设置会满足您的需求吗?

【讨论】:

  • 如果我不需要批量提交,将多个 committableSource 中的元素合并到一个流中会提供与 Consumer.committablePartitionedSource 提供的相同性能吗? committablePartitionedSource 将为每个分区创建源,因此它提供批量提交?
  • 不,committablePartitionedSource 在内部使用单个 Kafka 消费者。提交与它一起工作。
猜你喜欢
  • 2020-07-15
  • 1970-01-01
  • 1970-01-01
  • 2019-09-20
  • 1970-01-01
  • 1970-01-01
  • 2023-03-23
  • 2017-02-01
  • 1970-01-01
相关资源
最近更新 更多