双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组件的拉取阶段,避免实例先拉取文件再做判断的延迟问题:
- 保留你已定义的
idempotentRepoBean配置 - 修改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
相关产品推荐
相关产品推荐

