【问题标题】:How to start @KafkaListener based on start flag如何根据启动标志启动@KafkaListener
【发布时间】:2020-09-21 06:48:11
【问题描述】:

我仅在标志设置为 true 时才尝试启动我的 KafkaListener。

@Component
public class KafkaTopicConsumer {

//Somehow wrap the listener to only start when a property value is set to true

@KafkaListener(topics = "#{@consumerTopic}", groupId = "#{@groupName}")
public void consumeMessage(ConsumerRecord<String, String> message) throws IOException {
    logger.info("Consumed message from topic: {} with message: {}", message.topic(), message);
}

有没有办法确保监听器在诸如 start.consumer 属性设置为 true 时才启动?我不希望监听器在每次启动应用程序时仅在我指定我希望启动它时启动。有没有处理这个用例的好方法?

【问题讨论】:

    标签: java apache-kafka


    【解决方案1】:

    首先,您需要将autoStartup 设置为false 并为您的容器命名。然后你需要使用@EventListener根据标志手动启动它。

    @Component
    public class KafkaTopicConsumer {
        @Autowired
        private KafkaListenerEndpointRegistry kafkaListenerEndpointRegistry;
    
        @Value("${start.consumer}")
        private boolean shouldStart;
    
        @KafkaListener(id = "myListener", autoStartup = "false", topics = "#{@consumerTopic}", groupId = "#{@groupName}")
        public void consumeMessage(ConsumerRecord<String, String> message) throws IOException {
            logger.info("Consumed message from topic: {} with message: {}", message.topic(), message);
        }
    
        @EventListener
        public void onStarted(ApplicationStartedEvent event) {
            if (shouldStart) {
                MessageListenerContainer listenerContainer = kafkaListenerEndpointRegistry.getListenerContainer("myListener");
                listenerContainer.start();    
            }
        }
    }
    

    注意:@EventListener 将确保容器正确加载,如果您使用@PostConstruct,它可能无法正常工作。

    编辑

    使用@Value 注解添加了属性的实际读取。

    注意:这种方法具有额外的灵活性,只需进行一些更改即可动态调用startstop 方法(例如使用JMX)。这有助于我们希望禁用消费者并稍后启用它而不重新启动应用程序的场景。

    正如@Makoton's answer 中正确说明的那样,另一个好方法是使用@ConditionalOnProperty。请注意,在您的示例中,您可以将其与@Component 一起使用,而不是手动定义@Bean

    @Component
    @ConditionalOnProperty(
            value = "start.consumer",
            havingValue = "true")
    public class KafkaTopicConsumer { // ...
    

    这一切都取决于您需要的灵活性水平。

    【讨论】:

    • 太棒了,谢谢!一旦我能够运行和测试它,我一定会回头看看我是否需要更多信息并将其标记为答案!
    【解决方案2】:

    您可以将 ConditionalsBeans 与属性一起使用

    @Bean
    @ConditionalOnProperty(
      value="my.custom.flag", 
      havingValue = "true")
    public KafkaListener kafkaListener{
     .....
    }
    

    条件 bean 允许您基于属性或自定义条件启动 bean。 Reference

    【讨论】:

    • 太棒了,谢谢!一旦我能够运行和测试它,我一定会回头看看我是否需要更多信息并将其标记为答案!
    • 在我看来这是最好的方法。
    猜你喜欢
    • 2020-10-23
    • 1970-01-01
    • 1970-01-01
    • 2012-04-15
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多