【问题标题】:Nodejs MySql pool transaction with loop带循环的Nodejs MySql池事务
【发布时间】:2022-01-03 17:11:12
【问题描述】:

我目前正在使用 nodejs-mysql (https://www.npmjs.com/package/mysql)。我正在尝试将 mysql 连接从 mysql.createConnection(config) 更改为 mysql.createPool(config)。到目前为止,除了事务之外,大多数 API 都不需要做太多更改。

问题:有一个 for 循环将在池事务回调地狱中执行查询。但循环并未等待回调完成

  • 旧代码 - mysql.createConnection(config) [工作正常]
  /* Begin transaction */
  try {
    await db.beginTransaction();

    // insert sales
    var saleId = await new Promise((resolve, reject) => {
      db.query('INSERT INTO sales SET ?', saleData, (err, dbResult) => {
        if (err) {
          if (err.sqlMessage) console.log(err.sqlMessage);
          resolve(null);
        } else {
          resolve(Object.keys(dbResult).length ? dbResult.insertId : null);
        }
      });
    });

    if (!saleId) {
      throw "Error: sale id not found, fail to insert sale.";
    }

    // insert payment
    var paymentData = {
      method: "cash",
      sale_id: saleId
    }

    var paymentId = await new Promise((resolve, reject) => {
      db.query(`INSERT INTO payments SET ?`, paymentData, (err, dbResult) => {
        if (err) {
          if (err.sqlMessage) console.log(err.sqlMessage);
          resolve(null);
        } else {
          resolve(Object.keys(dbResult).length ? dbResult.insertId : null);
        }
      });
    });

    if (!paymentId) {
      throw "Error: payment id not found, fail to insert payment.";
    }

    // update stock
    // just an update query inside a loop
    for (var i = 0; i < cartList.length; i++) {
      var updateStockSql = `UPDATE stocks 
        SET ${branchId}_quantity = COALESCE(${branchId}_quantity, 0) - ${cartList[i].Total_stock} 
        WHERE product_id = ${cartList[i].Product_id}`;

      var updateStockResult = await new Promise((resolve, reject) => {
        db.query(updateStockSql, (err, dbResult) => {
          if (err) {
            if (err.sqlMessage) console.log(err.sqlMessage);
            resolve(null);
          } else {
            resolve(Object.keys(dbResult).length ? dbResult : null);
          }
        });
      });

      if (!updateStockResult) throw "Error: Fail to update stock.";
    }

    await db.commit();
  } catch (error) {
    await db.rollback();
    return { status: "error", msg: "- Transaction Fail -" };
  }
  /* End transaction */
  • 当前代码 - mysql.createPool(config) [遇到问题]
/* Begin transaction */
  try {
    var trxStatus = await new Promise((resolve, reject) => {
      db.getConnection((err, trxConn) => {

        trxConn.beginTransaction((err) => {

          //start
          if (err) {
            trxConn.rollback(() => trxConn.release());
            resolve({ status: "error", msg: `Transaction Fail: ${err}` });
          } else {

            // insert sales
            trxConn.query('INSERT INTO sales SET ?', saleData, (err, salesResult) => {
              if (err || !salesResult.insertId) {
                // rollback
                trxConn.rollback(() => trxConn.release());
                if (!salesResult.insertId) err = "Error: sale id not found, fail to insert sale.";
                else if (err.hasOwnProperty('sqlMessage')) err = err.sqlMessage;
                resolve({ status: "error", msg: `Transaction Fail: ${err}` });
              } else {

                // insert payments
                trxConn.query('INSERT INTO payments SET ?', paymentData, (err, paymentResult) => {
                  if (err || !paymentResult.insertId) {
                    // rollback
                    resolve({ status: "error", msg: `Transaction Fail: ${err}` });
                  } else {

                    // here is whr I'm stucked
                    // update stock
                    var updateStockStatus = true;
                    for (var i = 0; i < cartList.length; i++) {
                      var updateStockSql = `UPDATE stocks 
                      SET ${branchId}_quantity = COALESCE(${branchId}_quantity, 0) - ${cartList[i].Total_stock} 
                      WHERE product_id = ${cartList[i].Product_id}`;

                      trxConn.query(updateStockSql, (err, updateStockResult) => {
                        if (err || !Object.keys(updateStockResult).length) {
                          if (err.hasOwnProperty('sqlMessage')) console.log(err.sqlMessage);
                          updateStockStatus = false;
                        }
                      });
                      if (!updateStockStatus) break;
                    }

                    if (!updateStockStatus) {
                      // rollback
                      trxConn.rollback(() => trxConn.release());
                      resolve({ status: "error", msg: `Transaction Fail: fail to update stock` });
                    } else {
                      // commit
                      trxConn.commit((err) => {
                        if (err) trxConn.rollback(() => trxConn.release());
                        else trxConn.release();
                        resolve({ status: "ok", msg: `- Payment Successful -` });
                      });
                    }
                  }
                });
              }
            });
          }
          //end

        });
      });
    });
  } catch (error) {
    console.log(error);
    return { status: "error", msg: `Transaction Fail: please contact IT support` };
  }
  /* End transaction */

ps:抱歉问题太长了,我已经尽力简化了。

问题:有人知道如何解决这种异步问题吗?

任何帮助将不胜感激,在此先感谢。

【问题讨论】:

    标签: mysql node.js express connection-pooling


    【解决方案1】:

    我怀疑您的问题可能是您的代码的这一部分:

    var updateStockStatus = true;
    // Start Loop
    //  |_ Do Stuff
    //    |_Looped Call out
    //      |_Wait for each Response
    // End Loop
    
    if (!updateStockStatus) {
      rollback()
    } else {
      commit()
    }
    

    我认为问题在于异步执行。当你的Looped Call out代码正在执行时,js会继续,测试updateStockStatus然后执行commit()方法。

    尝试像这样将updateStockStatus 的测试移动到循环中:

    var updateStockStatus = true;
    var numOfCalls = cartList.length; // just for clarity
    // Start Loop
    for (i = 0; i <= numOfCalls; i++){
        // Do Stuff
        var sqlQuery = 'foobar'
        // Looped Query
        trxConn.query(sqlQuery, (err, callbackFn) =>{
            // On Response Do Stuff
            if (err) {
                // if there's a problem, initiate rollback
                trxConn.rollback()
                break; //Stop processing
            } else { 
                // if there's no error, do stuff
                var stuff = 'Just Do It'
                
                // check if this is the last call to make
                if (i == numOfCalls) {
                    //if it is, then initiate a commit
                    trxConn.commit()
                }
    
        })
        
    }
    

    【讨论】:

    • 嗨,肖恩,感谢您的回复。但是我不认为您提供的解决方案是我正在寻找的...也许您可以用代码提供您的答案,但只是为了确保我正确理解您的解决方案。
    • 嗨@Chris 道歉,我今天才看到这个。我已经更新了我的答案以包含代码,但主要的收获是从 .query() 方法内部而不是外部启动回滚/提交。这是因为 Javascript 会在等待 .query() 调用完成时继续处理,因此它会在您真正得到所有响应之前测试您的 if (!updateStockStatus) 语句。
    猜你喜欢
    • 1970-01-01
    • 2019-06-15
    • 2012-05-27
    • 2012-05-09
    • 1970-01-01
    • 2015-08-13
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多