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

循环中异步函数致多OPCUA设备连接会话创建失败求助

问题描述

我将从MySQL查询生成的多维数组用于存储多台OPCUA设备的连接参数,怀疑是循环内的异步函数引发了问题,但不知如何修复。运行脚本时无会话启动,返回如下错误:

myprj/node-opcua-assert/dist/index.js:11
        const err = new Error(message);
                    ^
Error
    at assert (myprj/node-opcua-assert/dist/index.js:11:21)
    at OPCUAClientImpl._createSession (myprj/node-opcua-client/dist/private/opcua_client_impl.js:856:40)
    at OPCUAClientImpl.createSession (myprj/node-opcua-client/dist/private/opcua_client_impl.js:273:14)
    at OPCUAClientImpl.createSession (myprj/thenify/index.js:65:46)
    at OPCUAClientImpl.createSession (myprj/node-opcua-client/dist/private/opcua_client_impl.js:267:25)
    at myprj/thenify/index.js:72:10
    at new Promise (<anonymous>)
    at OPCUAClientImpl.createSession (myprj/thenify/index.js:70:12)
    at myprj/read_node.js:8

相关代码如下:

/*
    此处省略导入配置及数据库连接设置的代码
*/

const endpointUrl = config.endpointUrl;
const opcua = require("node-opcua");
const ping = require('ping');
const {
    AttributeIds,
    OPCUAClient
} = require("node-opcua");
const chalk = require("chalk");
const mysql = require('mysql2');

/*接收命令行参数*/

var argvs = process.argv.slice(2);
var args = [];
argvs.forEach(function (val, index, array) {
    var arg = val.split('=');
    args[arg[0]] = arg[1];
});

/*
数据库连接
*/

var con = mysql.createConnection(config.mysql);

con.connect(function(err) {
    if (err) throw err;
    console.log(chalk.green("DB Connected!"));
});

/*
读取机器->配方->节点
*/

;(async () => {
    try {

        var sql = "";
        sql = "SELECT * FROM machine ...";

        con.query(sql, function(err, machines) {

            if (err) throw err;

            /*
            机器列表遍历
            */

            var m = 0;
            machines.forEach(async machine => {
                
                /*
                检查机器是否在线
               */              

                let HostInfo = await ping.promise.probe(machine.ip);
                console.log(HostInfo.alive);

                if(HostInfo.alive)
                {
                    console.log('Host ' + machine.ip + ' 可访问!');

                    var endpointUrl = machine.ip + ':' + machine.port;

                    console.log(endpointUrl);

                    client = OPCUAClient.create({
                        endpointMustExist: false,
                    });
            
                    console.log(" connecting to " , chalk.cyan(endpointUrl));
                    await client.connect(endpointUrl);
                    console.log(" connected to " , chalk.green(endpointUrl));

                    session = await client.createSession();
                    console.log(" session created".yellow);

                    var sql = "";
                    sql = "SELECT * FROM machine_recipes...";

                    /*
                    机器配方遍历
                    */                  

                    con.query(sql, function(err, machine_recipes) {

                        if (err) throw err;

                        machine_recipes.forEach(async machine_recipe => {

                            /*
                            读取机器配方待检查的节点
                            */

                            var sql = "";
                            sql = "SELECT * FROM machine_recipes_nodes...";

                            con.query(sql, function(err, nodes) {

                                var i = 0;
                                nodes.forEach(async node => {

                                    var params = JSON.parse(node.params);

                                    var dataValues  = await session.read([{
                                        nodeId: params.unitid
                                    }]);

                                    var node = params.unitid;
                                    var value = dataValues [0].value.value;
                                    
                                    
                                    /*
                                        将值保存到数据库...
                                    */

                                });
                            });
                        });
                    });
                }

                m++;
            });
        });

    } catch (err) {
        console.log('error !', err.message);
    }
})();
问题分析与修复方案

