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

双Camel服务从同一SFTP取文件的并发重复处理问题求解

解决Camel多实例SFTP下载重复处理文件的问题

针对你遇到的多实例同时处理同一SFTP文件、出现No such file exception的问题,以下是几种可行的解决方案:

方案1:利用SFTP原子重命名实现分布式锁

SFTP服务器的文件重命名操作是原子性的,我们可以通过先将待处理文件移动到临时目录的方式,确保只有一个实例能获取到文件处理权限:

修改SFTP Endpoint配置,添加move参数指定临时目录,同时开启自动创建目录的选项:

@Override
public void configure() {
    from("sftp?preferredAuthentications=password&password=${password}&scheduler=quartz&scheduler.cron=${cron}&useUserKnownHostsFile=false&delete=true&shuffle=true&move=.processing/${file:name}&moveCreateFolder=true")
            .transacted()
            .routeId(RouteDownloadResponses.class.getSimpleName())
            .convertBodyTo(byte[].class)
            .andSoOn();
}

原理:实例拉取文件时会先原子性地将原文件移动到.processing目录,只有成功完成移动的实例才能继续处理文件;其他实例会因为原文件已不存在而跳过。处理完成后,delete=true会自动删除临时目录下的文件。

方案2:将幂等判断提前到SFTP Consumer层面

把幂等校验逻辑从路由内移到SFTP组件的拉取阶段,避免实例先拉取文件再做判断的延迟问题:

  1. 保留你已定义的idempotentRepo Bean配置
  2. 修改SFTP URL,添加幂相关参数,同时移除路由内的idempotentConsumer配置:
@Override
public void configure() {
    from("sftp?preferredAuthentications=password&password=${password}&scheduler=quartz&scheduler.cron=${cron}&useUserKnownHostsFile=false&delete=true&shuffle=true&idempotent=true&idempotentRepository=#idempotentRepo&idempotentKey=${header.camelFileAbsolutePath}&eager=true")
            .transacted()
            .routeId(RouteDownloadResponses.class.getSimpleName())
            .convertBodyTo(byte[].class)
            .andSoOn();
}

原理:SFTP Consumer在拉取文件前,会先通过共享的数据库幂等仓库检查文件绝对路径是否已被处理。已处理过的文件会直接跳过,未处理的文件会被标记为已处理后再进入路由。eager=true确保拉取时就完成标记,避免多实例竞争。

方案3:引入分布式锁(如Redis锁)

如果上述方案不适用,可以通过分布式锁强制控制文件的处理权限:

以Redis锁为例(需引入Camel Redis组件):

@Override
public void configure() {
    from("sftp?preferredAuthentications=password&password=${password}&scheduler=quartz&scheduler.cron=${cron}&useUserKnownHostsFile=false&delete=true&shuffle=true")
            .transacted()
            .setHeader("lockKey", simple("sftp-file-lock:${header.camelFileAbsolutePath}"))
            // 获取锁,超时5分钟自动释放
            .to("redis:lock?command=LOCK&key=${header.lockKey}&timeout=300000")
            .choice()
                .when(header("CamelRedisLockSuccess").isEqualTo(true))
                    .convertBodyTo(byte[].class)
                    .andSoOn()
                    // 处理完成后释放锁
                    .to("redis:lock?command=UNLOCK&key=${header.lockKey}")
                .otherwise()
                    .log("Skipping file ${header.camelFileName}, lock not acquired")
            .endChoice()
            .routeId(RouteDownloadResponses.class.getSimpleName());
}

原理:每个实例处理文件前先尝试获取对应文件的分布式锁,只有拿到锁的实例才能继续处理;获取失败的实例直接跳过该文件,避免重复处理。

关键注意事项

  • 所有服务实例必须共享同一个幂等仓库(如同一数据库)或分布式锁资源,否则分布式控制逻辑不生效。
  • 确保SFTP服务器支持原子重命名操作(主流OpenSSH、VSFTPD等均支持),这是方案1的核心前提。
  • 结合transacted()使用时,需保证幂等仓库或锁操作能参与到事务中,避免异常场景下的状态不一致。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 00:42:14