Node.js使用ws库实现WebSocket收消息后回复与定向推送
解答
完全可以在wss.on('connection')代码块外部获取对应请求的socket实例,你已经维护了clients映射表存储连接,只要修正现有代码的几处逻辑错误就能实现需求,不用额外做特殊处理。
你现有代码的问题
- 用时间戳
Date.now()作为key存储socket连接,没有把连接和实际用户身份(user_id、是否管理员)绑定,后续根本找不到哪个是发起请求的用户、哪个是管理员账号 - 调用
BuyCoin时没把当前触发消息事件的ws实例传进去,函数内部根本拿不到当前请求对应的连接 - MySQL查询的回调写法完全错误:你在回调形参位置写了
ws,会直接覆盖外层作用域的ws变量,这个形参实际接的是数据库返回的错误对象,调用send会直接抛错 - 连接关闭、出错时没有从
clients里删掉对应记录,会留一堆无效连接造成内存泄漏,后续发消息很容易命中已断开的连接报错
实现步骤
- 调整连接存储逻辑:等客户端连接后完成身份校验(比如第一条消息传鉴权token,校验通过拿到user_id、管理员标识),用user_id作为key存socket,同时把身份信息挂在ws实例上
- 调用外部业务函数(比如
BuyCoin)时,直接把当前message事件回调里拿到的ws实例、请求参数一起传进去,不需要绕弯去全局找 - 数据库操作完成后,用传入的ws直接响应当前请求用户,再遍历
clients找到所有管理员的socket发通知 - 连接断开、出错时及时从
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
相关产品推荐
相关产品推荐

