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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 02:35:27