You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在基于文件的输入场景中使用Camel Aggregator?

在Camel文件输入场景中正确使用Aggregator的解决方案

默认情况下,Camel的File组件会在Exchange完成路由后,自动将文件移动到$inputLocation/.camel/目录。但使用Aggregator时,单个Exchange在经过聚合步骤后会继续执行聚合块外的剩余路由逻辑,导致文件被提前移动,后续聚合完成时,处理器无法找到原文件路径(因为文件已经被移走)。

问题示例路由代码:

@Override
public void configure() throws Exception {
  from("file:workdirs/in/")
    .aggregate(constant(true), AggregationStrategies.groupedExchange())
      .completionTimeout(10*1000) // 10秒超时聚合
      .process("doSomethingWithTheAggregatedFiles")
      .end()
  .log("file completed route: $simple{header[CamelFileName]}")
  .end();
}

上述代码中,doSomethingWithTheAggregatedFiles处理器会发现文件已被移动到workdirs/in/.camel/,无法访问原位置的文件。由于输入文件可能较大,不想通过读取文件内容(如ZipAggregationStrategy的方式)来规避问题。


解决方案1:调整路由结构,让聚合逻辑成为路由终点

把聚合块外的后续逻辑(比如日志输出)移到聚合内部,确保单个Exchange不会在聚合完成前执行后续路由,从而避免文件被提前移动。修改后的路由:

@Override
public void configure() throws Exception {
  from("file:workdirs/in/")
    .aggregate(constant(true), AggregationStrategies.groupedExchange())
      .completionTimeout(10*1000)
      .process("doSomethingWithTheAggregatedFiles")
      .log("aggregation completed, processed ${body.size()} files")
      .end();
}

这样,单个Exchange进入聚合后,会等待聚合完成才会执行后续逻辑,文件不会被提前移动,聚合处理器可以正常访问原文件路径。

解决方案2:禁用File组件的自动移动行为,手动管理文件生命周期

通过配置File组件的move参数,关闭自动移动(或设置为延迟执行),等聚合完成后再手动处理文件的移动或删除。

配置示例1:延迟移动到指定目录

from("file:workdirs/in/?move=done/${date:now:yyyyMMdd}/${file:name}")

这里的move配置会在整个路由完成后才移动文件,而不是单个Exchange完成聚合步骤就移动,确保聚合时文件仍在原位置。

配置示例2:完全禁用自动移动,手动处理

from("file:workdirs/in/?move=noop")

使用move=noop后,File组件不会自动移动文件,需要在聚合完成后的处理器中手动处理文件,避免重复读取:

.process(exchange -> {
    List<Exchange> aggregatedExchanges = exchange.getIn().getBody(List.class);
    for (Exchange ex : aggregatedExchanges) {
        String fileName = ex.getIn().getHeader(Exchange.FILE_NAME, String.class);
        File originalFile = new File("workdirs/in/" + fileName);
        File targetFile = new File("workdirs/done/" + fileName);
        // 创建目标目录(如果不存在)
        Files.createDirectories(targetFile.getParentFile().toPath());
        // 移动文件
        Files.move(originalFile.toPath(), targetFile.toPath(), StandardCopyOption.REPLACE_EXISTING);
    }
})

解决方案3:自定义AggregationStrategy保存文件引用

自定义聚合策略,在聚合时直接保存File对象,而不是依赖Exchange中的路径信息,确保聚合处理器能直接访问文件:

public class FileAggregationStrategy implements AggregationStrategy {
    @Override
    public Exchange aggregate(Exchange oldExchange, Exchange newExchange) {
        if (oldExchange == null) {
            // 初始化聚合结果,存储File对象列表
            List<File> fileList = new ArrayList<>();
            File currentFile = newExchange.getIn().getBody(File.class);
            fileList.add(currentFile);
            newExchange.getIn().setBody(fileList);
            return newExchange;
        } else {
            // 追加新的File对象到聚合列表
            List<File> fileList = oldExchange.getIn().getBody(List.class);
            File currentFile = newExchange.getIn().getBody(File.class);
            fileList.add(currentFile);
            oldExchange.getIn().setBody(fileList);
            return oldExchange;
        }
    }
}

路由中使用该策略:

from("file:workdirs/in/")
    .aggregate(constant(true), new FileAggregationStrategy())
      .completionTimeout(10*1000)
      .process("doSomethingWithTheAggregatedFiles")
      .end();

注意:此方案建议结合解决方案2使用,避免文件被自动移动导致File对象路径失效。


内容的提问来源于stack exchange,提问作者ncgisudo

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.23 11:16:03