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

基于PostgreSQL与Node.js实现两阶段提交回滚的方案咨询

PostgreSQL + Node.js 实现跨库数据迁移的两阶段提交回滚

核心逻辑

跨库挪数据最怕出现“半成功”状态——tb2插进去了但tb1没删掉,或者反过来。要避免这种情况,就得用PostgreSQL的两阶段提交(2PC),把操作拆成两步:

  1. 先让两个数据库的事务完成预操作,进入“待命提交”状态
  2. 确认两边都准备妥了再一起提交;只要有一边出问题,就把两边的操作全回滚,保证数据一致性

具体实现步骤

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 05:35:12