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

Node.js+Redis集成WebSockets报错排查与多服务器同步咨询

解决方案:Redis订阅与发布客户端分离

错误根源

你遇到的ERR Can't execute 'publish': only (P)SUBSCRIBE / (P)UNSUBSCRIBE / PING / QUIT / RESET are allowed in this context错误,本质是同一个Redis客户端不能同时承担订阅和发布角色。当客户端执行subscribe进入订阅模式后,Redis会限制该客户端只能执行订阅相关命令,无法再执行publish操作。同时你的代码里还存在redisPubClient未定义的问题,这也是导致消息监听无输出的原因。

修正方案

需要创建两个独立的Redis客户端:

  • 一个专门用于发布消息(pubClient)
  • 一个专门用于订阅频道(subClient)

当WebSocket收到客户端消息时,用pubClient发布到Redis频道;subClient监听频道消息,收到后广播给当前服务器上的所有WebSocket连接,以此实现多服务器间的消息同步(负载均衡场景下,不同服务器的客户端通过Redis共享消息)。

修正后的完整代码

const PackageJson = require("./package.json");
const redis = require('redis');

module.exports = (async () => {
  const {Config} = await require("./common/config");

  // ==============================================
  // Base setup
  // ==============================================

  process.env.TZ = "Etc/GMT";
  const express = require("express");
  const session = require("express-session");
  const app = express();
  const server = require('http').createServer(app);
  const WSServer = require("ws").Server;
  const wss = new WSServer({server: server,});
  const port = process.env.port || 3000;
  const bodyParser = require('body-parser');
  app.use(bodyParser.json({limit: '1500mb'}));
  app.use(bodyParser.urlencoded({limit: '1500mb', extended: true, parameterLimit:1500000}));
  app.use(require("express-useragent").express());
  app.use(session({secret: "keyboard cat",resave: false,saveUninitialized: true,cookie: { secure: true },}));
  app.use(async (req,res,next) => {
    res.header("Access-Control-Allow-Origin", "*");
    req.session.language = req.query.lg;
    res.header("Access-Control-Allow-Headers","Origin, X-Requested-With, Content-Type, Accept, auth-uid, auth-token,");
    res.header("Access-Control-Allow-Methods","PATCH, POST, GET, DELETE, OPTIONS");
    next();
  });

  // ===============================================================
  // Routes
  // ===============================================================

  app.get("/", (req,res) => {res.send(
    "<html><head></head><body>API is running <br>"+
      "App: "+Config.FrontEnd.AppName+"<br>"+
      "Env: "+Config.Env+"<br>"+
      "Version: "+PackageJson.version+"<br>"+
    "</body></html>"
  );});

  // ==============================================
  // Redis - 分离发布/订阅客户端
  // ==============================================

  // 1. 发布客户端:用于向Redis频道发送消息
  const redisPubClient = redis.createClient({
    password: Config.Keys.RedisPassword,
    socket: {
        host: Config.Keys.RedisUrl,
        port: Config.Keys.RedisPort
    }
  });

  redisPubClient.on('error', err => console.log('Redis Pub Client Error', err));
  await redisPubClient.connect();

  // 2. 订阅客户端:用于监听Redis频道消息
  const redisSubClient = redis.createClient({
    password: Config.Keys.RedisPassword,
    socket: {
        host: Config.Keys.RedisUrl,
        port: Config.Keys.RedisPort
    }
  });

  redisSubClient.on('error', err => console.log('Redis Sub Client Error', err));
  await redisSubClient.connect();

  // 订阅目标频道,收到消息后广播给当前服务器的WebSocket客户端
  await redisSubClient.subscribe('my-channel', (message) => {
    wss.clients.forEach((client) => {
      if (client.readyState === client.OPEN) {
        client.send(message);
      }
    });
  });

  // 订阅成功回调
  redisSubClient.on('subscribe', (channel, count) => {
    console.log(`Subscribed to Redis channel ${channel}, current subscription count: ${count}`);
  });

  // ==============================================
  // Sockets
  // ==============================================

  const wssRooms = {};
  wss.on("connection", function connection(ws) {
    console.log('New WebSocket client connected');

    ws.on("message", async function incoming(message) {
      // 使用发布客户端发送消息到Redis频道
      await redisPubClient.publish("my-channel", message.toString());
    });

    ws.on("close", () => {
      console.log('WebSocket client disconnected');
    });
  });

  // ==============================================
  // Start the server
  // ==============================================
  server.listen(process.env.PORT || port, () => {
    console.log(Config.FrontEnd.AppName+` Listening on port:${process.env.PORT || port} Env:${Config.Env}`);
  });
})();

关键说明

  1. 客户端分离:严格区分发布和订阅客户端,避免订阅模式下的命令限制
  2. 消息广播逻辑:订阅客户端收到Redis消息后,遍历当前服务器的所有WebSocket连接,将消息推送给在线客户端
  3. 多服务器同步:负载均衡场景下,任意一台服务器的WebSocket客户端发送消息,都会通过Redis广播到所有服务器的订阅客户端,进而推送给所有服务器上的客户端,实现跨节点的消息同步
  4. 错误处理:为两个Redis客户端分别添加错误监听,便于排查连接或命令执行问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 21:12:50