如何实现Oracle与MongoDB跨数据源的"全有或全无"分布式事务
实现方案说明
方案适配性
你的业务思路完全符合分布式事务的适用场景,核心逻辑可以通过JTA+XA的标准化分布式事务实现,不需要自行封装两阶段提交逻辑。
前置依赖
- Oracle侧:使用支持XA的JDBC驱动(ojdbc8及以上版本),确保数据库用户拥有XA事务相关权限
- MongoDB侧:版本≥4.0,部署形态为副本集或分片集群,已开启分布式事务能力
- Java侧:集成JTA事务管理器(推荐用Atomikos,轻量且支持分布式环境部署)
核心实现逻辑
全局事务管理器会将Oracle数据源、MongoDB数据源都注册为XA资源,所有操作均在全局事务上下文中执行:
- 全局事务开启后,先从Oracle的
sourcetable读取待迁移数据,加行锁避免其他进程重复读取 - 执行Oracle侧的两个存储过程,解析输出游标生成对应的Mongo文档
- 在MongoDB分布式事务上下文内向
collection1、collection2写入对应文档 - 更新Oracle侧
sourcetable对应记录的迁移状态为成功 - 所有操作无异常的情况下,由事务管理器统一触发两阶段提交:先向所有XA资源发起预提交请求,所有资源返回就绪后再发起正式提交
- 任意步骤抛出异常时,事务管理器直接向所有XA资源发起回滚请求,Oracle侧的存储过程执行、状态更新会回滚,MongoDB侧的文档写入也会完全回滚
伪代码示例
import com.atomikos.icatch.jta.UserTransactionImp; import javax.transaction.UserTransaction; import oracle.jdbc.xa.client.OracleXADataSource; import com.mongodb.client.MongoClients; import com.mongodb.client.MongoClient; import com.mongodb.client.MongoCollection; import org.bson.Document; import java.sql.*; public class DataMigrationService { // 初始化Oracle XA数据源 private OracleXADataSource initOracleXADataSource() throws SQLException { OracleXADataSource xaDataSource = new OracleXADataSource(); xaDataSource.setURL("jdbc:oracle:thin:@//oracle-host:1521/service-name"); xaDataSource.setUser("username"); xaDataSource.setPassword("password"); return xaDataSource; } // 初始化MongoDB客户端 private MongoClient initMongoClient() { return MongoClients.create("mongodb://mongo-host1:27017,mongo-host2:27017/?replicaSet=rs0"); } // 游标转Mongo文档的自定义方法,可根据业务逻辑调整 private Document buildDocFromResultSet(ResultSet rs) throws SQLException { Document doc = new Document(); // 字段映射逻辑 return doc; } public void migrateData(Long sourceRecordId) throws Exception { UserTransaction globalTransaction = new UserTransactionImp(); // 设置全局事务超时时间,根据业务执行时长调整 globalTransaction.setTransactionTimeout(300); OracleXADataSource oracleXADataSource = initOracleXADataSource(); MongoClient mongoClient = initMongoClient(); Connection oracleConn = null; try { // 开启全局事务 globalTransaction.begin(); // 获取Oracle连接,自动加入当前全局事务 oracleConn = oracleXADataSource.getXAConnection().getConnection(); oracleConn.setAutoCommit(false); // 步骤1:读取sourcetable对应记录,加行锁 String selectSql = "SELECT * FROM sourcetable WHERE id = ? FOR UPDATE"; PreparedStatement selectStmt = oracleConn.prepareStatement(selectSql); selectStmt.setLong(1, sourceRecordId); ResultSet sourceRs = selectStmt.executeQuery(); if (!sourceRs.next()) { throw new RuntimeException("待迁移记录不存在"); } // 步骤2:执行两个存储过程,处理游标生成Mongo文档 // 执行存储过程1 CallableStatement proc1Stmt = oracleConn.prepareCall("{call proc1(?, ?)}"); proc1Stmt.setLong(1, sourceRecordId); proc1Stmt.registerOutParameter(2, OracleTypes.CURSOR); proc1Stmt.execute(); ResultSet proc1Rs = (ResultSet) proc1Stmt.getObject(2); Document doc1 = buildDocFromResultSet(proc1Rs); // 执行存储过程2 CallableStatement proc2Stmt = oracleConn.prepareCall("{call proc2(?, ?)}"); proc2Stmt.setLong(1, sourceRecordId); proc2Stmt.registerOutParameter(2, OracleTypes.CURSOR); proc2Stmt.execute(); ResultSet proc2Rs = (ResultSet) proc2Stmt.getObject(2); Document doc2 = buildDocFromResultSet(proc2Rs); // 写入MongoDB两个集合,操作加入全局事务 var mongoDb = mongoClient.getDatabase("target_db"); MongoCollection<Document> coll1 = mongoDb.getCollection("collection1"); MongoCollection<Document> coll2 = mongoDb.getCollection("collection2"); coll1.insertOne(doc1); coll2.insertOne(doc2); // 步骤3:更新Oracle源表状态 String updateSql = "UPDATE sourcetable SET status = '迁移成功' WHERE id = ?"; PreparedStatement updateStmt = oracleConn.prepareStatement(updateSql); updateStmt.setLong(1, sourceRecordId); updateStmt.executeUpdate(); // 全局提交 globalTransaction.commit(); } catch (Exception e) { // 任意异常触发全局回滚 if (globalTransaction != null) { try { globalTransaction.rollback(); } catch (Exception rollbackEx) { // 记录回滚异常,后续人工兜底处理 System.err.printf("全局事务回滚失败,待处理记录id:%s%n", sourceRecordId); rollbackEx.printStackTrace(); } } throw e; } finally { if (oracleConn != null) { oracleConn.close(); } if (mongoClient != null) { mongoClient.close(); } } } }
注意事项
- 需添加幂等校验逻辑,避免重试时重复写入数据
- 回滚失败的记录建议落盘到异常表,新增定时任务做兜底重试
- 全局事务超时时间需大于单次迁移的最大耗时,避免事务提前中断
内容的提问来源于stack exchange,提问作者Mohan
相关产品推荐
相关产品推荐

