【问题标题】:How can I process @KafkaListener method in different threads?如何在不同的线程中处理 @KafkaListener 方法?
【发布时间】:2019-06-24 14:44:27
【问题描述】:

我在 Spring Boot 中有 kafka 处理程序:

    @KafkaListener(topics = "topic-one", groupId = "response")
    public void listen(String response) {
        myService.processResponse(response);
    }

例如生产者每秒发送一条消息。但是myService.processResponse 工作 10 秒。我需要处理每条消息并在新线程中启动myService.processResponse。我可以创建我的执行者并将每个响应委托给它。但我认为 kafka 中还有其他配置供他们使用。我找到了 2:

1) 将 concurrency = "5" 添加到 @KafkaListener 注释 - 它似乎正在工作。但我不确定有多正确,因为我有第二种方法:

2) 我可以创建ConcurrentKafkaListenerContainerFactory 并将其设置为ConsumerFactoryconcurrency

我不明白这些方法之间的区别?只需将concurrency = "5" 添加到@KafkaListener 注释就足够了,还是我需要创建ConcurrentKafkaListenerContainerFactory

或者我什么都不懂,还有其他办法吗?

【问题讨论】:

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


    【解决方案1】:

    在管理已提交的偏移方面,使用执行器会使事情变得复杂;不推荐。

    使用@KafkaListener,框架会为您创建一个ConcurrentKafkaListenerContainerFactory

    注解上的concurrency只是为了方便;它会覆盖出厂设置。

    这允许您使用具有多个侦听器的同一工厂,每个侦听器具有不同的并发性。

    您可以使用启动属性设置容器并发(默认);该值被注释值覆盖;请参阅 javadocs...

    /**
     * Override the container factory's {@code concurrency} setting for this listener. May
     * be a property placeholder or SpEL expression that evaluates to a {@link Number}, in
     * which case {@link Number#intValue()} is used to obtain the value.
     * <p>SpEL {@code #{...}} and property place holders {@code ${...}} are supported.
     * @return the concurrency.
     * @since 2.2
     */
    String concurrency() default "";
    

    【讨论】:

    • 还是不完全明白。即,如果我用@KafkaListener 注释该方法 - 此方法中的每个句柄将在新线程中?因为 spring 创建默认 ConcurrentKafkaListenerContainerFactory 而我不需要创建 ConcurrentKafkaListenerContainerFactory ?
    • concurrency = 5,主题上至少需要5个partition;分区将分布在线程中 - 一个分区只能由一个线程处理。您不需要创建工厂,因为引导会通过自动配置为您完成。
    • 在管理提交的偏移方面,使用执行器会使事情变得复杂;不推荐。
    • 如果更容易,请告诉我我应该怎么做才能得到这个逻辑 - 1)我在 @KafkaListener 方法中处理消息。 2) 开始在新线程中处理此消息。 3) 处理新消息....等等
    • 再一次;只有在主题上至少有该数量的分区时,才能增加并发性; Kafka 不允许同一组的分区上有多个消费者。如果您移交给另一个线程,则在管理为分区提交的偏移量方面会增加很多复杂性。您需要了解更多关于 Kafka 的信息。
    【解决方案2】:

    concurrency 选项与并发处理同一消费者收到的消息无关。当您有多个消费者,每个消费者都处理自己的分区时,它适用于消费者组。

    将处理传递给单独的线程非常复杂,我相信 Spring-Kafka 团队决定不“按设计”这样做。您甚至无需深入研究 Spring-Kafka 即可了解原因。查看KafkaConsumer's Detecting Consumer Failures 文档:

    必须注意确保提交的偏移量不会得到 领先于实际位置。通常,您必须禁用自动 仅在之后提交和手动提交记录的已处理偏移量 线程已完成处理它们(取决于交付 你需要的语义)。另请注意,您需要暂停 分区,这样直到之后才从轮询中收到新记录 线程已经完成了之前返回的处理。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2020-02-03
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2011-01-11
      • 1970-01-01
      • 2011-11-16
      相关资源
      最近更新 更多