基于JavaScript与Node.js自研传输协议的实时聊天开发指引求助
自研类WebSocket实时通信协议的入门方向与可行方案
核心思路参考
WebSocket的本质是基于TCP的全双工通信协议,核心包含HTTP握手升级、自定义帧格式、心跳保活三个关键环节。自研时无需完全照搬WebSocket的标准,可根据需求简化或调整,但这三个环节是实现实时双向通信的基础。
入门方向
- 吃透TCP字节流特性:重点理解TCP的粘包/拆包问题,这是自定义协议的核心难点——因为TCP会把数据拆分成多个包发送,或把多个消息合并成一个包,必须通过自定义帧格式来区分独立消息。
- 参考WebSocket设计逻辑:不用看具体库实现,只需要了解它的握手流程(HTTP升级到TCP连接)、帧结构(长度前缀+内容)、心跳机制的作用,这些是解决浏览器端无法直接发起TCP连接、连接保活的关键。
- 掌握Node.js原生模块:熟练使用
net(TCP服务端/客户端)、http(处理HTTP握手升级)模块,这是服务端实现的基础。 - 熟悉浏览器端流式API:浏览器无法直接操作TCP socket,需要用
fetch流式响应、ReadableStream来接收实时数据,结合POST请求或其他方式发送数据(HTTP/2可实现单连接全双工,HTTP/1.1需用双连接或复用连接)。
可行实现方案
1. 基于HTTP握手升级建立双向连接
浏览器不能直接发起TCP连接,所以先通过HTTP请求触发握手,将连接升级为可双向读写的TCP连接:
- 服务端:用
http模块监听upgrade事件,校验自定义的升级头(比如Upgrade: my-chat-protocol),返回101 Switching Protocols响应,之后直接操作底层socket收发数据。 - 客户端:发送带升级头的HTTP请求,收到101响应后,通过流式API读取服务端推送的数据,同时用独立请求或复用连接发送数据。
2. 自定义帧格式解决粘包问题
必须定义统一的帧结构来拆分字节流,推荐最简单的长度前缀格式:
[4字节消息长度(大端序)] + [消息内容(字符串/二进制)]
- 编码:发送消息时,先计算内容的字节长度,写入4字节Uint32,再拼接消息内容。
- 解码:接收数据时,先缓存字节流,当缓存长度≥4字节时,读取长度值,再等待缓存长度≥4+长度值,拆分出完整消息,剩余缓存继续等待下一个帧。
3. 实现心跳与连接保活
TCP连接可能因网络超时、路由变化静默断开,需添加心跳机制:
- 客户端定期发送
PING帧(比如每30秒),服务端收到后回复PONG帧。 - 若超过指定时间(比如90秒)未收到对方的心跳响应,判定连接断开,触发重连逻辑。
4. 服务端连接池管理
维护一个全局的socket集合,当收到客户端消息时,遍历集合广播给其他在线客户端,实现群聊功能。
代码示例(简化版)
服务端(Node.js)
const http = require('http'); const connections = new Set(); // 连接池 const server = http.createServer((req, res) => { if (req.url === '/send') { // 处理客户端发送的消息 let body = ''; req.on('data', chunk => body += chunk); req.on('end', () => { const { message } = JSON.parse(body); // 广播消息给所有连接 const frame = encodeMessage(message); connections.forEach(socket => socket.write(frame)); }); res.writeHead(200); res.end(); return; } res.writeHead(200); res.end('Chat Server Running'); }); // 处理升级请求 server.on('upgrade', (req, socket, head) => { if (req.headers.upgrade !== 'my-chat-protocol') { socket.end('HTTP/1.1 400 Bad Request'); return; } // 发送升级响应 socket.write('HTTP/1.1 101 Switching Protocols\r\n'); socket.write('Upgrade: my-chat-protocol\r\n'); socket.write('Connection: Upgrade\r\n\r\n'); connections.add(socket); console.log('新客户端连接,当前在线:', connections.size); socket.on('close', () => { connections.delete(socket); console.log('客户端断开,当前在线:', connections.size); }); // 监听心跳(示例) socket.on('data', data => { const msg = decodeMessage(data); if (msg === 'PING') { socket.write(encodeMessage('PONG')); } }); }); // 编码消息为帧 function encodeMessage(message) { const content = Buffer.from(message, 'utf8'); const frame = Buffer.alloc(4 + content.length); frame.writeUInt32BE(content.length, 0); // 大端序写入长度 frame.write(message, 4); return frame; } // 解码帧为消息(简化版,未处理粘包) function decodeMessage(data) { const length = data.readUInt32BE(0); return data.slice(4, 4 + length).toString('utf8'); } server.listen(3000, () => console.log('服务端运行在 http://localhost:3000'));
客户端(浏览器JS)
class CustomChatClient { constructor(baseUrl) { this.baseUrl = baseUrl; this.onMessage = null; this.onClose = null; this.heartbeatTimer = null; } async connect() { // 建立流式连接接收消息 const res = await fetch(this.baseUrl, { headers: { 'Upgrade': 'my-chat-protocol', 'Connection': 'Upgrade' } }); if (res.status !== 101) throw new Error('连接升级失败'); this._startReadingStream(res.body); this._startHeartbeat(); console.log('连接成功'); } async _startReadingStream(body) { const reader = body.getReader(); let buffer = new Uint8Array(); while (true) { const { done, value } = await reader.read(); if (done) { this.onClose?.(); this._stopHeartbeat(); break; } // 拼接缓存 const newBuffer = new Uint8Array(buffer.length + value.length); newBuffer.set(buffer); newBuffer.set(value, buffer.length); buffer = newBuffer; // 尝试解码帧 while (buffer.length >= 4) { const view = new DataView(buffer.buffer); const length = view.getUint32(0); if (buffer.length >= 4 + length) { const messageBytes = buffer.slice(4, 4 + length); const message = new TextDecoder('utf8').decode(messageBytes); this.onMessage?.(message); // 裁剪缓存 buffer = buffer.slice(4 + length); } else { break; } } } } send(message) { fetch(`${this.baseUrl}/send`, { method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify({ message }) }); } _startHeartbeat() { this.heartbeatTimer = setInterval(() => { // 这里用send发送PING,实际可优化为直接通过流式连接发送 this.send('PING'); }, 30000); } _stopHeartbeat() { clearInterval(this.heartbeatTimer); } } // 使用示例 const client = new CustomChatClient('http://localhost:3000'); client.onMessage = msg => console.log('收到消息:', msg); client.onClose = () => console.log('连接断开'); client.connect().catch(err => console.error('连接失败:', err)); // 发送消息 client.send('大家好,这是自定义协议的消息!');
关键注意事项
- 粘包处理:示例中的解码逻辑是简化版,实际生产环境需要更严谨的缓存拼接和帧拆分逻辑,避免消息截断或重复解析。
- 跨域问题:服务端需要添加CORS头,允许客户端的升级请求和POST请求,比如设置
Access-Control-Allow-Origin: *、Access-Control-Allow-Headers: Upgrade, Connection。 - 错误处理:添加socket异常、网络错误的捕获逻辑,实现自动重连。
- 性能优化:对于大消息,可分帧发送;对于高频消息,可合并发送减少TCP包数量。
内容的提问来源于stack exchange,提问作者ndubs18
相关产品推荐
相关产品推荐

