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

Node.js操作MySQL报错Cannot enqueue Query after invoking quit求助

Node.js MySQL循环操作时连接提前关闭的问题解决

我用Node.js编写循环执行MySQL查询的逻辑,需要根据记录是否存在执行删除或插入操作,但目前遇到报错:Error: Cannot enqueue Query after invoking quit.,原因是回调内的查询还未执行,数据库连接就已被关闭。


现有代码

代码1:BBDDP类的insertOrUpdate方法

class BBDDP extends BBDD {
    insertOrUpdate = async (profile) => {
        let bdc = await this.connect();
        await new Promise(async (resolve, reject) => {
            await profile.tweets.forEach(async tweet => {
                await new Promise(async (resolve, reject) => {
                    bdc.query(
                        ("SELECT x FROM table_x WHERE id = ?"), [profile.id],
                        function (err, row) {
                            if (err) {
                                throw err;
                            } else {
                                if (row && row.length) {
                                    //bdc.connect()
                                    bdc.query("DELETE * FROM table_x where id = ?";[profile.id])
                                    console.log('Case row was found!');
                                } else {
                                    bdc.query("INSERT INTO table_x set name = ?, addres = ? where id = ?;",[profile.name, profile.addres, profile.id])
                                    console.log('No case row was found :( !');
                                }
                            }
                        }
                    );
                    resolve();
                })
                resolve();
            });
        });
        bdc.end()
    }
}

代码2:BBDD基类

class BBDD {
    /**
     * 用于执行数据库查询的工具类
     * @param {*} ip 服务器IP地址
     * @param {*} user 认证用户名
     * @param {*} password 数据库密码
     * @param {*} dataBase 数据库名称
     * @param {*} port 端口号
     */

    constructor(ip, user, password, dataBase, port = 3306) {
        this.ip = ip;
        this.user = user;
        this.password = password;
        this.dataBase = dataBase;
        this.port = port;
        this.data = null;
        console.log(`数据库连接信息:IP ${this.ip},端口 ${this.port},数据库 ${this.dataBase}.`);
    }

    connect = async () => {
        let con = await new Promise(async (resolve) => {
            this.con = mysql.createConnection({
                host: this.ip,
                user: this.user,
                password: this.password,
                database: this.dataBase,
                port: this.port,
                charset: "utf8"
            });

            this.con.connect(function (err) {
                if (err) throw err;
                console.log("Connected!");
            });
            resolve(this.con);
            return this.con;
        });
        return con;
    }

    query = async (query, values = null) => {
        /**
         * @param {*} query 要执行的SQL语句
         * @returns MySQL查询结果
         */
        await this.connect();
        return await new Promise(async resolve => {
            if (values != null) {
                let x = this.con.query(query, values, function (error, results, fields) {
                    if (error) throw error;
                    resolve(results);
                })
                await this.closeConnection();
            } else {
                let x = this.con.query(query, function (error, results, fields) {
                    if (error) throw error;
                    resolve(results);
                })
                await this.closeConnection();
            }
        });
    }

    closeConnection = async () => {
        return await new Promise(async (resolve) => {
            this.con.end();
            console.log("连接已关闭.")
            resolve();
        })
    }
}

代码3:Main类执行逻辑

class Main {
    constructor() {}

    static async run() {
        this.api = MyApi('xxx', "xxx","xxx");
        this.bdCargas = new BBDDQ("xxx", "xxx", "xxx", "xxx");
        this.bd = new BBDDP("xxx", "xxx", "xxx", "xxx");
        let profiles = await this.bdCargas.getQueue();

        profiles.forEach(async profile => {
            await this.api.setFollowers(profile);
            await this.api.setUserContents(profile);
            await this.bd.insertOrUpdate(profile);
        })
    }
}

错误原因分析

  1. forEach不支持异步等待:profile.tweets.forEach(async tweet => {})中的async函数不会阻塞循环,循环会直接执行完毕,导致提前调用bdc.end()关闭连接,而此时回调内的查询还未执行。
  2. Promise resolve时机错误:嵌套的Promise在调用bdc.query后立刻执行resolve(),没有等待查询回调完成,导致流程提前推进到关闭连接的步骤。
  3. SQL语法错误:DELETE * FROM table_x中的*是多余的;INSERT INTO ... WHERE id = ?语法非法,INSERT语句不能添加WHERE条件。
  4. 基类连接管理冲突:BBDD基类的query方法每次调用都会重新连接并在查询后立刻关闭连接,与子类自行管理连接的逻辑冲突,导致连接被提前关闭。

