【问题标题】:Consuming from Camel queue every x minutes每 x 分钟从 Camel 队列中消费
【发布时间】:2019-11-14 00:33:44
【问题描述】:

尝试实现一种方法,让我的消费者每 30 分钟左右从队列接收消息。 就上下文而言,我的错误队列中有 20 条消息,直到 x 分钟过去,然后我的路由消耗队列上的所有消息并继续“休眠”,直到又过去了 30 分钟。

不确定实现这一点的最佳方法,我尝试过 spring @Scheduled、camel timer 等,但都没有达到我的预期。我一直试图让它与路由策略一起使用,但在正确的功能中没有骰子。它似乎立即从队列中消耗。

路由策略是正确的路径还是有其他东西可以使用?

【问题讨论】:

    标签: spring-boot apache-camel


    【解决方案1】:

    从队列中读取的路由总是尽可能快地读取任何消息。

    您可以做的一件事是启动/停止或暂停使用消息的路由,因此请进行这种设置:

    route 1: error_q_reader, which goes from('jms').
    route 2: a timed route that fires every 20 mins
    

    路由 2 可以使用 control bus 组件来启动路由。

     from('timer?20mins') // or whatever the timer syntax is...
        .to("controlbus:route?routeId=route1&action=start")
    

    这里的棘手部分是知道何时停止路线。你让它运行5分钟吗?一旦消息全部消耗,您想停止它吗?可能有一种方法可以运行另一条可以检查队列深度的路由(比如每 1 分钟左右),如果它是 0,则关闭route 1,你可能会让它工作,但我可以向你保证这会变得混乱您尝试处理许多异步操作。

    您还可以尝试一些更奇特的东西,例如自定义QueueBrowseStrategy,它可以在队列中没有消息时触发关闭route 1 的事件。

    【讨论】:

      【解决方案2】:

      我构建了一个客户 bean 来排空队列并关闭,但这不是一个非常优雅的解决方案,我很想找到一个更好的解决方案。

      public class TriggeredPollingConsumer {
      
          private ConsumerTemplate consumer;
          private Endpoint consumerEndpoint;
          private String endpointUri;
          private ProducerTemplate producer;
          private static final Logger logger = Logger.getLogger( TriggeredPollingConsumer.class );
      
          public TriggeredPollingConsumer() {};
      
          public TriggeredPollingConsumer( ConsumerTemplate consumer, String endpoint, ProducerTemplate producer ) {
              this.consumer = consumer;
              this.endpointUri = endpoint;
              this.producer = producer;
          }
      
          public void setConsumer( ConsumerTemplate consumer) {
              this.consumer = consumer;
          }
      
          public void setProducer( ProducerTemplate producer ) {
              this.producer = producer;
          }
      
          public void setConsumerEndpoint( Endpoint endpoint ) {
              consumerEndpoint = endpoint;
          }
      
          public void pollConsumer() throws Exception {
              long count = 0;
              try {
                  if ( consumerEndpoint == null ) consumerEndpoint = consumer.getCamelContext().getEndpoint( endpointUri );
                  logger.debug( "Consuming: " + consumerEndpoint.getEndpointUri() );
                  consumer.start();
                  producer.start();
                  while ( true ) {
                      logger.trace("Awaiting message: " + ++count );
                      Exchange exchange = consumer.receive( consumerEndpoint, 60000 );
                      if ( exchange == null ) break;
                      logger.trace("Processing message: " + count );
                      producer.send( exchange );
                      consumer.doneUoW( exchange );
                      logger.trace("Processed message: " + count );
                  }
                  producer.stop();
                  consumer.stop();
                  logger.debug( "Consumed " + (count - 1) + " message" + ( count == 2 ? "." : "s." ) );
              } catch ( Throwable t ) {
                  logger.error("Something went wrong!", t );
                  throw t;
              }
          }
      
      }
      

      您配置 bean,然后从计时器调用 bean 方法,并配置直接路由来处理队列中的条目。

      from("timer:...")
        .beanRef("consumerBean", "pollConsumer");
      from("direct:myRoute")
        .to(...);
      

      然后它将读取队列中的所有内容,并在一分钟内没有条目到达时立即停止。您可能想减少分钟,但我发现一秒钟意味着如果 JMS 有点慢,它会在排空队列的中途超时。

      我也一直在研究 sjms-batch 组件,以及它如何与 pollEnrich 模式一起使用,但到目前为止我还没有能够让它工作。

      【讨论】:

        【解决方案3】:

        我通过在微服务方法中将我的应用程序用作 CronJob 解决了这个问题,并赋予它优雅地自行关闭的能力,我们可以设置属性 camel.springboot.duration-max-idle-seconds。因此,您的 JMS 消费者路线保持简单。

        另一种方法是声明一个路由来控制 JMS 消费者路由的“生命周期”(启动、睡眠和恢复)。

        我强烈建议您使用第一种方法。

        【讨论】:

          【解决方案4】:

          如果你使用ActiveMQ,你可以利用它的Scheduler feature

          您可以延迟在代理上传递消息,只需将 JMS 属性 AMQ_SCHEDULED_DELAY 设置为延迟的毫秒数。骆驼路线很容易

          .setHeader("AMQ_SCHEDULED_DELAY", 60000)
          

          这并不完全符合您的要求,因为它不是每 30 分钟排空一个队列,而是将每条单独的消息延迟 30 分钟。

          请注意,您必须在代理配置中启用schedulerSupport。否则延迟属性将被忽略。

          <broker brokerName="localhost" dataDirectory="${activemq.data}" schedulerSupport="true">
          ...
          </broker>
          

          【讨论】:

            【解决方案5】:

            你应该考虑Aggregation EIP

            from(URI_WAITING_QUEUE)
                .aggregate(new GroupedExchangeAggregationStrategy())
                    .constant(true)
                    .completionInterval(TIMEOUT)
                .to(URI_PROCESSING_BATCH_OF_EXCEPTIONS);
            

            此示例描述了以下规则:所有传入URI_WAITING_QUEUE 的对象将被分组到Listconstant(true) 是一个分组条件(没有任何)。并且每个TIMEOUT 周期(以毫秒为单位)所有分组对象都将被传递到URI_PROCESSING_BATCH_OF_EXCEPTIONS 队列中。

            所以URI_PROCESSING_BATCH_OF_EXCEPTIONS 队列将处理要处理的对象的List。可以申请Split EIP进行拆分,一一处理。

            【讨论】:

              猜你喜欢
              • 1970-01-01
              • 2015-12-13
              • 1970-01-01
              • 1970-01-01
              • 1970-01-01
              • 1970-01-01
              • 1970-01-01
              • 2011-02-05
              • 1970-01-01
              相关资源
              最近更新 更多