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

Node.js使用ws库实现WebSocket收消息后回复与定向推送

解答

完全可以在wss.on('connection')代码块外部获取对应请求的socket实例,你已经维护了clients映射表存储连接,只要修正现有代码的几处逻辑错误就能实现需求,不用额外做特殊处理。

你现有代码的问题

  • 用时间戳Date.now()作为key存储socket连接,没有把连接和实际用户身份(user_id、是否管理员)绑定,后续根本找不到哪个是发起请求的用户、哪个是管理员账号
  • 调用BuyCoin时没把当前触发消息事件的ws实例传进去,函数内部根本拿不到当前请求对应的连接
  • MySQL查询的回调写法完全错误:你在回调形参位置写了ws,会直接覆盖外层作用域的ws变量,这个形参实际接的是数据库返回的错误对象,调用send会直接抛错
  • 连接关闭、出错时没有从clients里删掉对应记录,会留一堆无效连接造成内存泄漏,后续发消息很容易命中已断开的连接报错

实现步骤

  1. 调整连接存储逻辑:等客户端连接后完成身份校验(比如第一条消息传鉴权token,校验通过拿到user_id、管理员标识),用user_id作为key存socket,同时把身份信息挂在ws实例上
  2. 调用外部业务函数(比如BuyCoin)时,直接把当前message事件回调里拿到的ws实例、请求参数一起传进去,不需要绕弯去全局找
  3. 数据库操作完成后,用传入的ws直接响应当前请求用户,再遍历clients找到所有管理员的socket发通知
  4. 连接断开、出错时及时从clients里清理无效记录,发消息前先判断连接处于可用状态

修正后的参考代码

const WebSocketServer = require('ws');
const mysql = require('mysql');
// 初始化数据库连接,替换成你自己的配置
const connection = mysql.createConnection({
  host: 'localhost',
  user: 'db_user',
  password: 'db_pwd',
  database: 'your_db_name'
});
connection.connect();

// 启动ws服务
const wss = new WebSocketServer.Server({ port: 8080 });
// 存储在线连接,key为用户ID,value为对应的ws实例
const clients = new Map();

wss.on("connection", ws => {
    console.log("新客户端已连接");
    // 临时存当前连接的用户信息,鉴权后赋值
    ws.userInfo = null;

    ws.on('message', rawMsg => {
      const msg = JSON.parse(rawMsg);
      // 处理鉴权逻辑
      if (msg.type === 'auth') {
        // 这里替换成你实际的鉴权流程,比如校验token拿到用户ID、是否为管理员
        const userId = msg.user_id;
        const isAdmin = !!msg.is_admin;
        ws.userInfo = { userId, isAdmin };
        clients.set(userId, ws);
        ws.send(JSON.stringify({type: 'auth', code: 0, msg: '连接成功'}));
        return;
      }

      // 处理买币请求,直接把当前ws传入业务函数
      if (msg.type === 'buy_coin') {
        BuyCoin(ws, msg);
      }
    });

    ws.on("close", () => {
        console.log("客户端已断开连接");
        // 清理无效连接
        if (ws.userInfo?.userId) clients.delete(ws.userInfo.userId);
    });

    ws.onerror = function () {
        console.log("连接发生错误");
        // 出错也清理连接
        if (ws.userInfo?.userId) clients.delete(ws.userInfo.userId);
    }
});

// 买币业务逻辑
function BuyCoin(currentWs, reqParams){
  console.log(`处理用户${reqParams.user_id}的买币请求`);
  const selectSql = 'SELECT * FROM users WHERE id = ? LIMIT 1';
  connection.query(selectSql, [reqParams.user_id], (err, rows) => {
    if (err) {
      currentWs.send(JSON.stringify({type: 'buy_coin', code: -1, msg: '服务异常'}));
      return;
    }
    // 先响应给发起请求的用户
    currentWs.send(JSON.stringify({
      type: 'buy_coin',
      code: 0,
      data: rows[0]
    }));

    // 遍历连接列表给所有在线管理员发通知
    for (const [userId, ws] of clients) {
      // 只给已标记为管理员、连接正常的实例发消息
      if (ws.userInfo?.isAdmin && ws.readyState === WebSocketServer.OPEN) {
        ws.send(JSON.stringify({
          type: 'admin_notification',
          content: `用户${reqParams.user_id}已完成买币操作`
        }));
      }
    }
  });
}

注意点

  • 不要用时间戳这类和业务身份无关的值当连接的key,不然你永远没法精准定位到指定用户的连接
  • 给客户端发消息前一定要判断ws.readyState === WebSocketServer.OPEN,避免给已关闭的连接发消息抛错
  • SQL语句的参数一定要放在数组里传入,避免SQL注入风险,也能兼容多参数的查询场景

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 00:57:26