【问题标题】:How to implement consumer Thread Safety with Spring-Kafka如何使用 Spring-Kafka 实现消费者线程安全
【发布时间】:2020-02-21 20:34:27
【问题描述】:

我正在使用 spring boot 2.1.7.RELEASE 和 spring-kafka 2.2.7.RELEASE。并且我正在使用 @KafkaListener 注释来创建消费者,并且我正在使用消费者的所有默认设置。

这是我的消费者配置:

@Configuration
@EnableKafka
public class KafkaConsumerCommonConfig implements KafkaListenerConfigurer {

      @Bean
      public <K,V> ConcurrentKafkaListenerContainerFactory<String, Object> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(primaryConsumerFactory());
        factory.getContainerProperties().setMissingTopicsFatal(false);
        return factory;
      }

      @Bean
      public DefaultKafkaConsumerFactory<Object, Object> primaryConsumerFactory() {
        return new DefaultKafkaConsumerFactory<>(sapphireKafkaConsumerConfig.getConfigs());
      }

}

由于某些原因,我在同一个应用程序中有多个消费者,如下所示。

@KafkaListener(topics = "TEST_TOPIC1")
public void consumer1(){

}

@KafkaListener(topics = "TEST_TOPIC2")
public void consumer2(){

}

@KafkaListener(topics = "TEST_TOPIC3")
public void consumer3(){

}

话虽如此,根据关于“消费者线程安全”的融合文档

你不能有多个属于同一组的消费者 线程,你不能让多个线程安全地使用同一个线程 消费者。每个线程一个消费者是规则。运行多个 同一组中的消费者在一个应用程序中,您将需要运行 每个都在自己的线程中。将消费者逻辑包装在 自己的对象然后用Java的ExecutorService启动多个 每个线程都有自己的消费者。

现在我的问题是,在使用上面的代码使用 spring-kafka 来解决这种情况时,我是否应该做任何额外的事情,或者没关系,因为组会随机生成,因为我没有指定?请提出建议。

【问题讨论】:

    标签: spring-kafka


    【解决方案1】:

    您的所有听众都针对不同的主题进行了配置。所以,他们的组完全无关紧要。当同一主题有多个消费者时,该小组就会出现。

    【讨论】:

    • Spring 负责“每个线程一个消费者”语义 - 但是,如果您使用任何更高级的 API 公开原始 Consumer 对象(事件、侦听器等),您 必须遵守 Kafka 指令 - 不使用来自调用侦听器的线程以外的线程对该使用者对象的引用(默认情况下,它将调用任何应用程序) Consumer 在事件中暴露的事件侦听器)。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-08-13
    • 2020-08-20
    • 1970-01-01
    • 1970-01-01
    • 2020-10-19
    • 2016-07-16
    相关资源
    最近更新 更多