核心问题

  1. 异步流程混乱:同时混用MySQL的回调式con.query和async/await,forEach中的异步函数不会阻塞循环,导致多台设备的OPCUA连接、会话创建操作并行执行,资源竞争引发断言错误。此外,回调函数脱离了外层async上下文,错误无法被try/catch捕获。
  2. 全局变量污染:client和session未用let/const声明,属于全局变量,循环中会被后续迭代覆盖,导致之前的会话被意外篡改。
  3. 回调地狱:多层嵌套的con.query导致代码可读性极差,异步流程完全失控。

修复步骤

  1. Promise化MySQL查询:用util.promisify将con.query转为Promise形式,配合async/await统一异步写法。
  2. 改用for...of遍历:替代forEach,确保异步操作按顺序执行(若需并行可使用Promise.all,但串行更适合设备连接场景,避免资源过载)。
  3. 声明局部变量:将client、session声明为局部变量,避免全局污染。
  4. 增加局部错误捕获:为每个设备的连接、节点读取操作添加try/catch,避免单个设备出错导致整个脚本崩溃。

修复后的代码

/*
    此处省略导入配置及数据库连接设置的代码
*/

const endpointUrl = config.endpointUrl;
const opcua = require("node-opcua");
const ping = require('ping');
const { AttributeIds, OPCUAClient } = require("node-opcua");
const chalk = require("chalk");
const mysql = require('mysql2');
const util = require('util'); // 引入util用于promisify

/*接收命令行参数*/

const argvs = process.argv.slice(2);
const args = {}; // 改用对象存储参数更合理
argvs.forEach(val => {
    const [key, value] = val.split('=');
    args[key] = value;
});

/*
数据库连接 - Promise化query方法
*/

const con = mysql.createConnection(config.mysql);
const query = util.promisify(con.query).bind(con); // 封装为Promise形式的query

con.connect(function(err) {
    if (err) throw err;
    console.log(chalk.green("DB Connected!"));
});

/*
读取机器->配方->节点
*/

;(async () => {
    try {
        // 查询机器列表
        const machines = await query("SELECT * FROM machine ...");

        // 用for...of遍历,确保异步操作串行执行
        for (const machine of machines) {
            try {
                // 检查机器是否在线
                const hostInfo = await ping.promise.probe(machine.ip);
                console.log(hostInfo.alive);

                if (!hostInfo.alive) {
                    console.log(`Host ${machine.ip} 不可访问,跳过`);
                    continue;
                }
                console.log(`Host ${machine.ip} 可访问!`);

                const endpointUrl = `${machine.ip}:${machine.port}`;
                console.log(endpointUrl);

                // 声明局部client变量,避免全局污染
                const client = OPCUAClient.create({
                    endpointMustExist: false,
                });
        
                console.log(" connecting to " , chalk.cyan(endpointUrl));
                await client.connect(endpointUrl);
                console.log(" connected to " , chalk.green(endpointUrl));

                // 声明局部session变量
                const session = await client.createSession();
                console.log(chalk.yellow(" session created"));

                // 查询机器配方
                const machineRecipes = await query("SELECT * FROM machine_recipes...");

                for (const machineRecipe of machineRecipes) {
                    // 查询配方对应的节点
                    const nodes = await query("SELECT * FROM machine_recipes_nodes...");

                    for (const node of nodes) {
                        try {
                            const params = JSON.parse(node.params);
                            const dataValues = await session.read([{
                                nodeId: params.unitid
                            }]);

                            const nodeId = params.unitid;
                            const value = dataValues[0].value.value;
                            
                            /*
                                将值保存到数据库...
                                这里也建议用await调用Promise化的query
                            */
                            // 示例:await query("INSERT INTO ... VALUES (?, ?)", [nodeId, value]);

                        } catch (nodeErr) {
                            console.log(`节点 ${node.id} 读取失败:`, nodeErr.message);
                        }
                    }
                }

                // 用完会话和客户端后关闭
                await session.close();
                await client.disconnect();

            } catch (machineErr) {
                console.log(`机器 ${machine.id} 处理失败:`, machineErr.message);
            }
        }

    } catch (err) {
        console.log('全局错误:', err.message);
    } finally {
        // 脚本结束后关闭数据库连接
        con.end();
    }
})();

内容的提问来源于stack exchange,提问作者user1636103

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 18:50:26