【问题标题】:Apache Camel: async operation and backpressureApache Camel:异步操作和背压
【发布时间】:2017-10-31 22:46:30
【问题描述】:

在 Apache Camel 2.19.0 中,我想在并发 seda 队列上异步生成消息并使用结果,同时在 seda 队列上的执行程序已满时阻塞。 它背后的用例:我需要处理多行的大文件,并且需要为它创建批处理,因为每一行的单个消息开销太大,而我无法将整个文件放入堆中。但最后,我需要知道我触发的所有批次是否都已成功完成。 如此有效,我需要一个背压机制来向队列发送垃圾邮件,同时又想利用多线程处理。

这是 Camel 和 Spring 中的一个简单示例。我配置的路由:

package com.test;

import org.apache.camel.builder.RouteBuilder;
import org.springframework.stereotype.Component;

@Component
public class AsyncCamelRoute extends RouteBuilder {

    public static final String ENDPOINT = "seda:async-queue?concurrentConsumers=2&size=2&blockWhenFull=true";

    @Override
    public void configure() throws Exception {
        from(ENDPOINT)
                .process(exchange -> {
                    System.out.println("Processing message " + (String)exchange.getIn().getBody());
                    Thread.sleep(10_000);
                });
    }
}

生产者长这样:

package com.test;

import org.apache.camel.ProducerTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.event.ContextRefreshedEvent;
import org.springframework.context.event.EventListener;
import org.springframework.stereotype.Component;

import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.CompletableFuture;

@Component
public class AsyncProducer {

    public static final int MAX_MESSAGES = 100;

    @Autowired
    private ProducerTemplate producerTemplate;

    @EventListener
    public void handleContextRefresh(ContextRefreshedEvent event) throws Exception {
        new Thread(() -> {
            // Just wait a bit so everything is initialized
            try {
                Thread.sleep(5_000);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            List<CompletableFuture> futures = new ArrayList<>();

            System.out.println("Producing messages");
            for (int i = 0; i < MAX_MESSAGES; i++) {
                CompletableFuture future = producerTemplate.asyncRequestBody(AsyncCamelRoute.ENDPOINT, String.valueOf(i));
                futures.add(future);
            }
            System.out.println("All messages produced");

            System.out.println("Waiting for subtasks to finish");
            futures.forEach(CompletableFuture::join);
            System.out.println("Subtasks finished");
        }).start();

    }
}

这段代码的输出如下:

Producing messages
All messages produced
Waiting for subtasks to finish
Processing message 6
Processing message 1
Processing message 2
Processing message 5
Processing message 8
Processing message 7
Processing message 9
...
Subtasks finished

因此,blockIfFull 似乎被忽略,所有消息都在处理之前创建并放入队列。

有什么方法可以创建消息,以便我可以在骆驼中使用异步处理,同时确保如果有太多未处理的元素,将元素放入队列会阻塞?

【问题讨论】:

  • 你能用requestBody(..)代替asyncRequestBody(..)吗?您最终可能会在用于执行异步消息发送的池中遇到大量阻塞线程。而不是阻塞你的客户端线程。
  • 嗨@Ralf,我不太明白你的方法 - requestBody 让客户端(生产者)阻塞,直到消费者完成。虽然我想阻止客户端向消费者发送垃圾邮件,但只要有消费者,它就应该创建消息。但是,我使用不同的方法解决了它。
  • 没错。但是,如果您执行任何异步操作,那么另一个线程正在执行提交到 seda 并等待响应的工作。除非处理异步任务的线程池耗尽,否则您在其中运行循环并调用 asyncRequestBody(..) 的线程不会被阻塞。但是,如果线程是根据需要在池中创建的,那么您将永远不会看到循环线程被阻塞。
  • 感谢您的解释。我有点希望有一些我忽略的功能,其行为类似于使用普通的 Java ExecutorService。这意味着生产者基本上可以将任务放到 ExecutorService 上,直到底层队列已满,然后阻塞直到再次有可用空间。但是从您的解释看来,似乎只有同步和异步的可能性,它们的行为都与我的预期不同。但如前所述,我现在解决了我自己的答案中描述的问题,这似乎可以按预期工作。

标签: java asynchronous apache-camel


【解决方案1】:

我通过使用流媒体和自定义拆分器解决了这个问题。通过这样做,我可以使用返回行列表而不是仅单行的迭代器将源代码行拆分为块。有了这个,在我看来,我可以根据需要使用 Camel。

所以路由包含以下部分:

.split().method(new SplitterBean(), "splitBody").streaming().parallelProcessing().executorService(customExecutorService)

使用具有上述行为的定制拆分器。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2023-03-22
    • 2020-10-27
    • 2018-12-16
    • 2023-03-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多