【问题标题】:Spring Boot kafkaTemplate consumer message load and processing messageSpring Boot kafkaTemplate 消费者消息加载和处理消息
【发布时间】:2021-09-15 21:06:03
【问题描述】:

在我的应用程序中,我使用Spring Boot kafkaTemplate 来使用消息。我是使用 Spring Boot 的 kafka 新手。我添加了一个消费者如下 -

 @KafkaListener(topics = "#{'${app.kafka.consumer.topic}'.split(',')}")
 public void receivedMessage(ConsumerRecord<String, String> cr, @Payload String message){
    log.info("Message received from topic {} ", cr.topic());
    //TODO
}

在topic 上,我们每秒将收到近 20 万条消息。我收到的message 将发送到另一个处理方法,该方法根据特定条件过滤消息​​,然后将过滤后的message 发布到另一个topic。

我的问题是,上述@KafkaListener 方法是否会处理此负载,或者我是否需要进行任何特殊处理,例如threading 或concurrency。

【问题讨论】:

    标签: java spring-boot apache-kafka spring-kafka


    【解决方案1】:

    这完全取决于您的工作负载、CPU 中的内核数量等。一般来说,增加并发性将提供更多的吞吐量(只要您在主题上至少有那么多分区)。

    但是,您的下游代码必须是完全线程安全的。

    即便如此,您的代码中也可能存在其他瓶颈(数据库等)。

    如果您所做的只是计算、转换和发布到另一个主题,而不进行任何其他 I/O,那么增加并发肯定会有所帮助。

    唯一真正的解决方案是实验,如果您没有获得所需的吞吐量,请分析您的应用程序。

    如果不实际操作,您无法学习这些技能。

    【讨论】:

    • 嗨@Gary,感谢您的回复。我们决定轮询来自 kafka 主题的一批消息,因为我正在使用 nack() 以防万一,但我们使用的是 Spring Boot 2.1.8 并且看起来 nack() 在 2.1.8 版本中不受支持。你能告诉我Spring Boot 2.1.8中的等效nack()功能吗
    • 不要在cmets中问无关的问题;它不能帮助人们找到问题/答案。 Spring Boot 2.1.x 已经报废近一年; nack() 是在 spring-kafka 2.3.x 中引入的(Boot 2.2.x 附带,这也是生命终结)。请参阅github.com/spring-projects/spring-boot/wiki/Supported-Versions 如您所见,Boot 2.3.x 也即将结束生命周期。在早期版本中没有 nack() 的等价物。有关 spring-kafka/spring-boot 版本的矩阵,请参见项目页面。 spring.io/projects/spring-kafka#overview
    猜你喜欢
    • 1970-01-01
    • 2023-04-11
    • 2016-05-04
    • 1970-01-01
    • 2020-08-11
    • 1970-01-01
    • 2020-11-11
    • 1970-01-01
    • 2016-11-11
    相关资源
    最近更新 更多