【问题标题】:Kafka consumer reading messages parallelKafka消费者并行阅读消息
【发布时间】:2021-03-07 04:25:17
【问题描述】:

我们是否可以让多个消费者从一个主题中消费以在 kafka 中实现并行处理。 我的用例是从单个分区并行读取消息。

【问题讨论】:

  • 请编辑问题以将其限制为具有足够详细信息的特定问题,以确定适当的答案。避免一次问多个不同的问题。有关澄清此问题的帮助,请参阅如何提问页面。 stackoverflow.com/help/how-to-ask
  • 另外,请说明您已经完成了哪些研究,例如阅读 kafka 文档。

标签: apache-kafka


【解决方案1】:

简单地说,默认情况下我们无法为消费者实现分区级别的并行性。

但是你可以试试Akka Streams Kafka (Reactive kafka)。一旦通过这些文档。

【讨论】:

    【解决方案2】:

    分区数定义了从 kafka 主题中读取的并行级别。但是阅读(或多或少)仅受您的网络容量限制。

    一个好的模式是将消息的读取和处理分开(每个主题分区一个线程用于读取,多个线程用于处理此消息)。

    【讨论】:

    • 是的。这是我们目前计划做的。但是这里的问题是,我们不能在处理时再次重新处理失败的消息
    • “我们不能再重新处理消息”是什么意思?您的意思是消息处理可能由于某种原因而失败,您需要稍后再试吗?我可以看到两个可能的操作:关闭自动偏移提交,并在您处理某个批次时手动执行此操作。选项二:将此事件发送给您也使用的死信主题。
    • 第二个选项看起来不错。我会试试的。谢谢
    • 请谨慎使用此选项 - 我建议为您发送到死信主题的消息添加一些计数器或最后处理时间戳(以避免一次又一次地处理相同的消息)跨度>
    【解决方案3】:

    是的,您可以使用多个 Kafka 消费者并行处理消息,但不,如果您只有一个分区,这是不可能的。

    Kafka 消费中的并行度由分区数量定义,您可以随时轻松地重新分区主题以创建更多分区。

    下面是如何使用rapids-kafka-client 并行处理消息的示例,这是一个使 Kafka 并行消费更容易的库。

    public static void main(String[] args){
      ConsumerConfig.<String, String>builder()
          .prop(KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName())
          .prop(VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName())
          .prop(GROUP_ID_CONFIG, "stocks")
          .topics("stock_changed")
          .consumers(7)
          .callback((ctx, record) -> {
            System.out.printf("status=consumed, value=%s%n", record.value());
          })
          .build()
          .consume()
          .waitFor();
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-11-09
      • 2020-04-13
      • 1970-01-01
      • 2017-09-23
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多