如何在基于文件的输入场景中使用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

