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

Node.js通过SSH连接MySQL时stream查询多次触发问题如何解决?

问题原因分析
  • 核心触发原因是重复使用了全局的sshClient实例,且没有在每次逻辑执行完后移除ready事件监听:每次调用新增终端逻辑时,都会给同一个sshClient追加新的ready事件回调,下次触发ready时,所有历史注册的回调都会依次执行,会把之前的终端数据再次插入到数据库中,就会出现删除的旧终端自动恢复的问题。
  • 其次是在数据库插入、更新操作完成后,没有主动关闭MySQL连接和SSH隧道资源,残留的连接上下文会被后续操作复用,也可能导致重复执行历史SQL。
  • 当前直接拼接SQL字符串的写法存在SQL注入风险,也建议同步修改。
修复方案

按照以下步骤修改即可解决问题:

  1. 每次执行新增终端逻辑时,创建独立的sshClient实例,避免事件监听重复绑定
  2. 数据库操作完成后主动关闭MySQL连接、SSH流和SSH客户端
  3. 改用参数化查询替换SQL字符串拼接
  4. 逻辑执行出错时也要主动销毁资源,避免内存泄漏

修改后的代码参考:

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 17:36:01