基于NodeJS API实现SQL Server表数据迁移至MySQL的问题排查
问题:SQL Server数据同步到MySQL时列值为null导致写入失败
我正在尝试通过NodeJS API将SQL Server中的表数据复制到MySQL,计划使用setInterval实现约30分钟执行一次的定时任务。但写入MySQL时出现问题:列值返回null,导致无数据写入并抛出错误,然而使用console.log输出该数据时,能正常获取SQL Server的整张表数据。两端数据库的表结构(含数据类型)完全一致,以下是API的完整代码:
import { createPool } from "mysql2"; import pkg from "mssql"; const { connect, query } = pkg; const sql = pkg; // MySQL Server export const db = createPool({ host: "IP Address", user: "root", password: "Password", database: "demo", port: "3306", connectionLimit: 300, waitForConnections: true, multipleStatements: true, }); db.getConnection((err, connection) => { if (err) throw err; console.log("Database connected successfully"); connection.release(); }); // SQL Server const config = { user: "user", password: "Password", server: "Server Name", database: "demo", port: 1433, options: { trustedConnection: true, encrypt: true, enableArithAbort: true, trustServerCertificate: true, }, }; // Transfer function const Transfer = () => { sql.connect(config, function (err) { if (err) console.log(err); let sqlRequest = new sql.Request(); let sqlQuery = "Select * From access"; sqlRequest.query(sqlQuery, function (err, data) { if (err) console.log(err); if (data) { db.query( "INSERT IGNORE INTO attendance ID = ?, dateTime = ?, date = ?, time = ?, direction = ?, deviceName = ?, deviceSerial = ?, personName = ?, cardNumber = ?", [ data.userID, data.dateTime, data.date, data.time, data.direction, data.deviceName, data.deviceSerial, data.personName, data.cardNumber, ], (err) => { if (err) { console.log(err); } else { return; } } ); } }); }); }; setInterval(Transfer, 2000); const port = 3001; app.listen(port, () => console.log(`Listening on port ${port}`));
解决方案
1. 修复mssql查询结果的读取方式
mssql的query方法返回的data对象中,实际表数据存储在data.recordset数组里,而非直接挂载在data本身。你之前直接取data.userID会得到undefined,导致MySQL插入时传入null。需要遍历recordset中的每一条数据执行插入操作。
2. 修正MySQL INSERT语句的语法
你的INSERT语句写法错误,正确的INSERT IGNORE语法需要先指定列名,再通过VALUES子句传入参数,而非列名=?的形式。正确格式:
INSERT IGNORE INTO attendance (ID, dateTime, date, time, direction, deviceName, deviceSerial, personName, cardNumber) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
3. 优化异步逻辑(使用async/await避免回调地狱)
原代码嵌套过多回调,容易引发异步时序问题,改用async/await可让逻辑更清晰,同时便于捕获和处理错误。
4. 补充缺失的express依赖
代码中app变量未定义,需要添加express的导入和初始化。
修复后的完整代码
import { createPool } from "mysql2/promise"; // 使用promise版本适配async/await import pkg from "mssql"; import express from "express"; // 导入express const { connect } = pkg; const sql = pkg; const app = express(); // 初始化app实例 // MySQL Server export const db = createPool({ host: "IP Address", user: "root", password: "Password", database: "demo", port: 3306, // 端口改用数字类型 connectionLimit: 300, waitForConnections: true, multipleStatements: true, }); // 测试MySQL连接 (async () => { try { const connection = await db.getConnection(); console.log("Database connected successfully"); connection.release(); } catch (err) { console.error("MySQL连接失败:", err); } })(); // SQL Server const config = { user: "user", password: "Password", server: "Server Name", database: "demo", port: 1433, options: { trustedConnection: true, encrypt: true, enableArithAbort: true, trustServerCertificate: true, }, }; // Transfer函数改用async/await重构 const Transfer = async () => { let sqlServerConnection; try { // 连接SQL Server sqlServerConnection = await connect(config); const sqlRequest = new sql.Request(sqlServerConnection); const sqlQuery = "Select * From access"; const result = await sqlRequest.query(sqlQuery); // 遍历查询结果,逐条插入MySQL for (const row of result.recordset) { await db.query( `INSERT IGNORE INTO attendance (ID, dateTime, date, time, direction, deviceName, deviceSerial, personName, cardNumber) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`, [ row.userID, // 确保SQL Server表中列名与此处匹配,若主键列名为ID则改为row.ID row.dateTime, row.date, row.time, row.direction, row.deviceName, row.deviceSerial, row.personName, row.cardNumber, ] ); } console.log("数据同步完成"); } catch (err) { console.error("同步过程出错:", err); } finally { // 关闭SQL Server连接 if (sqlServerConnection) { await sqlServerConnection.close(); } } }; // 设置30分钟定时任务(1800000毫秒) setInterval(Transfer, 1800000); const port = 3001; app.listen(port, () => console.log(`Listening on port ${port}`));
额外注意事项
- 确保SQL Server查询结果的列名与你插入时使用的
row.xxx完全匹配,比如若SQL Server主键列名为ID,需对应修改为row.ID。 - 使用
mysql2/promise版本可更好地配合async/await,避免回调嵌套问题。 - 定时任务已改为30分钟(1800000毫秒),原代码中的2000毫秒为测试用,确认功能正常后保留此设置即可。
- 新增的错误捕获逻辑可避免单个同步任务失败导致整个进程崩溃。
内容的提问来源于stack exchange,提问作者cameronErasmus
相关产品推荐
相关产品推荐

