如何使用Node.js实现MySQL跨服务器数据传输?求替代RabbitMQ方案
基于Node.js的MySQL跨服务器数据传输简易方案
方案1:直接读写同步(零中间件,最快上手)
适合小批量数据或定时同步场景,无需额外中间件,直接连接两个数据库完成数据迁移。
实现步骤
- 安装依赖
npm install mysql2
- 核心代码示例
const mysql = require('mysql2/promise'); // 源数据库配置 const sourceConfig = { host: '源服务器IP', user: '用户名', password: '密码', database: '源数据库名' }; // 目标数据库配置 const targetConfig = { host: '目标服务器IP', user: '用户名', password: '密码', database: '目标数据库名' }; async function syncData() { // 建立两个数据库连接 const sourceConn = await mysql.createConnection(sourceConfig); const targetConn = await mysql.createConnection(targetConfig); try { // 从源库查询数据(可根据需求加分页,避免一次性拉取过多数据) const [rows] = await sourceConn.execute('SELECT * FROM 源表名 LIMIT 1000'); if (rows.length === 0) { console.log('无数据需要同步'); return; } // 批量插入目标库(根据表结构调整字段) const fields = Object.keys(rows[0]).join(','); const placeholders = rows.map(() => `(${Object.keys(rows[0]).map(() => '?').join(',')})`).join(','); const values = rows.flatMap(row => Object.values(row)); await targetConn.execute(`INSERT INTO 目标表名 (${fields}) VALUES ${placeholders}`, values); console.log(`成功同步${rows.length}条数据`); } catch (err) { console.error('同步失败:', err); } finally { // 关闭连接 await sourceConn.end(); await targetConn.end(); } } // 执行同步 syncData();
注意事项
- 数据量大时需添加分页逻辑(比如按ID分段查询),避免内存溢出
- 可配合
node-schedule实现定时自动同步 - 建议在事务中执行插入操作,保证数据一致性
方案2:Binlog增量同步(实时同步场景)
通过监听源MySQL的Binlog日志,捕获数据变更(新增/修改/删除),实时同步到目标库,适合需要增量同步的场景。
实现步骤
- 开启源库的Binlog:修改MySQL配置文件
my.cnf(或my.ini),添加以下配置后重启MySQL
server-id = 1 log-bin = mysql-bin binlog-format = ROW
- 安装依赖
npm install node-mysql-binlog mysql2
- 核心代码示例
const Binlog = require('node-mysql-binlog'); const mysql = require('mysql2/promise'); // 源数据库配置(需拥有REPLICATION SLAVE权限的账号) const binlog = new Binlog({ host: '源服务器IP', user: '用户名', password: '密码', database: '源数据库名' }); // 目标数据库连接 const targetConfig = { host: '目标服务器IP', user: '用户名', password: '密码', database: '目标数据库名' }; const targetConn = await mysql.createConnection(targetConfig); // 监听指定表的变更事件 binlog.addListener('event', async (event) => { if (event.table !== '源表名') return; switch (event.type) { case 'INSERT': const insertFields = Object.keys(event.rows[0]).join(','); const insertValues = Object.values(event.rows[0]); await targetConn.execute(`INSERT INTO 目标表名 (${insertFields}) VALUES (?)`, [insertValues]); break; case 'UPDATE': const updateFields = Object.keys(event.rows[0].after).map(key => `${key}=?`).join(','); const updateValues = [...Object.values(event.rows[0].after), event.rows[0].before.id]; await targetConn.execute(`UPDATE 目标表名 SET ${updateFields} WHERE id=?`, updateValues); break; case 'DELETE': await targetConn.execute(`DELETE FROM 目标表名 WHERE id=?`, [event.rows[0].before.id]); break; } }); // 启动监听 binlog.start();
注意事项
- 源库账号需要授予
REPLICATION SLAVE和REPLICATION CLIENT权限 - 需确保目标表结构与源表一致
- 可添加异常重试机制,避免单次同步失败导致数据丢失
方案3:ORM封装同步(代码更优雅)
如果项目已经使用ORM(如Sequelize),可以通过双实例连接实现同步,代码可读性更强。
实现步骤
- 安装依赖
npm install sequelize mysql2
- 核心代码示例
const { Sequelize, Model, DataTypes } = require('sequelize'); // 源库Sequelize实例 const sourceSeq = new Sequelize('源数据库名', '用户名', '密码', { host: '源服务器IP', dialect: 'mysql' }); // 目标库Sequelize实例 const targetSeq = new Sequelize('目标数据库名', '用户名', '密码', { host: '目标服务器IP', dialect: 'mysql' }); // 定义源库模型(需与表结构一致) class SourceModel extends Model {} SourceModel.init({ id: { type: DataTypes.INTEGER, primaryKey: true, autoIncrement: true }, name: DataTypes.STRING, email: DataTypes.STRING }, { sequelize: sourceSeq, modelName: '表名' }); // 目标库模型(复用源库模型定义) class TargetModel extends Model {} TargetModel.init(SourceModel.rawAttributes, { sequelize: targetSeq, modelName: '表名' }); async function syncData() { // 查询源库数据 const dataList = await SourceModel.findAll({ limit: 1000 }); // 批量创建目标库数据 await TargetModel.bulkCreate(dataList.map(item => item.toJSON()), { ignoreDuplicates: true }); console.log(`同步完成,共${dataList.length}条数据`); } syncData();
注意事项
bulkCreate的ignoreDuplicates参数可避免重复插入- 可利用Sequelize的事务功能保证数据一致性
- 适合已有ORM项目的无缝集成
内容的提问来源于stack exchange,提问作者Prakash Kumar Gupta
相关产品推荐
相关产品推荐

