基于PostgreSQL与Node.js实现两阶段提交回滚的方案咨询
PostgreSQL + Node.js 实现跨库数据迁移的两阶段提交回滚
核心逻辑
跨库挪数据最怕出现“半成功”状态——tb2插进去了但tb1没删掉,或者反过来。要避免这种情况,就得用PostgreSQL的两阶段提交(2PC),把操作拆成两步:
- 先让两个数据库的事务完成预操作,进入“待命提交”状态
- 确认两边都准备妥了再一起提交;只要有一边出问题,就把两边的操作全回滚,保证数据一致性
具体实现步骤
1. 安装依赖
使用Node.js官方PostgreSQL驱动pg:
npm install pg
2. 代码实现
先配置两个数据库的连接池,再按2PC流程执行操作:
const { Pool } = require('pg'); // 初始化两个数据库连接池 const db1Pool = new Pool({ user: '你的数据库账号', host: 'localhost', // 替换为实际数据库地址 database: 'db1', password: '你的数据库密码', port: 5432, }); const db2Pool = new Pool({ user: '你的数据库账号', host: 'localhost', database: 'db2', password: '你的数据库密码', port: 5432, }); // 生成唯一事务ID,标记本次2PC的两个事务 function generateTxId() { return `tx_${Date.now()}_${Math.random().toString(36).slice(2, 10)}`; } async function migrateRow(rowId) { const txId = generateTxId(); let db1Client, db2Client; let db2Prepared = false; let db1Prepared = false; try { // 从连接池获取客户端 db1Client = await db1Pool.connect(); db2Client = await db2Pool.connect(); // 第一步:先在db2插入数据(先插后删,避免数据丢失) await db2Client.query('BEGIN'); // 查询tb1目标行并加行锁,防止其他操作修改该行 const { rows } = await db1Client.query('SELECT * FROM tb1 WHERE id = $1 FOR UPDATE', [rowId]); if (rows.length === 0) { throw new Error('要迁移的行不存在'); } const targetRow = rows[0]; // 替换为tb2的实际字段 await db2Client.query('INSERT INTO tb2 (id, col1, col2) VALUES ($1, $2, $3)', [targetRow.id, targetRow.col1, targetRow.col2]); // 第二步:在db1执行删除预操作 await db1Client.query('BEGIN'); await db1Client.query('DELETE FROM tb1 WHERE id = $1', [rowId]); // 第三步:将两个事务置为准备状态(此步骤后数据不会因事务中断丢失) await db2Client.query(`PREPARE TRANSACTION '${txId}_db2'`); db2Prepared = true; await db1Client.query(`PREPARE TRANSACTION '${txId}_db1'`); db1Prepared = true; // 第四步:正式提交两个已准备的事务 await db2Client.query(`COMMIT PREPARED '${txId}_db2'`); await db1Client.query(`COMMIT PREPARED '${txId}_db1'`); console.log('数据迁移成功!'); } catch (err) { console.error('操作失败,开始回滚:', err.message); // 回滚已准备的事务 if (db2Prepared) { await db2Client.query(`ROLLBACK PREPARED '${txId}_db2'`).catch(e => console.error('db2回滚失败:', e)); } if (db1Prepared) { await db1Client.query(`ROLLBACK PREPARED '${txId}_db1'`).catch(e => console.error('db1回滚失败:', e)); } // 回滚未进入准备状态的临时事务 if (db2Client && !db2Prepared) { await db2Client.query('ROLLBACK').catch(e => console.error('db2临时事务回滚失败:', e)); } if (db1Client && !db1Prepared) { await db1Client.query('ROLLBACK').catch(e => console.error('db1临时事务回滚失败:', e)); } throw err; } finally { // 释放客户端回连接池,避免连接泄漏 if (db1Client) db1Client.release(); if (db2Client) db2Client.release(); } } // 调用示例:迁移id为1的行 moveRow(1).catch(err => console.error('最终失败:', err));
关键注意事项
- 事务ID唯一性:每次生成的ID必须唯一,避免与历史准备事务冲突
- 行锁机制:查询tb1时用
FOR UPDATE加行锁,防止迁移过程中其他操作修改目标行 - 异常兜底:回滚操作本身可能失败,需单独捕获处理,避免连锁报错
- 连接池清理:无论操作成功与否,都要将客户端放回连接池,防止连接耗尽
- 崩溃恢复:PostgreSQL的准备事务会持久化到磁盘,若服务崩溃,重启后可通过
SELECT * FROM pg_prepared_xacts查看未处理事务,手动提交或回滚
内容的提问来源于stack exchange,提问作者Nazmul Haque
相关产品推荐
相关产品推荐

