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

Camel SFTP结合Hazelcast幂等仓库配置后仍出现重复消息处理问题问询

问题原因及遗漏配置说明

你遇到的重复消费问题是集群场景下Camel SFTP组件的默认配置不兼容分布式部署导致的,主要遗漏了以下核心配置项:

  • 缺失分布式处理中仓库(InProgressRepository)配置
    Camel SFTP组件默认使用内存级MemoryIdempotentRepository作为处理中文件的标记仓库,仅对当前节点可见。集群多节点同时轮询时,会同时检测到同一文件可用,在分布式幂等检查生效前就会并行启动处理流程。你需要将共用的Hazelcast幂等仓库也配置为处理中仓库,追加配置:
    .inProgressRepository(idempotentRepository)
    你也可以单独创建一个专属的Hazelcast Map存储处理中文件,和幂等校验的存储隔离。

  • readLock集群兼容配置缺失
    仅配置readLock=rename不足以适配集群场景,还需要补充两个关联配置:

    • readLockCheckInterval:设置合理的锁检查间隔(建议≥2000ms),避免多个节点同时发起重命名锁操作产生冲突
    • readLockTimeout:设置大于单文件最大处理时长的超时时间(建议≥30000ms),避免处理未完成时锁提前被释放
      另外需要确认你的SFTP账号拥有文件重命名权限,否则readLock=rename会静默失效,等同于没有配置锁。
  • 可选补充配置(避免版本兼容或异常场景问题)

    • 显式指定幂等Key:SFTP组件默认仅用文件名作为幂等Key,建议用文件名+文件大小+修改时间的组合作为唯一标识,避免重名文件误判:
      .idempotentKey(simple("${file:name}-${file:size}-${file:modified}"))
    • 显式开启预检查模式:配置.idempotentEager(true),确保在文件处理前就完成幂等Key的写入校验,避免处理后再写入幂等Key导致的时间差重复消费。
    • 检查Hazelcast集群状态:确认所有消费节点的Hazelcast实例都加入了同一个集群,配置的ftp-consume-reception-map为分布式共享Map,所有节点读写的是同一份存储数据。如果Hazelcast节点间网络不通,每个节点使用本地独立仓库,幂等校验自然失效。

提示:如果业务场景允许同名文件后续重新消费,可以给Hazelcast的ftp-consume-reception-map配置TTL自动过期时间,或配置.idempotentRemoveOnSuccess(true)在文件处理成功后移除幂等Key。

修正后完整配置示例

IdempotentRepository idempotentRepository = new HazelcastIdempotentRepository(hazelcastInstance, "ftp-consume-reception-map");
// 如需隔离处理中存储,可单独初始化:IdempotentRepository inProgressRepo = new HazelcastIdempotentRepository(hazelcastInstance, "ftp-in-progress-map");

from(sftp(hostname + ":" + port + "/" + inputDirectory)
   .username(username)
   .password(password)
   .delay(delay)
   .maxMessagesPerPoll(maxMessagesPerPoll)
   .move(archiveDirectory)
   .moveFailed(errorDirectory)
   .idempotent(true)
   .idempotentEager(true)
   .idempotentRepository(idempotentRepository)
   .inProgressRepository(idempotentRepository) // 新增分布式处理中标记
   .idempotentKey(simple("${file:name}-${file:size}-${file:modified}")) // 可选 明确幂等唯一Key
   .passiveMode(true)
   .readLock("rename")
   .readLockCheckInterval(2000) // 新增锁检查间隔2s
   .readLockTimeout(30000) // 新增锁超时30s
)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 17:06:08