修复方案

  1. 用for...of循环替代forEach,确保每个异步操作完成后再执行下一个,避免循环提前结束。
  2. 将回调式的bdc.query封装为Promise,统一用async/await管理异步流程,避免回调嵌套和resolve时机错误。
  3. 修改BBDD基类的连接逻辑,让connect只创建一次连接,query方法不自动关闭连接,由外部手动控制连接生命周期。
  4. 修正SQL语句的语法错误。

修改后的代码

优化后的BBDD基类

class BBDD {
    constructor(ip, user, password, dataBase, port = 3306) {
        this.ip = ip;
        this.user = user;
        this.password = password;
        this.dataBase = dataBase;
        this.port = port;
        this.data = null;
        this.con = null; // 保存连接实例,避免重复连接
        console.log(`数据库连接信息:IP ${this.ip},端口 ${this.port},数据库 ${this.dataBase}.`);
    }

    connect = async () => {
        if (this.con) return this.con; // 已连接则直接返回实例
        return new Promise((resolve, reject) => {
            this.con = mysql.createConnection({
                host: this.ip,
                user: this.user,
                password: this.password,
                database: this.dataBase,
                port: this.port,
                charset: "utf8"
            });

            this.con.connect(err => {
                if (err) {
                    console.error("连接失败:", err);
                    return reject(err);
                }
                console.log("Connected!");
                resolve(this.con);
            });
        });
    }

    // 封装query为Promise,不自动关闭连接
    query = async (query, values = null) => {
        await this.connect();
        return new Promise((resolve, reject) => {
            const handler = (error, results, fields) => {
                if (error) {
                    console.error("查询失败:", error);
                    return reject(error);
                }
                resolve(results);
            };
            values ? this.con.query(query, values, handler) : this.con.query(query, handler);
        });
    }

    closeConnection = async () => {
        if (!this.con) return;
        return new Promise((resolve, reject) => {
            this.con.end(err => {
                if (err) {
                    console.error("关闭连接失败:", err);
                    return reject(err);
                }
                this.con = null; // 清空连接实例
                console.log("连接已关闭.");
                resolve();
            });
        });
    }
}

优化后的BBDDP insertOrUpdate方法

class BBDDP extends BBDD {
    insertOrUpdate = async (profile) => {
        try {
            await this.connect();

            // 用for...of确保每个tweet的操作顺序执行
            for (const tweet of profile.tweets) {
                // 查询记录是否存在
                const rows = await this.query("SELECT x FROM table_x WHERE id = ?", [profile.id]);
                
                if (rows.length > 0) {
                    // 存在则删除,修正SQL语法
                    await this.query("DELETE FROM table_x WHERE id = ?", [profile.id]);
                    console.log('记录已找到,执行删除操作!');
                } else {
                    // 不存在则插入,修正SQL语法(去掉无效的WHERE)
                    await this.query(
                        "INSERT INTO table_x (name, addres, id) VALUES (?, ?, ?)",
                        [profile.name, profile.addres, profile.id]
                    );
                    console.log('未找到记录,执行插入操作!');
                }
            }
        } catch (err) {
            console.error("操作执行失败:", err);
            throw err; // 抛出错误给上层处理
        } finally {
            // 无论成功失败,最终关闭连接
            await this.closeConnection();
        }
    }
}

优化后的Main类run方法

class Main {
    constructor() {}

    static async run() {
        this.api = MyApi('xxx', "xxx","xxx");
        this.bdCargas = new BBDDQ("xxx", "xxx", "xxx", "xxx");
        this.bd = new BBDDP("xxx", "xxx", "xxx", "xxx");
        let profiles = await this.bdCargas.getQueue();

        // 用for...of替代forEach,确保每个profile的操作顺序执行
        for (const profile of profiles) {
            await this.api.setFollowers(profile);
            await this.api.setUserContents(profile);
            await this.bd.insertOrUpdate(profile);
        }
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 12:50:25