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

如何在含SFTP与数据库的Apache Camel路由中实现事务

解决方案:基于补偿机制实现SFTP与数据库的一致性

核心问题分析

Camel的transacted()仅支持XA事务兼容的资源(如关系型数据库、JMS),但SFTP属于非事务性资源——文件一旦上传到SFTP服务器,无法通过事务回滚直接撤销,因此必须通过补偿逻辑实现最终一致性。

具体实现方案

方案1:onException+补偿路由回滚已上传文件

思路:捕获数据库保存(步骤4)的异常,触发补偿逻辑删除SFTP上已上传的文件。

示例路由代码:

from("jdbc:myDb?query=SELECT * FROM data_to_process")
    .process(dataGenerator) // 步骤2:生成文件
    .to("sftp://sftp-server/path?fileName=${header.fileName}") // 步骤3:上传SFTP
    .doTry()
        .to("jdbc:myDb?sql=INSERT INTO file_logs (file_name, status) VALUES (?, 'SUCCESS')") // 步骤4:存数据库
    .doCatch(SQLException.class)
        // 补偿逻辑:删除SFTP上的目标文件
        .setHeader("CamelFileName", simple("${header.fileName}"))
        .to("sftp://sftp-server/path?operation=delete")
        .log("补偿执行:SFTP文件 ${header.fileName} 已删除")
    .endDoTry();

方案2:利用Camel补偿事务组件管理流程

通过compensateWith绑定上传操作的补偿逻辑,让Camel在事务失败时自动触发回滚动作。

示例路由代码:

from("jdbc:myDb?query=SELECT * FROM data_to_process")
    .process(dataGenerator)
    // 绑定SFTP上传的补偿路由
    .compensateWith("direct:deleteSftpFile")
    .to("sftp://sftp-server/path?fileName=${header.fileName}")
    .transacted() // 事务仅管理数据库操作
    .to("jdbc:myDb?sql=INSERT INTO file_logs (file_name, status) VALUES (?, 'SUCCESS')")
    .log("文件日志已成功保存");

// 补偿路由:删除SFTP文件
from("direct:deleteSftpFile")
    .setHeader("CamelFileName", simple("${exchangeProperty.CamelFileName}"))
    .to("sftp://sftp-server/path?operation=delete")
    .log("补偿:已删除SFTP文件 ${header.fileName}");

注:此处transacted()仅覆盖数据库插入操作,当数据库操作失败时,Camel会自动触发绑定的补偿路由。

方案3:调整步骤顺序,先写数据库再传SFTP

若业务允许,可反转步骤3和4的顺序,先完成数据库日志插入,再执行SFTP上传:

  • 数据库插入成功后再上传SFTP;
  • 若SFTP上传失败,删除已插入的数据库日志。

示例路由代码:

from("jdbc:myDb?query=SELECT * FROM data_to_process")
    .process(dataGenerator)
    .transacted()
    .to("jdbc:myDb?sql=INSERT INTO file_logs (file_name, status) VALUES (?, 'PENDING')")
    .doTry()
        .to("sftp://sftp-server/path?fileName=${header.fileName}")
        // 更新数据库状态为成功
        .to("jdbc:myDb?sql=UPDATE file_logs SET status='SUCCESS' WHERE file_name=?")
    .doCatch(SftpException.class)
        // 补偿:删除数据库中的待处理日志
        .to("jdbc:myDb?sql=DELETE FROM file_logs WHERE file_name=?")
        .log("补偿:已删除数据库中未完成的文件日志")
    .endDoTry();

额外注意事项

  • 幂等性保障:添加idempotentConsumer,用数据库或文件系统存储已处理的数据ID,防止异常重试时重复生成文件或上传。
  • SFTP权限:确保执行删除操作的SFTP账号拥有目标路径的删除权限,否则补偿逻辑会失效。
  • 异常范围:明确捕获特定异常类型(如SQLException、SftpException),避免误触发补偿。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 13:52:47