【问题标题】:Reactor Kafka - Bottleneck if the number of consumer is greater than the number of partitionReactor Kafka - 消费者数量大于分区数量时的瓶颈
【发布时间】:2023-02-17 19:18:25
【问题描述】:

如果 kafka 消费者的数量(远)大于分区数量,分区数量是否会成为性能瓶颈?

假设我有一个名为 the-topic 的主题,只有三个分区。

现在,我有这个下面的应用程序,以便从主题中消费:

@Service
public class MyConsumer implements CommandLineRunner {

    @Autowired
    private KafkaReceiver<String, String> kafkaReceiver;


    @Override
    public void run(String... args) {
        myConsumer().subscribe();
    }

    public Flux<String> myConsumer() {
        return kafkaReceiver.receive()
                .flatMap(oneMessage -> consume(oneMessage))
                .doOnNext(abc -> System.out.println("successfully consumed {}={}" + abc))
                .doOnError(throwable -> System.out.println("something bad happened while consuming : {}" + throwable.getMessage()));
    }

    private Mono<String> consume(ConsumerRecord<String, String> oneMessage) {
        // this first line is a heavy in memory computation which transforms the incoming message to a data to be saved.
        // it is very intensive computation, but has been tested NON BLOCKING by different tools, and takes 1 second :D
        String transformedStringCPUIntensiveNonButNonBLocking = transformDataNonBlockingWithIntensiveOperation(oneMessage);
        //then, just saved the correct transformed data into any REACTIVE repository :)
        return myReactiveRepository.save(transformedStringCPUIntensiveNonButNonBLocking);
    }

}

我将应用程序 docker 化并部署在 Kubernetes 中。

借助云提供商,我能够轻松部署其中的 60 个容器和 60 个应用程序。

并假设为了这个问题,我的每个应用程序都具有超强的弹性,从不崩溃。

这是否意味着,由于主题只有三个分区,任何时候都会浪费 57 个其他消费者?

当分区数量较少时,如何从扩展容器数量中获益?

【问题讨论】:

  • 为什么需要 60 个消费者?这背后的逻辑是什么。增加分区数量会增加吞吐量,但也会有一些缺点,例如在进行新领导者选举时会增加停机时间。您目前有 3 个分区,如果有 60 个消费者,那么这些消费者中的大多数将处于非活动状态。通常你应该有和分区一样多的消费者

标签: java apache-kafka reactor-kafka


【解决方案1】:

既然主题只有三个分区,那在任何时候,57个其他消费者会被浪费掉吗?

是的。这就是 Kafka 消费者 API 的工作方式。您围绕它使用的框架无关紧要。

当分区数量较少时,增加容器数量的好处

您需要将事件处理(保存到存储库)与实际消费/轮询循环分开。例如,将转换后的事件推送到非阻塞的外部队列/外部 API,而无需等待响应。然后在该 API 端点上设置自动缩放器。

【讨论】:

  • 感谢@OneCricketeer 的回答。感谢您确认浪费资源,赞成。在我接受之前,我想请您对您的回答的第二部分进行一些澄清。如果您查看我的 sn-p,.flatMap(oneMessage -&gt; consume(oneMessage)) + private Mono&lt;String&gt; consume(ConsumerRecord&lt;String, String&gt; oneMessage) 事件处理(保存到存储库)以及事件转换已经是非阻塞的。反应堆核心/事件循环是否可以在不依赖外部 API 的情况下“分离事件处理”?
  • 不过,它并没有解耦。在函数结束并处理下一个事件之前,您正在等待 save 方法的结果。然后你的 doOnNext / doOnError 正在阻止调用打印到 IO 流
  • 明白了。对于保存方法,我使用了 ReactiveCassandraRepository,并测试了尝试插入非常重要的数据。我计时并使用了 BlockHound,我相信它是非阻塞的。会仔细检查。至于 doOnNext / doOnError,我使用的是异步 LOGGER。但是,让我再次检查一下
【解决方案2】:

这是否意味着,由于主题只有三个分区,任何时候都会浪费 57 个其他消费者?

是的。对于单个消费者组,您可以拥有与分区数量一样多的并发消费者。

当分区数量较少时,如何从扩大容器数量中获益?

您可能会尝试将这些容器中的每一个注册到分区数量不同的消费者组中。那可行。但是正如@OneCricketeer 所提到的,您应该有一个单独的事件处理管道。如果您不想多次处理同一个事件,那将是最好的方法。

【讨论】:

  • 单独的消费者组在这里效果不佳,因为这会在看似数据库保存操作的情况下复制数据
  • 是的,多个消费者群体会多次读取相同的数据。但是,当添加到数据库时,可以检查数据是否已经保存到数据库中,如果是则不保存。
  • 也可以尝试类似的方法,一个消费者组从偶数偏移量读取,而其他消费者组从奇数偏移量读取。这将节省对 db 的额外调用以检查它是否已被处理。
  • 这不是 Kafka 消费者群体的工作方式。偏移量按顺序读取,而不是选择性地读取
  • 我并不是说我们必须有选择地读取偏移量,但是当消费者从主题的分区中消费消息时,它可以根据其所在的消费者组跳过分区中奇数或偶数偏移量的消息。我相信api 可以提供消息的偏移量。其他消费者组可以继续使用剩余的偏移量,因为我相信每个消费者组都维护了偏移量。
猜你喜欢
  • 2014-02-13
  • 2020-09-19
  • 2018-10-03
  • 1970-01-01
  • 1970-01-01
  • 2017-01-04
  • 1970-01-01
  • 2018-12-28
  • 2023-01-31
相关资源
最近更新 更多