简短的回答是没有答案:Camel 文件组件的sortBy 选项内存效率太低,无法容纳我的用例:
- 独特性:如果文件已经存在,我不想将其放入队列中。
- 优先级:应首先处理标记为高优先级的文件。
- 性能:拥有几十万甚至几百万个文件应该没问题。
- FIFO:(奖励)最旧的文件(按优先级)应该首先被拾取。
问题似乎是,如果我正确阅读了source code 和documentation,所有文件详细信息都在内存中以执行排序,无论是内置语言还是自定义可插入sorter用来。文件组件总是创建一个包含所有详细信息的对象列表,这显然会在经常轮询许多文件时导致疯狂大量垃圾收集开销。
我主要通过以下步骤使我的用例工作,而无需求助于使用数据库或编写自定义组件:
- 从父目录
cameldest/queue 上的一个文件使用者移动到对子目录中的文件进行递归排序(cameldest/queue/high/ 在cameldest/queue/low/ 之前)到两个使用者,每个目录一个,没有完全排序。
-
仅设置来自
/cameldest/queue/high/ 的消费者以通过我的实际业务逻辑处理文件。
- 从
/cameldest/queue/low 设置消费者,以简单地将文件从“低”提升到“高”(复制它们,即.to("file://cameldest/queue/high");)
- 至关重要的是,为了仅在 high 不忙时从“low”提升到“high”,将路由策略附加到“high”以 限制其他路由,即“低”,如果有任何消息在“高”中进行中
- 此外,我在“high”中添加了
ThrottlingInflightRoutePolicy,以防止它同时进行过多的交换。
想象一下这就像在机场办理登机手续时,如果商务舱通道是空的,游客会被邀请进入商务舱通道。
这在负载下就像一个魅力,即使数十万个文件在“低”队列中,新消息(文件)直接进入“高”也可以在几秒钟内得到处理。
此解决方案未涵盖的唯一要求是有序性:无法保证较旧的文件首先被拾取,而是随机拾取。可以想象这样一种情况:源源不断的传入文件可能导致某个特定文件 X 总是不走运并且永远不会被拾取。但是,发生这种情况的可能性非常低。
可能的改进: 目前允许/暂停将文件从“低”提升到“高”的阈值设置为“高”中的 0 条正在传输的消息。一方面,这个保证文件落入“高”的文件将在从“低”执行另一次提升之前被处理,另一方面它导致了一些停止-启动模式,尤其是在多线程场景中。不过,这不是一个真正的问题,其表现令人印象深刻。
来源:
我的路线定义:
ThrottlingInflightRoutePolicy trp = new ThrottlingInflightRoutePolicy();
trp.setMaxInflightExchanges(50);
SuspendOtherRoutePolicy sorp = new SuspendOtherRoutePolicy("lowPriority");
from("file://cameldest/queue/low?delay=500&maxMessagesPerPoll=25&preMove=inprogress&delete=true")
.routeId("lowPriority")
.log("Copying over to high priority: ${in.headers."+Exchange.FILE_PATH+"}")
.to("file://cameldest/queue/high");
from("file://cameldest/queue/high?delay=500&maxMessagesPerPoll=25&preMove=inprogress&delete=true")
.routeId("highPriority")
.routePolicy(trp)
.routePolicy(sorp)
.threads(20)
.log("Before: ${in.headers."+Exchange.FILE_PATH+"}")
.delay(2000) // This is where business logic would happen
.log("After: ${in.headers."+Exchange.FILE_PATH+"}")
.stop();
我的SuspendOtherRoutePolicy,像ThrottlingInflightRoutePolicy一样松散构建
public class SuspendOtherRoutePolicy extends RoutePolicySupport implements CamelContextAware {
private CamelContext camelContext;
private final Lock lock = new ReentrantLock();
private String otherRouteId;
public SuspendOtherRoutePolicy(String otherRouteId) {
super();
this.otherRouteId = otherRouteId;
}
@Override
public CamelContext getCamelContext() {
return camelContext;
}
@Override
public void onStart(Route route) {
super.onStart(route);
if (camelContext.getRoute(otherRouteId) == null) {
throw new IllegalArgumentException("There is no route with the id '" + otherRouteId + "'");
}
}
@Override
public void setCamelContext(CamelContext context) {
camelContext = context;
}
@Override
public void onExchangeDone(Route route, Exchange exchange) {
//log.info("Exchange done on route " + route);
Route otherRoute = camelContext.getRoute(otherRouteId);
//log.info("Other route: " + otherRoute);
throttle(route, otherRoute, exchange);
}
protected void throttle(Route route, Route otherRoute, Exchange exchange) {
// this works the best when this logic is executed when the exchange is done
Consumer consumer = otherRoute.getConsumer();
int size = getSize(route, exchange);
boolean stop = size > 0;
if (stop) {
try {
lock.lock();
stopConsumer(size, consumer);
} catch (Exception e) {
handleException(e);
} finally {
lock.unlock();
}
}
// reload size in case a race condition with too many at once being invoked
// so we need to ensure that we read the most current size and start the consumer if we are already to low
size = getSize(route, exchange);
boolean start = size == 0;
if (start) {
try {
lock.lock();
startConsumer(size, consumer);
} catch (Exception e) {
handleException(e);
} finally {
lock.unlock();
}
}
}
private int getSize(Route route, Exchange exchange) {
return exchange.getContext().getInflightRepository().size(route.getId());
}
private void startConsumer(int size, Consumer consumer) throws Exception {
boolean started = super.startConsumer(consumer);
if (started) {
log.info("Resuming the other consumer " + consumer);
}
}
private void stopConsumer(int size, Consumer consumer) throws Exception {
boolean stopped = super.stopConsumer(consumer);
if (stopped) {
log.info("Suspending the other consumer " + consumer);
}
}
}