WebSocket实时通信实现 websocket-implementation

该技能用于构建可扩展的实时通信系统,通过WebSocket技术实现服务器端连接管理、房间机制消息传递和水平扩展。适用于开发聊天应用、实时通知、协作工具或实时仪表板等场景。关键词包括WebSocket、Socket.IO、实时通信、连接管理、房间消息、水平扩展、Redis适配、JavaScript、Node.js、前后端集成。

后端开发 0 次安装 0 次浏览 更新于 3/7/2026

名称: websocket-implementation 描述: 实现实时WebSocket通信,包括连接管理、基于房间的消息传递和水平扩展。在构建聊天系统、实时通知、协作工具或实时仪表板时使用。

WebSocket实现

构建具有适当连接管理的可扩展实时通信系统。

服务器实现 (Socket.IO)

const { Server } = require('socket.io');
const { createAdapter } = require('@socket.io/redis-adapter');
const { createClient } = require('redis');

const io = new Server(server, {
  cors: { origin: process.env.CLIENT_URL, credentials: true }
});

// Redis适配器用于水平扩展
const pubClient = createClient({ url: process.env.REDIS_URL });
const subClient = pubClient.duplicate();
Promise.all([pubClient.connect(), subClient.connect()]).then(() => {
  io.adapter(createAdapter(pubClient, subClient));
});

// 连接管理
const users = new Map();

io.use((socket, next) => {
  const token = socket.handshake.auth.token;
  try {
    socket.user = verifyToken(token);
    next();
  } catch (err) {
    next(new Error('认证失败'));
  }
});

io.on('connection', (socket) => {
  users.set(socket.user.id, socket.id);
  console.log(`用户 ${socket.user.id} 已连接`);

  socket.on('join-room', (roomId) => {
    socket.join(roomId);
    socket.to(roomId).emit('user-joined', socket.user);
  });

  socket.on('message', ({ roomId, content }) => {
    io.to(roomId).emit('message', {
      userId: socket.user.id,
      content,
      timestamp: Date.now()
    });
  });

  socket.on('disconnect', () => {
    users.delete(socket.user.id);
  });
});

// 消息分发的实用方法
function broadcastUserUpdate(userId, data) {
  io.to(`user:${userId}`).emit('user-update', data);
}

function notifyRoom(roomId, event, data) {
  io.to(`room:${roomId}`).emit(event, data);
}

function sendDirectMessage(userId, message) {
  const socketId = users.get(userId);
  if (socketId) {
    io.to(socketId).emit('direct-message', message);
  }
}

客户端实现

import { io } from 'socket.io-client';

class WebSocketClient {
  constructor(url, token) {
    this.socket = io(url, {
      auth: { token },
      reconnection: true,
      reconnectionDelay: 1000,
      reconnectionAttempts: 5
    });

    this.messageQueue = [];
    this.setupListeners();
  }

  setupListeners() {
    this.socket.on('connect', () => {
      console.log('已连接');
      this.flushQueue();
    });

    this.socket.on('disconnect', (reason) => {
      console.log('断开连接:', reason);
    });

    this.socket.on('message', (msg) => {
      this.onMessage?.(msg);
    });
  }

  joinRoom(roomId) {
    this.socket.emit('join-room', roomId);
  }

  send(roomId, content) {
    if (this.socket.connected) {
      this.socket.emit('message', { roomId, content });
    } else {
      this.messageQueue.push({ roomId, content });
    }
  }

  flushQueue() {
    while (this.messageQueue.length > 0) {
      const msg = this.messageQueue.shift();
      this.socket.emit('message', msg);
    }
  }
}

React钩子

function useWebSocket(url) {
  const [socket, setSocket] = useState(null);
  const [connected, setConnected] = useState(false);
  const [messages, setMessages] = useState([]);

  useEffect(() => {
    // getToken() 是用户提供的辅助函数,返回当前认证令牌
    // 示例实现:
    // - 从localStorage: () => localStorage.getItem('authToken')
    // - 从上下文: () => authContext.token
    // - 从cookie: () => document.cookie.split('token=')[1]
    const ws = io(url, { auth: { token: getToken() } });

    ws.on('connect', () => setConnected(true));
    ws.on('disconnect', () => setConnected(false));
    ws.on('message', (msg) => {
      setMessages(prev => [...prev, msg]);
    });

    setSocket(ws);
    return () => ws.disconnect();
  }, [url]);

  const send = useCallback((roomId, content) => {
    socket?.emit('message', { roomId, content });
  }, [socket]);

  return { connected, messages, send };
}

消息协议

interface Message {
  type: 'message' | 'typing' | 'presence';
  roomId: string;
  userId: string;
  content?: string;
  timestamp: number;
}

// 确认传递
socket.emit('message', data, (ack) => {
  if (ack.success) console.log('已传递');
});

附加实现

参见 references/python-websocket.md 获取:

  • Python aiohttp WebSocket 服务器
  • FastAPI WebSocket 端点
  • 异步连接处理

扩展考虑

连接数 架构
<10K 单服务器
10K-100K Redis 发布/订阅
>100K 分片Redis + 负载均衡器

监控端点

// Express 端点用于操作可见性
app.get('/api/ws/stats', (req, res) => {
  res.json({
    activeConnections: io.sockets.sockets.size,
    rooms: [...io.sockets.adapter.rooms.keys()],
    users: users.size
  });
});

app.get('/api/ws/health', (req, res) => {
  res.json({
    status: '健康',
    uptime: process.uptime(),
    memoryUsage: process.memoryUsage()
  });
});

最佳实践

  • 允许操作前进行认证
  • 实现指数退避重连
  • 使用房间和频道进行定向广播
  • 添加心跳/健康检查
  • 持久化重要消息到数据库
  • 监控活跃连接计数
  • 显示用户在线/可用状态
  • 对传入消息实施速率限制
  • 使用确认机制确认消息传递
  • 利用Redis进行分布式部署
  • 实现全面的错误处理

切勿做

  • 发送未加密的敏感数据
  • 在内存中存储无限消息
  • 跳过房间加入的授权
  • 忽略连接错误处理
  • 允许无限制的房间订阅
  • 忽视断开连接用户数据的清理
  • 发送频繁的大负载消息
  • 在消息体中包含认证凭证
  • 部署前未进行安全验证
  • 允许不受控制的连接累积
  • 构建不考虑可扩展性