Camel SFTP多实例部署下如何避免重复处理同一文件?
多实例部署下避免Camel SFTP重复处理文件的方案
你的Camel SFTP路由如下:
from("sftp://userName:password@ip:22/<my-folder>?move=.done") .routeId("my-route-1") .<processing-logic>
针对多实例部署场景,可通过以下几种方式避免重复处理同一文件:
1. 启用SFTP组件自带的独占锁机制
Camel SFTP组件内置多种文件锁策略,读取文件前锁定资源,阻止其他实例访问:
- 配置
readLock参数,常用值包括:changed:通过检查文件修改时间和大小判断是否被锁定,兼容性强;fileLock:利用SFTP服务器原生文件锁(需服务器支持该功能);
- 配合
readLockTimeout设置锁超时时间,避免实例崩溃导致死锁。
示例路由:
from("sftp://userName:password@ip:22/<my-folder>?move=.done&readLock=changed&readLockTimeout=30000") .routeId("my-route-1") .<processing-logic>
2. 预移动文件到临时隔离目录
正式处理前,将文件移动到临时目录并添加唯一标识,让其他实例无法读取该文件:
- 使用
preMove参数,结合实例唯一ID(如主机名、环境变量)或时间戳生成唯一文件名; - 处理完成后,原路由的
move=.done会将文件移至完成目录,即使实例崩溃,后续可清理临时目录残留文件重新处理。
示例路由:
// 假设INSTANCE_ID是每个实例的唯一标识(可通过环境变量注入) from("sftp://userName:password@ip:22/<my-folder>?move=.done&preMove=.inprogress/${file:name}-${INSTANCE_ID}") .routeId("my-route-1") .<processing-logic>
3. 引入分布式锁协调跨实例访问
如果SFTP服务器不支持文件锁,或需要更可靠的跨实例协调,可使用分布式锁(如Redis、ZooKeeper):
- 处理文件前,以文件名作为锁Key获取分布式锁,获取失败则跳过该文件;
- 处理完成后释放锁,确保其他实例可处理其他文件。
示例伪代码:
from("sftp://userName:password@ip:22/<my-folder>?move=.done") .routeId("my-route-1") .process(exchange -> { String fileName = exchange.getIn().getHeader(Exchange.FILE_NAME, String.class); // 调用分布式锁服务获取锁,超时时间设为60秒 boolean lockAcquired = distributedLockService.acquire(fileName, 60000); if (!lockAcquired) { // 停止当前文件的路由处理 exchange.setProperty(Exchange.ROUTE_STOP, true); } }) .<processing-logic> .process(exchange -> { String fileName = exchange.getIn().getHeader(Exchange.FILE_NAME, String.class); distributedLockService.release(fileName); });
4. 用共享存储跟踪已处理文件
维护一个共享的已处理文件列表(如数据库表、Redis集合),处理前先检查文件是否已被处理:
- 利用原子操作保证检查和标记的一致性,避免并发冲突;
- 比如用Redis的
SADD命令,返回1表示文件未被处理,可继续;返回0表示已处理,直接跳过。
示例伪代码:
from("sftp://userName:password@ip:22/<my-folder>?move=.done") .routeId("my-route-1") .process(exchange -> { String fileName = exchange.getIn().getHeader(Exchange.FILE_NAME, String.class); // 检查文件是否已在已处理集合中 boolean isProcessed = redisTemplate.opsForSet().isMember("processed_files", fileName); if (isProcessed) { exchange.setProperty(Exchange.ROUTE_STOP, true); return; } // 原子性添加到已处理集合 boolean added = redisTemplate.opsForSet().add("processed_files", fileName); if (!added) { // 并发情况下被其他实例抢先标记,停止处理 exchange.setProperty(Exchange.ROUTE_STOP, true); } }) .<processing-logic>;
内容的提问来源于stack exchange,提问作者dark ninja
相关产品推荐
相关产品推荐

