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

如何使用Node.js实现MySQL跨服务器数据传输?求替代RabbitMQ方案

基于Node.js的MySQL跨服务器数据传输简易方案

方案1:直接读写同步(零中间件,最快上手)

适合小批量数据或定时同步场景,无需额外中间件,直接连接两个数据库完成数据迁移。

实现步骤

  1. 安装依赖
npm install mysql2
  1. 核心代码示例
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日志,捕获数据变更(新增/修改/删除),实时同步到目标库,适合需要增量同步的场景。

实现步骤

  1. 开启源库的Binlog:修改MySQL配置文件my.cnf(或my.ini),添加以下配置后重启MySQL
server-id = 1
log-bin = mysql-bin
binlog-format = ROW
  1. 安装依赖
npm install node-mysql-binlog mysql2
  1. 核心代码示例
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),可以通过双实例连接实现同步,代码可读性更强。

实现步骤

  1. 安装依赖
npm install sequelize mysql2
  1. 核心代码示例
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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 11:40:53