【问题标题】:Spring Integration multiple queue consumersSpring集成多个队列消费者
【发布时间】:2020-08-03 15:19:22
【问题描述】:

Spring integration MessageQueue without polling 的结果是,我有一个轮询器,它使用自定义 TaskScheduler 立即使用队列中的消息:

ThreadPoolTaskScheduler scheduler = new ThreadPoolTaskScheduler();
scheduler.setThreadNamePrefix("resultProcessor-");

IntegrationFlows
  .from("inbound")
  .channel(MessageChannels.priority().get())
  .bridge(bridge -> bridge
    .taskScheduler(taskScheduler)
    .poller(Pollers.fixedDelay(0).receiveTimeout(Long.MAX_VALUE)))
  .fixedSubscriberChannel()
  .route(inboundRouter())
  .get()

现在我想让多个线程同时使用,所以我尝试了:

ThreadPoolTaskScheduler scheduler = new ThreadPoolTaskScheduler();
scheduler.setThreadNamePrefix("resultProcessor-");
scheduler.setPoolSize(4);

但是,由于在AbstractPollingEndpoint 中任务调度器调度了一个同步轮询器(有点复杂),所以只创建了一个线程。如果我将 TaskExecutor 设置为 SyncTaskExecutor(默认)以外的任何值,我会遇到大量计划任务(请参阅 Spring integration MessageQueue without polling)。

如何在 Spring Integration 中同时从队列中消费?这似乎很基本,但我找不到解决方案。

我可以使用ExecutorChannel 代替队列,但是,(AFAIK)然后我会丢失队列功能,例如优先级、队列大小和我所依赖的指标。

【问题讨论】:

    标签: spring-integration


    【解决方案1】:

    见PollerSpec.taskExecutor():

    /**
     * Specify an {@link Executor} to perform the {@code pollingTask}.
     * @param taskExecutor the {@link Executor} to use.
     * @return the spec.
     */
    public PollerSpec taskExecutor(Executor taskExecutor) {
    

    这样,在根据您的taskScheduler 和delay 定期安排任务后,真正的任务将在提供的执行程序的线程上执行。默认情况下,它确实在调度程序的线程上执行任务。

    更新

    我不确定这是否满足您的要求,但这是保持您的队列逻辑和处理并行处理的唯一方法:

     .bridge(bridge -> bridge
        .taskScheduler(taskScheduler)
        .poller(Pollers.fixedDelay(0).receiveTimeout(Long.MAX_VALUE)))
     .channel(channels -> channel.executor(threadPoolExecutor()))   
     .fixedSubscriberChannel()
    

    【讨论】:

    • 但是如果任务调度器的延迟为0,它会在几秒钟内向任务执行器发送数千个轮询任务?
    • 嗯,从技术上讲,是的,但这是您的配置。你还能期待什么?另一种方法是在bridge() 之后使用ExecutorChannel,因此只有真正的拉取消息才会被并行处理。但同样:如果您想并行处理,那么优先级是什么?
    • 目前,大约有 70 万条消息排队,应该首先处理高优先级消息。我期望的是一种让多个消费者处理队列中的这些消息的方法,而无需定期轮询——就像 JMS 消费者一样。如果我桥接到一个 ExecutorChannel,所有消息都会立即发送到它,而违背了优先级队列的目的,对吧?
    • 在我的回答中查看更新。
    • 顺便说一句。我不确定为什么在引用的问题中我有一个 fixedSubscriberChannel()。我不再使用它了。没有它,taskScheduler 和 taskExecutor 就不起作用(并不意味着它可以 with 它工作 - 我没有测试)。同时,我尝试了另一种方法并取得了成功,该方法将我发布为答案。感谢您的帮助!
    【解决方案2】:

    我可以这样解决:

    • 执行轮询的单线程任务调度程序
    • 具有同步队列的线程池执行器

    这样,任务调度器可以给每个执行器 1 个任务并在没有执行器空闲时阻塞,从而不会耗尽源队列或垃圾邮件任务。

      @Bean
      public IntegrationFlow extractTaskResultFlow() {
        return IntegrationFlows
          .from(ChannelNames.TASK_RESULT_QUEUE)
          .bridge(bridge -> bridge
            .taskScheduler(taskResultTaskScheduler())
            .poller(Pollers
              .fixedDelay(0)
              .taskExecutor(taskResultExecutor())
              .receiveTimeout(Long.MAX_VALUE)))
          .handle(resultProcessor)
          .channel(ChannelNames.TASK_FINALIZER_CHANNEL)
          .get();
      }
    
      @Bean
      public TaskExecutor taskResultExecutor() {
        ThreadPoolExecutor executor = new ThreadPoolExecutor(
          1, // corePoolSize
          8, // maximumPoolSize
          1L, // keepAliveTime
          TimeUnit.MINUTES,
          new SynchronousQueue<>(),
          new CustomizableThreadFactory("resultProcessor-")
        );
        executor.setRejectedExecutionHandler(new CallerBlocksPolicy(Long.MAX_VALUE));
        return new ErrorHandlingTaskExecutor(executor, errorHandler);
      }
    
      @Bean
      public TaskScheduler taskResultTaskScheduler() {
        ThreadPoolTaskScheduler scheduler = new ThreadPoolTaskScheduler();
        scheduler.setThreadNamePrefix("resultPoller-");
        return scheduler;
      }
    

    (最初的示例是从链接的问题中复制的,这个现在类似于我的实际解决方案)

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-06-02
      • 2013-04-17
      • 1970-01-01
      • 2016-11-11
      相关资源
      最近更新 更多