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

如何使用node-pg-migrate实现独立的迁移流水线?

基于node-pg-migrate从零构建新系统数据库(含异构旧数据导入)

我来分享一套完整的实现方案,专门解决新旧系统数据结构完全不同时,用node-pg-migrate从零搭建新库并完成旧数据迁移的需求——这类新旧系统切换的场景我帮不少团队落地过,亲测这套流程靠谱。

整体思路

核心逻辑是做分层迁移:先搭一个“导入缓冲层”存旧系统原始数据,再构建新系统的业务结构,最后把缓冲层的数据做映射转换,导入到新系统中。流程拆解下来是这四步:

  • 用node-pg-migrate创建专属的导入模式(schema)
  • 通过旧系统Web API拉取原始数据,写入导入模式的表
  • 用迁移脚本创建新系统的所有业务模式(表、函数、约束等)
  • 编写转换逻辑,把导入层的数据映射到新系统结构中

具体实现步骤

1. 先搞定node-pg-migrate的基础环境

首先得把依赖和配置弄好:

npm install node-pg-migrate pg --save-dev

在项目根目录建migrations/.config.json,填好你的数据库连接信息:

{
  "dev": {
    "driver": "pg",
    "user": "你的数据库用户名",
    "password": "你的数据库密码",
    "host": "localhost",
    "database": "新系统数据库名",
    "port": 5432
  }
}

2. 创建导入模式的迁移脚本

第一个迁移脚本用来建专门存旧数据的schema和临时表,比如migrations/1_create_import_schema.js:

exports.up = async (pgm) => {
  // 创建导入专用的schema,避免和新系统结构混在一起
  pgm.createSchema('import', { ifNotExists: true });

  // 建和旧系统数据结构对应的表(这里以用户表为例)
  pgm.createTable('import.old_users', {
    old_user_id: { type: 'integer', primaryKey: true },
    old_username: { type: 'varchar(100)', notNull: true },
    old_email: { type: 'varchar(255)' },
    created_at_old: { type: 'timestamp' }
  }, { schema: 'import' });

  // 其他旧系统的表也按这个方式创建
};

exports.down = async (pgm) => {
  // 回滚时清理导入层的表和schema
  pgm.dropTable('import.old_users', { ifExists: true });
  pgm.dropSchema('import', { ifExists: true, cascade: true });
};

运行这个迁移,把导入层结构建好:

npx node-pg-migrate up --env dev

3. 从旧系统Web服务拉取并导入原始数据

写个独立的脚本(比如scripts/import-old-data.js)来调用旧系统API,把数据写入导入表:

const { Pool } = require('pg');
const axios = require('axios'); // 用axios调用旧系统API,也可以用其他请求库

// 初始化数据库连接池
const pool = new Pool({
  user: "你的数据库用户名",
  password: "你的数据库密码",
  host: "localhost",
  database: "新系统数据库名",
  port: 5432
});

async function importOldUsers() {
  // 调用旧系统API拿用户数据
  const response = await axios.get('https://旧系统API地址/users');
  const oldUsers = response.data;

  // 批量插入到导入表,重复数据跳过
  const insertQuery = `
    INSERT INTO import.old_users (old_user_id, old_username, old_email, created_at_old)
    VALUES ($1, $2, $3, $4)
    ON CONFLICT (old_user_id) DO NOTHING;
  `;

  // 循环插入(数据量极大的话建议用批量插入优化)
  for (const user of oldUsers) {
    await pool.query(insertQuery, [
      user.id,
      user.username,
      user.email,
      user.created_at
    ]);
  }
}

// 执行导入逻辑
importOldUsers()
  .then(() => console.log('旧系统用户数据导入完成!'))
  .catch(err => console.error('导入失败:', err))
  .finally(() => pool.end());

运行脚本完成数据导入:

node scripts/import-old-data.js

4. 创建新系统的核心业务结构

接下来写迁移脚本搭建新系统的schema、表、函数等,比如migrations/2_create_core_schema.js:

exports.up = async (pgm) => {
  // 创建新系统的核心业务schema
  pgm.createSchema('core', { ifNotExists: true });

  // 建新系统的用户表(和旧结构完全不同)
  pgm.createTable('core.users', {
    user_id: { type: 'uuid', primaryKey: true, default: pgm.func('gen_random_uuid()') },
    full_name: { type: 'varchar(255)', notNull: true },
    email_address: { type: 'varchar(255)', notNull: true, unique: true },
    signup_date: { type: 'timestamp', notNull: true, default: pgm.func('current_timestamp') },
    is_active: { type: 'boolean', default: true }
  }, { schema: 'core' });

  // 新建系统的业务函数示例
  pgm.createFunction(
    'core.get_active_users_count',
    [],
    { returns: 'integer', language: 'plpgsql' },
    `
      BEGIN
        RETURN COUNT(*) FROM core.users WHERE is_active = true;
      END;
    `
  );

  // 其他业务表、约束、索引都按需求添加
};

exports.down = async (pgm) => {
  // 回滚时清理核心业务层
  pgm.dropFunction('core.get_active_users_count', [], { ifExists: true });
  pgm.dropTable('core.users', { ifExists: true });
  pgm.dropSchema('core', { ifExists: true, cascade: true });
};

运行这个迁移,把新系统的基础结构搭好:

npx node-pg-migrate up --env dev

5. 数据转换并迁移到新系统

最后写迁移脚本,把导入层的旧数据做映射转换,导入到新系统表中,比如migrations/3_transform_and_migrate_data.js:

exports.up = async (pgm) => {
  // 转换旧用户数据到新用户表,这里可以加各种清洗逻辑
  pgm.sql(`
    INSERT INTO core.users (full_name, email_address, signup_date)
    SELECT 
      old_username AS full_name,
      old_email AS email_address,
      created_at_old AS signup_date
    FROM import.old_users
    WHERE old_email IS NOT NULL; -- 过滤掉无效的邮箱数据
  `);

  // 其他表的数据转换逻辑也按这个方式写
};

exports.down = async (pgm) => {
  // 回滚时清空新系统的用户表(根据实际需求调整,比如保留初始化数据)
  pgm.sql(`TRUNCATE core.users CASCADE;`);
};

运行这个迁移,完成最终的数据迁移:

npx node-pg-migrate up --env dev

几个关键注意事项

  • 数据校验要做足:导入旧数据后,一定要先校验数据的完整性和合法性,避免脏数据进入新系统——可以在导入脚本或迁移中加校验逻辑。
  • 利用事务保障原子性:node-pg-migrate默认每个迁移脚本都在事务中执行,只要某一步失败,整个迁移会自动回滚,不用担心半吊子状态。
  • 增量迁移的考虑:如果旧系统还在运行,后续可以写增量导入脚本,只拉取新增/修改的数据,再同步到新系统。
  • 清理导入层:确认数据迁移完成且验证无误后,可以删除import模式节省空间,或者保留一段时间作为备份。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 07:56:00