【问题标题】:Camel: File consumer component "bites off more than it can chew", pipeline dies from out-of-memory error骆驼:文件消费者组件“咬得比它可以咀嚼的多”,管道因内存不足错误而死
【发布时间】:2017-07-08 08:37:03
【问题描述】:

我在 Camel 中定义了一个类似这样的路由: GET 请求进来,在文件系统中创建一个文件。文件消费者拾取它,从外部 Web 服务获取数据,并通过 POST 将生成的消息发送到其他 Web 服务。

下面的简化代码:

    // Update request goes on queue:
    from("restlet:http://localhost:9191/update?restletMethod=post")
    .routeId("Update via POST")
    [...some magic that defines a directory and file name based on request headers...]
    .to("file://cameldest/queue?allowNullBody=true&fileExist=Ignore")

    // Update gets processed
    from("file://cameldest/queue?delay=500&recursive=true&maxDepth=2&sortBy=file:parent;file:modified&preMove=inprogress&delete=true")
    .routeId("Update main route")
    .streamCaching() //otherwise stuff can't be sent to multiple endpoints
    [...enrich message from some web service using http4 component...]
    .multicast()
        .stopOnException()
        .to("direct:sendUpdate", "direct:dependencyCheck", "direct:saveXML")
    .end();

多播中的三个端点只是将生成的消息发布到其他 Web 服务。

当队列(即文件目录cameldest)相当空时,这一切都很好。文件正在cameldest/<subdir> 中创建,由文件使用者拾取并移动到cameldest/<subdir>/inprogress,并且正在将内容发送到三个传出 POST 没有问题。

但是,一旦传入的请求堆积到大约 300,000 个文件,进度就会减慢,最终管道由于内存不足错误而失败(超出了 GC 开销限制)。

通过增加日志记录,我可以看到文件消费者轮询基本上从不运行,因为它似乎每次都对它看到的所有文件负责,等待它们完成处理,并且只然后开始另一轮投票。除了(我假设)导致资源瓶颈之外,这也干扰了我的排序要求:一旦队列被成千上万条等待处理的消息堵塞,那些天真地被排序更高的新消息 - 如果它们仍然被拾取- 仍在那些已经“开始”的人后面等待。

现在,我尝试了 maxMessagesPerPoll 和 eagerMaxMessagesPerPoll 选项。一开始它们似乎缓解了这个问题,但经过几轮投票后,我仍然发现有数千个文件处于“已启动”的边缘。

唯一有效的方法是使delay 和maxMessages... 的瓶颈变得如此之窄,以至于平均而言,处理完成的速度比文件轮询周期快。

显然,这不是我想要的。我希望我的管道尽可能快地处理文件,但不是更快。我希望文件使用者在路由繁忙时等待。

我犯了一个明显的错误吗?

(我正在使用 XFS 的 Redhat 7 机器上运行稍旧的 Camel 2.14.0,如果这是问题的一部分。)

【问题讨论】:

    标签: java apache-camel enterprise-integration


    【解决方案1】:

    简短的回答是没有答案: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);
            }
        }
    }
    

    【讨论】:

      【解决方案2】:

      除非您确实需要将数据保存为文件,否则我会提出一个替代解决方案。

      从您的 restlet 消费者,将每个请求发送到消息队列应用程序,例如 activemq 或 rabbitmq 或类似的东西。您很快就会在该队列上收到大量消息,但这没关系。

      然后将您的文件使用者替换为队列使用者。这将需要一些时间,但每条消息都应单独处理并发送到您想要的任何地方。我用大约 500 000 条消息测试了 rabbitmq,效果很好。这也应该减少消费者的负担。

      【讨论】:

      • 是的,我们的第一个实现是从 activemq 开始的,但有一个要求我只是顺便提到的,那就是队列必须是有序且唯一的:我将对象的 ID 发送到与优先级标志一起被更新。我希望队列(a)首先返回高优先级对象,并且(b)不将已经排队的对象排入队列。当我在 SO 上询问该用例时,MQ 大神对我非常生气,因为消息队列显然不是队列而是代理。有人告诉我将对象保存在数据库或磁盘上。后者似乎是一个更简单的解决方案。
      • 啊哈,在您的用例中,我会使用 DB,因为这样您就可以根据您的优先级字段进行查询,这更容易,因为消息队列不适用于排序,而且排序可能会很棘手。如果顺序无关紧要,那么 MQ 将是一个很好的用例。
      • 如果您想避免使用数据库,另一个选择可能是仍然将所有内容放在队列中。但是当你消费时,你把所有东西都放在一些键值存储中,比如 hashmap 或类似的东西,键是优先级标志。完成后,下一步可能是先处理具有高优先级标志的键值存储,然后再处理具有低优先级标志的键值存储。
      • 是的,但是如果在我将队列清空到 hashmap 后服务器关闭,我的数据就会丢失。这就是为什么我认为我只使用文件系统的原因:它速度快,根据定义是独一无二的,它显然是持久的,而且我可以使用日期和文件夹轻松排序。它运行良好(并且没有增加数据库的复杂性),直到我在队列中达到 300,000 个对象。
      【解决方案3】:

      尝试在 from 文件端点上将 maxMessagesPerPoll 设置为较低的值,以便每次轮询最多只拾取 X 个文件,这也限制了您在 Camel 应用程序中将拥有的飞行消息的总数。

      您可以在文件组件的 Camel 文档中找到有关该选项的更多信息

      【讨论】:

      • 正如我的帖子所说,我试过了,但仍然内存不足。但我会再试一次,使用更好的日志记录,看看发生了什么。谢谢,克劳斯!
      • 您可以执行 JVM 转储并使用分析器查看是否可以找出占用内存的内容。还要确保为 JVM 配置其内存设置等敏感值。
      • 看起来确实是排序太占用内存了。一旦我关闭排序(或将其限制为仅几个文件),一切正常,内存方面,现在我的队列不再优先。
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2015-08-27
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2013-06-03
      • 2013-05-23
      相关资源
      最近更新 更多