Node.js使用mysql驱动串行执行多个MySQL事务报错问题咨询
问题根因
- 同步for循环不会等待异步的数据库回调执行完成,会一次性触发所有事务的创建逻辑,导致事务并发执行
- 每个事务的处理逻辑中都调用了
res.send()方法,HTTP响应仅允许发送一次,第一个事务完成发送响应后,后续事务再尝试调用res.send()就会抛出Error [ERR_HTTP_HEADERS_SENT]错误
解决方案(无需引入Sequelize)
推荐使用async/await + Promise化的MySQL接口实现串行事务,逻辑清晰易维护,改动量极小:
步骤1:安装mysql2(兼容原mysql API,原生支持Promise)
npm install mysql2
步骤2:改写逻辑
首先初始化连接池的时候改用mysql2的promise版本:
const mysql = require('mysql2/promise'); // 连接池配置和原来完全一致 const pool = mysql.createPool({ host: '你的数据库地址', user: '你的数据库用户名', password: '你的数据库密码', database: '你的库名', waitForConnections: true, connectionLimit: 10, queueLimit: 0 });
然后改写路由代码:
router.post('/production/add', async (req, res) => { // 遍历所有请求参数,串行执行事务 for (const obj of req.body) { let connection; try { // 从连接池获取连接 connection = await pool.getConnection(); // 开启事务 await connection.beginTransaction(); // 执行查询1 改用参数化查询避免SQL注入 const [result1] = await connection.query( 'select qty from production where prc_id = ?', [obj.prc_id] ); const new_prod_qty = result1[0].qty - obj.auth_prod_qty; // 执行查询2 let query2, params2; if (new_prod_qty > 0) { query2 = 'update production set qty = ? where prc_id = ?'; params2 = [new_prod_qty, obj.prc_id]; } else { query2 = 'delete from production where prc_id = ?'; params2 = [obj.prc_id]; } const [results2] = await connection.query(query2, params2); // 执行查询3 await connection.query( 'update prc set auth_prod_qty = ? where prc_id = ?', [obj.auth_prod_qty, obj.prc_id] ); // 执行查询4 const [storeResult] = await connection.query('select * from store'); // 提交事务 await connection.commit(); // 释放连接回连接池 connection.release(); } catch (error) { // 出错回滚事务 if (connection) { await connection.rollback(); connection.release(); } // 仅发送一次错误响应,直接结束请求 return res.status(400).send({ message: error.sqlMessage || error.message, code: error.code, errno: error.errno }); } } // 所有事务都执行成功后,统一发送成功响应 res.status(201).send({ message: '操作完成' }); });
可选:不更换mysql包的方案
如果你不想更换原有的mysql包,可以用Node.js内置的util.promisify将回调式的API转为Promise式,后续使用逻辑和上述mysql2的写法完全一致:
const { promisify } = require('util'); // 把mysql的回调方法转为Promise pool.getConnection = promisify(pool.getConnection);
额外说明
- 原有代码直接拼接SQL字符串存在SQL注入风险,上述方案中统一改用参数化查询规避该问题
- 串行执行时默认某一个事务失败就终止后续执行,你也可以根据业务需求调整为失败后仅记录错误、继续执行剩余事务
内容的提问来源于stack exchange,提问作者Rohan
相关产品推荐
相关产品推荐

