【问题标题】:Consume messages from Kafka only if some condition true仅在某些条件为真时才使用来自 Kafka 的消息
【发布时间】:2019-12-16 05:29:48
【问题描述】:

我们有特定的主题,只有在条件 consumeEnabled=true 时才需要消费消息。 所以,它应该像这样工作:

  1. 如果应用程序正在启动并且consumeEnabled=true,则分配 分区给消费者并使用来自主题的消息。
  2. 如果应用程序正在启动并且consumeEnabled=false,则不要将分区分配给消费者,也不要使用来自主题的消息。
  3. 如果应用程序已在使用 consumeEnabled=false 的情况下运行,但在运行时属性变为 consumeEnabled=true,则在运行时将分区分配给消费者并从主题消费消息。

应用正在消费消息的情况,但随后consumeEnabled变为false,无需考虑。

请您定义使用 Spring Kafka 和/或 Kafka Java 客户端实现决策的最佳方式

【问题讨论】:

    标签: java apache-kafka spring-kafka


    【解决方案1】:

    如果您使用的是@KafkaListener 那么

    @KafkaListener(id = "foo", ... , autoStartup="${consume.enabled}")
    

    consume.enabled 是一个属性。

    要在运行时启动/停止容器,请使用 KafkaListenerEndpointRegistry bean。

    registry.getListenerContainer("foo").start();
    

    【讨论】:

      【解决方案2】:

      您可以将您的消费者置于一个简单的线程中,以切换消费者对象的轮询状态。

      public class EnabledConsumer implements Runnable {
      
          private Consumer consumer;
          private boolean enabled;
      
          public EnabledConsumer(Consumer consumer, boolean enabled) {
              this.consumer = consumer;
              this.enabled = enabled;
          }
      
          public void setEnabled(boolean enable) {
              this.enabled = enable;
          }
      
          @Override
          public void run() {
              while(enabled) {
                  ConsumerRecords records = consumer.poll(...);
                  ...
              }
      
      }
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2014-04-28
        • 1970-01-01
        • 1970-01-01
        • 2014-05-31
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多