如何使用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
相关产品推荐
相关产品推荐

