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); }) } }
错误原因分析
- forEach不支持异步等待:
profile.tweets.forEach(async tweet => {})中的async函数不会阻塞循环,循环会直接执行完毕,导致提前调用bdc.end()关闭连接,而此时回调内的查询还未执行。 - Promise resolve时机错误:嵌套的Promise在调用
bdc.query后立刻执行resolve(),没有等待查询回调完成,导致流程提前推进到关闭连接的步骤。 - SQL语法错误:
DELETE * FROM table_x中的*是多余的;INSERT INTO ... WHERE id = ?语法非法,INSERT语句不能添加WHERE条件。 - 基类连接管理冲突:BBDD基类的
query方法每次调用都会重新连接并在查询后立刻关闭连接,与子类自行管理连接的逻辑冲突,导致连接被提前关闭。
修复方案
- 用
for...of循环替代forEach,确保每个异步操作完成后再执行下一个,避免循环提前结束。 - 将回调式的
bdc.query封装为Promise,统一用async/await管理异步流程,避免回调嵌套和resolve时机错误。 - 修改BBDD基类的连接逻辑,让
connect只创建一次连接,query方法不自动关闭连接,由外部手动控制连接生命周期。 - 修正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
相关产品推荐
相关产品推荐

