名称: 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进行分布式部署
- 实现全面的错误处理
切勿做
- 发送未加密的敏感数据
- 在内存中存储无限消息
- 跳过房间加入的授权
- 忽略连接错误处理
- 允许无限制的房间订阅
- 忽视断开连接用户数据的清理
- 发送频繁的大负载消息
- 在消息体中包含认证凭证
- 部署前未进行安全验证
- 允许不受控制的连接累积
- 构建不考虑可扩展性