Node.js通过SSH连接MySQL时stream查询多次触发问题如何解决?
问题原因分析
- 核心触发原因是重复使用了全局的
sshClient实例,且没有在每次逻辑执行完后移除ready事件监听:每次调用新增终端逻辑时,都会给同一个sshClient追加新的ready事件回调,下次触发ready时,所有历史注册的回调都会依次执行,会把之前的终端数据再次插入到数据库中,就会出现删除的旧终端自动恢复的问题。 - 其次是在数据库插入、更新操作完成后,没有主动关闭MySQL连接和SSH隧道资源,残留的连接上下文会被后续操作复用,也可能导致重复执行历史SQL。
- 当前直接拼接SQL字符串的写法存在SQL注入风险,也建议同步修改。
修复方案
按照以下步骤修改即可解决问题:
- 每次执行新增终端逻辑时,创建独立的sshClient实例,避免事件监听重复绑定
- 数据库操作完成后主动关闭MySQL连接、SSH流和SSH客户端
- 改用参数化查询替换SQL字符串拼接
- 逻辑执行出错时也要主动销毁资源,避免内存泄漏
修改后的代码参考:
const SSHConnection = new Promise(async (resolve, reject) => { // 每次创建新的sshClient实例,不复用全局实例 const sshClient = new require('ssh2').Client(); let connection = null; let stream = null; // 统一资源清理逻辑 const cleanUp = () => { try { connection?.end(); } catch(e) {} try { stream?.destroy(); } catch(e) {} try { sshClient.end(); } catch(e) {} }; sshClient.on('ready', async () => { sshClient.forwardOut( forwardConfig.srcHost, forwardConfig.srcPort, forwardConfig.dstHost, forwardConfig.dstPort, async (err, createdStream) => { stream = createdStream; if (err) { cleanUp(); return reject(err); } const updatedDbServer = { ...dbServer, stream }; connection = mysql.createConnection(updatedDbServer); connection.connect(async (error) => { if (error) { cleanUp(); return reject(error); } try { // 改用参数化查询,避免SQL注入 await connection.promise().query( "INSERT INTO `mandator_1`.`terminal` (`serialNumber`, `name`, `isActive`, `status`, `profileId`, `businessId`, `divisionId`, `ipAddress`, `groupId`, `setupVersion`, `macAddress`) VALUES (?, ?, '1', NULL, '1', '1', NULL, '172.45.17.197', '1', NULL, ?)", [terminal.serial, terminal.serial, terminal.macAddress] ); await connection.promise().query( "UPDATE `terminal` SET `status` = UNIX_TIMESTAMP(NOW()) * 1000, `setupVersion` = '005', `hardwareType` = '10', `update` = 0, `lastUpdated` = UNIX_TIMESTAMP(NOW()) * 1000, `dataReceived` = UNIX_TIMESTAMP(NOW()) * 1000 WHERE `setupVersion` IS NULL AND `hardwareType` IS NULL AND `lastUpdated` IS NULL AND `dataReceived` IS NULL AND `serialNumber` = ?", [terminal.serial] ); resolve(); } catch (e) { reject(e); } finally { // 操作完成后主动清理所有资源 cleanUp(); } }); }); }).on('error', err => { cleanUp(); reject(err); }).connect(tunnelConfig); });
内容的提问来源于stack exchange,提问作者Sithys
相关产品推荐
相关产品推荐

