Camel独立应用如何处理所有匹配文件后自动退出?
处理完目录中所有匹配文件后让Camel独立应用退出的解决方案
核心问题分析
你当前使用的-dm 1(即--durationMax 1)参数是强制Camel在处理1条消息后立即停止上下文,但目录中的每个文件对应一条消息,所以会出现只处理一个文件就退出的情况。要实现处理完所有匹配文件后退出,需要替换这个参数,改用以下几种方案:
方案1:使用IdleTimeout空闲超时机制(推荐)
这个方案让Camel持续监听目录,当目录中没有新文件且达到设定的空闲时间后,自动停止上下文,能覆盖静态文件批量处理和任务执行期间动态新增文件的场景。
- 移除启动参数中的
-dm 1。 - 在文件消费路由中配置
idleTimeout参数,设置空闲超时时间(单位:毫秒),并配合onCompletion触发上下文停止。
Java DSL示例
import org.apache.camel.builder.RouteBuilder; import org.apache.camel.CamelContext; import org.apache.camel.impl.DefaultCamelContext; public class FileProcessorApp { public static void main(String[] args) throws Exception { CamelContext context = new DefaultCamelContext(); context.addRoutes(new RouteBuilder() { @Override public void configure() throws Exception { from("file:///your/target/directory?fileName=TestFile.*&delete=true&idleTimeout=30000") .process(exchange -> { // 这里写你的文件处理逻辑 String fileName = exchange.getIn().getHeader("CamelFileName", String.class); System.out.println("处理文件:" + fileName); }) .onCompletion() .process(exchange -> { // 空闲超时触发时停止Camel上下文 exchange.getContext().stop(); }); } }); context.start(); // 阻塞主线程直到上下文停止 context.getExecutorServiceManager().awaitTermination(); } }
配置说明
idleTimeout=30000:表示目录连续30秒没有新文件进入时,触发空闲事件。delete=true:处理完文件后删除(根据需求调整为move等其他策略)。onCompletion:捕获空闲超时的完成事件,调用stop()终止上下文。
方案2:一次性拉取所有匹配文件批量处理
如果确定任务执行时目录中的文件是静态的(不会动态新增),可以用pollEnrich一次性拉取所有匹配文件,拆分后逐个处理,完成后直接停止上下文。
Java DSL示例
from("timer:startOnce?repeatCount=1") // 仅触发一次路由 .pollEnrich("file:///your/target/directory?fileName=TestFile.*&delete=true&maxMessagesPerPoll=-1") .split(body()) // 拆分拉取到的所有文件消息 .process(exchange -> { // 单个文件处理逻辑 String fileName = exchange.getIn().getHeader("CamelFileName", String.class); System.out.println("处理文件:" + fileName); }) .end() // 结束拆分流程 .process(exchange -> { // 所有文件处理完成后停止上下文 exchange.getContext().stop(); });
配置说明
timer:startOnce?repeatCount=1:确保路由只执行一次,拉取所有文件。maxMessagesPerPoll=-1:表示一次拉取目录中所有匹配的文件(默认是10)。
方案3:自定义剩余文件检查逻辑
通过代码手动检查目录中是否还有未处理的匹配文件,若没有则停止上下文。适合对文件处理时机有精确控制的场景。
Java DSL示例
import java.io.File; public class FileProcessorApp { private static final String TARGET_DIR = "/your/target/directory"; private static final File DIRECTORY = new File(TARGET_DIR); public static void main(String[] args) throws Exception { CamelContext context = new DefaultCamelContext(); context.addRoutes(new RouteBuilder() { @Override public void configure() throws Exception { from("file://" + TARGET_DIR + "?fileName=TestFile.*&delete=true") .process(exchange -> { // 文件处理逻辑 String fileName = exchange.getIn().getHeader("CamelFileName", String.class); System.out.println("处理文件:" + fileName); }) .process(exchange -> { // 检查目录中是否还有匹配的文件 File[] remainingFiles = DIRECTORY.listFiles((dir, name) -> name.startsWith("TestFile.")); if (remainingFiles == null || remainingFiles.length == 0) { exchange.getContext().stop(); } }); } }); context.start(); context.getExecutorServiceManager().awaitTermination(); } }
注意事项
- 若有其他进程在任务执行期间向目录写入文件,这种方法可能会提前停止,导致新文件未被处理,因此更适合静态文件场景。
内容的提问来源于stack exchange,提问作者Makaque
相关产品推荐
相关产品推荐

