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

Node.js Express中如何实现MySQL事务互斥避免并发执行冲突?

问题背景

需要实现MySQL数据库事务互斥执行:第二个接口调用必须等待第一个调用的事务完全执行结束后才能启动。此前尝试过npm的async-mutex框架、全链路Promise封装、MySQLSELECT ... FOR UPDATE行锁语法,均未达到预期效果。

异常表现

当2个HTTP请求同时到达服务端调用对应模型函数时,两个请求会同时进入canEntryBeInserted校验逻辑,最终产生错误业务结果。

  • 期望控制台日志输出顺序:1 2 1 2
  • 实际控制台日志输出顺序:1 1 2 2

原业务代码如下:

exports.addEntry = function (entry_id, type) {
    return new Promise(function (resolve, reject) {

        console.log("1")
        beginTransactionP().then(function () {
            canEntryBeInserted(entry_id, type).then(function (canBeInserted) {
                if (canBeInserted) {
                    addEntrySQL(entry_id, type).then(function(){
                        finishTransactionP().then(function() {
                            console.log("2")
                            resolve();
                        })
                    })
                } else {
                    finishTransactionP().then(function() {
                        console.log("2")
                        resolve();
                    })
                }
            })
        })
    })

}

addEntrySQL = function(entry_id, type){
    return new Promise(function (resolve, reject) {
        sql.query("INSERT INTO entry (entry_id, typ, time) VALUES (?, ?, NOW());", [entry_id, type], function (err, rows, fields) {
            if (err) {
                console.log("error: ", err);
                reject(err);
            }
            else {
                console.log("Inserted successfully.")
                resolve();
            }
        });
    })
}

canEntryBeInserted = function (entry_id, type) {
    return new Promise(function (resolve, reject) {
        sql.query("SELECT * FROM entry WHERE entry_id = ? AND typ = 1 UNION SELECT * FROM entry WHERE entry_id = ? ORDER BY entry_id DESC LIMIT 1;", [entry_id, entry_id], function (err0, rows0) {
            console.log(JSON.stringify(rows0), type)
            if (err0) {
                console.log("error: ", err0);
                reject(err0)
            }
            else if ((rows0.length == 1 && (rows0[0].typ == type || rows0[0].typ == 1)) || rows0.length == 2) {
                console.log("Error: Not a valid operation.");
                resolve(false)
            }
            else {
                resolve(true)
            }
        })
    })
}
根因分析
  • Promise链式调用存在语法错误:所有内层异步Promise没有return,外层逻辑无法感知内层异步操作的执行进度,事务执行顺序完全失控,甚至可能出现事务提前提交、连接提前释放的问题。
  • 数据库锁逻辑失效:canEntryBeInserted中的查询语句未加排他锁,且使用UNION语法会导致FOR UPDATE无法正确锁定目标行;配合MySQL默认的REPEATABLE READ事务隔离级别,两个并发事务会读取到相同的历史快照,完全感知不到其他事务的插入操作。
  • 应用层锁适用范围有限:async-mutex只能实现单Node.js进程内的互斥,多进程、多服务实例部署场景下根本无法跨实例锁请求,必须依赖数据库层的锁机制才能实现全局互斥。
  • 大概率存在事务连接复用问题:如果beginTransactionP、canEntryBeInserted、addEntrySQL、finishTransactionP各自从连接池获取独立的数据库连接,所有操作根本不在同一个事务上下文里,锁和事务机制完全不生效。
修复方案
  1. 重构异步写法,统一使用async/await保证执行顺序,所有事务操作必须复用同一个数据库连接,禁止跨连接执行事务内操作。
  2. 去掉UNION语法,校验查询直接对目标entry_id加排他锁,同时确认entry_id字段已建立索引,避免锁全表。
  3. 补充异常回滚逻辑,避免事务挂死导致锁长时间不释放。

修复后参考代码:

// 注意:beginTransactionP需要返回绑定了当前事务的连接实例,所有事务内操作都用该连接执行
exports.addEntry = async function (entry_id, type) {
    const conn = await beginTransactionP();
    try {
        console.log("1");
        // 直接加排他锁查目标entry_id的记录,不用UNION
        const [rows0] = await conn.query(
            "SELECT * FROM entry WHERE entry_id = ? FOR UPDATE;",
            [entry_id]
        );
        console.log(JSON.stringify(rows0), type);
        const canBeInserted = !((rows0.length == 1 && (rows0[0].typ == type || rows0[0].typ == 1)) || rows0.length == 2);
        
        if (canBeInserted) {
            await conn.query(
                "INSERT INTO entry (entry_id, typ, time) VALUES (?, ?, NOW());",
                [entry_id, type]
            );
            console.log("Inserted successfully.");
        } else {
            console.log("Error: Not a valid operation.");
        }
        // 提交事务
        await conn.commit();
        console.log("2");
    } catch (err) {
        console.log("error: ", err);
        // 异常必回滚
        await conn.rollback();
        throw err;
    } finally {
        // 释放连接回连接池
        conn.release();
    }
}

修复后并发请求会在SELECT ... FOR UPDATE处阻塞,前一个事务提交释放锁之后,后一个请求才能读到最新的记录继续执行,日志顺序会符合1 2 1 2的预期。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 04:06:06