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

如何在async/await外部获取mysql2查询结果并通过WebSocket发送

问题分析与修复方案

核心问题

你的代码里的async逻辑存在错误未捕获、WebSocket对象引用不当的问题,导致拿到数据库结果后无法发送到客户端。


修复步骤

1. 直接使用当前WebSocket连接对象,避免间接引用

在connection回调里,你已经持有ws对象的直接引用,没必要通过clients[ws.id]发送消息,直接用ws.send()更可靠,避免字典引用失效的问题。

2. 必须捕获async函数内的所有错误

自执行async函数如果抛出错误(比如selectRows返回undefined后访问result[0]报错),会直接静默失败,导致后续的send代码根本不执行。要给所有async逻辑套try/catch。

3. 修正数据库查询函数的错误处理

如果数据库查询出错,你的selectRows只打印错误但返回undefined,后续代码访问result[0]会抛出未捕获错误。应该让错误冒泡,或者明确返回错误标识,方便上层处理。


修改后的完整代码示例

(async function main(){
  try{
    const mysql2 = require('mysql2/promise');
    const http = require('http');
    const WebSocket = require('ws');
    const { v4: uuidv4 } = require('uuid');
    const url = require('url');
    const token = '你的token';
    const dbConfig = '你的配置';
    
    const clients = {}; // 加上const,避免全局变量污染

    // 创建web服务器和WebSocket服务
    const server = http.createServer().listen(6); // listen无需await,同步启动即可
    const wss = new WebSocket.Server({ server});

    // 数据库连接池
    const poolConfig = {
      host: "127.0.0.1",
      user: "arealuser",
      password: "qwertyuiop",
      database: "example"
    }
    const pool = mysql2.createPool(poolConfig);

    console.log("WebSocket服务器已启动...")

    wss.on('connection', async function connection(ws, req) { // 把回调改成async,避免嵌套自执行函数
      const tokenMatch = url.parse(req.url, true).query.token;
      const service = url.parse(req.url, true).query.service;

      if(tokenMatch !== token || service !== "encryption") {
        ws.close(401, '权限验证失败或服务类型不匹配'); // 验证失败直接关闭连接
        return;
      }

      ws.id = uuidv4();
      ws.ip = ws._socket.remoteAddress;
      console.log(`UUID ${ws.id} 已连接(加密服务): ${ws.ip}`);
      clients[ws.id] = ws;
      
      // 先发送连接成功消息
      ws.send("connected");
      ws.send('test 1');

      try {
        // 直接在async回调里执行数据库查询
        const result = await selectRows(4, pool);
        console.log(result[0]);
        // 直接用当前ws对象发送
        ws.send('test result');
        // 发送JSON格式结果,方便客户端解析
        ws.send(JSON.stringify({ type: 'queryResult', data: result }));
      } catch (err) {
        console.error('查询或发送失败:', err);
        ws.send(JSON.stringify({ type: 'error', message: err.message }));
      }
    });

    // 修正selectRows的错误处理:出错时抛出错误,让上层捕获
    async function selectRows(uid,pool){
      const sql = "SELECT idpasswords, password, salt, (SELECT ops_memlimit FROM users WHERE idusers =?) AS ops_memlimit FROM passwords WHERE uid = ?";
      const [rows, fields] = await pool.query(sql,[uid,uid]);
      return rows;
    }

    // 其他Websocket服务逻辑...

  } catch(error){
    console.error('服务器启动失败:', error);
  }
})();

针对reEncryptDb函数的修复

同样遵循直接用ws对象、捕获所有错误的原则:

async function reEncryptDb(ws, uid, pool){ // 传入连接池而非mysql2模块,更高效
  try{
    ws.send("starting the function");
    console.log("starting the function");
    const result = await selectPasswords(uid, pool);
    console.log("after await");
    ws.send("after await");
    // WebSocket只能发字符串/Buffer,对象要转JSON
    ws.send(JSON.stringify({ id: result[0].idpasswords }));
    ws.send("after result");
  }
  catch(error){
    console.error(`执行加密时出错: ${error}`);
    ws.send(`执行加密时出错: ${error.message}`);
  }
}

// 对应的selectPasswords修正错误处理
async function selectPasswords(uid, pool){
  const sql = "你的查询语句";
  const [rows] = await pool.query(sql, [uid]);
  if(rows.length === 0) throw new Error('未找到数据');
  return rows;
}

关键注意事项

  • WebSocket数据格式限制:只能发送字符串或Buffer,发送对象必须用JSON.stringify()转换,客户端再用JSON.parse()解析。
  • 检查连接状态:发送前可通过ws.readyState === WebSocket.OPEN判断连接是否正常,避免连接已关闭时报错。
  • 避免全局变量:clients要加const声明,防止全局变量污染引发意外问题。
  • async错误必须捕获:所有async逻辑都要套try/catch,否则错误会静默丢失,难以排查。

内容的提问来源于stack exchange,提问作者Él Sea

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 02:26:03