从WebSocket走到WebRTC:实时通信笔记
产品要做即时通讯,第一版用 HTTP 轮询,消息延迟肉眼可见,服务器还被空请求刷得够呛。换成 WebSocket 后聊天顺了,但音视频一上,又碰到 NAT 穿透和信令服务器这些 WebRTC 的老问题。
为什么需要实时通信
HTTP 轮询
客户端定期发起 HTTP 请求。
// 轮询实现
function pollMessages() {
fetch('/api/messages')
.then(res => res.json())
.then(messages => {
updateUI(messages);
})
.finally(() => {
setTimeout(pollMessages, 5000); // 5 秒后再次请求
});
}
pollMessages();
问题:
- 延迟高(轮询间隔内消息不会及时送达)
- 服务器压力大(即使没有消息也要响应)
- 带宽浪费
WebSocket
WebSocket 原理
WebSocket 是全双工通信协议,建立连接后可以双向通信。
建立 WebSocket 连接
const socket = new WebSocket('ws://localhost:8080/ws');
// 连接成功
socket.onopen = (event) => {
console.log('WebSocket 连接建立');
};
// 收到消息
socket.onmessage = (event) => {
const message = JSON.parse(event.data);
handleMessage(message);
};
// 连接关闭
socket.onclose = (event) => {
console.log('WebSocket 连接关闭');
};
// 连接错误
socket.onerror = (event) => {
console.error('WebSocket 错误');
};
// 发送消息
function sendMessage(text) {
socket.send(JSON.stringify({ type: 'message', text }));
}
// 关闭连接
function closeConnection() {
socket.close();
}
服务端实现
Node.js 实现:
const WebSocket = require('ws');
const wss = new WebSocket.Server({ port: 8080 });
const clients = new Set();
wss.on('connection', (ws) => {
console.log('新客户端连接');
clients.add(ws);
ws.on('message', (message) => {
const data = JSON.parse(message);
// 广播消息给所有客户端
clients.forEach(client => {
if (client.readyState === WebSocket.OPEN) {
client.send(JSON.stringify(data));
}
});
});
ws.on('close', () => {
console.log('客户端断开');
clients.delete(ws);
});
ws.on('error', (error) => {
console.error('WebSocket 错误:', error);
});
});
Go 实现:
package main
import (
"encoding/json"
"log"
"net/http"
"github.com/gorilla/websocket"
)
var upgrader = websocket.Upgrader{
CheckOrigin: func(r *http.Request) bool {
return true
},
}
type Message struct {
Type string `json:"type"`
Text string `json:"text"`
}
type Client struct {
conn *websocket.Conn
send chan Message
}
var clients = make(map[*Client]bool)
var broadcast = make(chan Message)
func handleWebSocket(w http.ResponseWriter, r *http.Request) {
conn, err := upgrader.Upgrade(w, r, nil)
if err != nil {
log.Println(err)
return
}
client := &Client{
conn: conn,
send: make(chan Message, 256),
}
clients[client] = true
go client.readPump()
go client.writePump()
}
func (c *Client) readPump() {
defer func() {
delete(clients, c)
c.conn.Close()
}()
for {
_, message, err := c.conn.ReadMessage()
if err != nil {
break
}
var msg Message
if err := json.Unmarshal(message, &msg); err != nil {
log.Println(err)
continue
}
broadcast <- msg
}
}
func (c *Client) writePump() {
defer c.conn.Close()
for {
select {
case message, ok := <-c.send:
if !ok {
return
}
data, _ := json.Marshal(message)
c.conn.WriteMessage(websocket.TextMessage, data)
}
}
}
func main() {
go func() {
for {
msg := <-broadcast
for client := range clients {
select {
case client.send <- msg:
default:
delete(clients, client)
close(client.send)
}
}
}
}()
http.HandleFunc("/ws", handleWebSocket)
log.Fatal(http.ListenAndServe(":8080", nil))
}
心跳检测
// 客户端心跳
const HEARTBEAT_INTERVAL = 30000; // 30 秒
function sendHeartbeat() {
if (socket.readyState === WebSocket.OPEN) {
socket.send(JSON.stringify({ type: 'heartbeat' }));
}
}
setInterval(sendHeartbeat, HEARTBEAT_INTERVAL);
// 服务端心跳检测
const HEARTBEAT_TIMEOUT = 60000; // 60 秒
const clientHeartbeats = new Map();
wss.on('connection', (ws) => {
const clientId = generateClientId();
clientHeartbeats.set(clientId, Date.now());
ws.on('message', (message) => {
const data = JSON.parse(message);
if (data.type === 'heartbeat') {
clientHeartbeats.set(clientId, Date.now());
}
});
ws.on('close', () => {
clientHeartbeats.delete(clientId);
});
});
// 定期检查心跳
setInterval(() => {
const now = Date.now();
clientHeartbeats.forEach((lastHeartbeat, clientId) => {
if (now - lastHeartbeat > HEARTBEAT_TIMEOUT) {
// 断开连接
const ws = getClientById(clientId);
if (ws) {
ws.close();
}
}
});
}, HEARTBEAT_TIMEOUT / 2);
断线重连
let socket = null;
let reconnectAttempts = 0;
const MAX_RECONNECT_ATTEMPTS = 5;
const RECONNECT_INTERVAL = 3000;
function connect() {
socket = new WebSocket('ws://localhost:8080/ws');
socket.onopen = () => {
console.log('WebSocket 连接建立');
reconnectAttempts = 0;
};
socket.onclose = () => {
console.log('WebSocket 连接关闭');
reconnect();
};
socket.onerror = (error) => {
console.error('WebSocket 错误:', error);
};
}
function reconnect() {
if (reconnectAttempts >= MAX_RECONNECT_ATTEMPTS) {
console.error('达到最大重连次数');
return;
}
reconnectAttempts++;
console.log(`尝试重连 (${reconnectAttempts}/${MAX_RECONNECT_ATTEMPTS})`);
setTimeout(connect, RECONNECT_INTERVAL);
}
connect();
WebRTC
WebRTC 原理
WebRTC 用于浏览器之间的点对点通信,不需要服务器中转。
建立 WebRTC 连接
// 获取本地媒体流
navigator.mediaDevices.getUserMedia({
video: true,
audio: true
})
.then(stream => {
// 显示本地视频
document.getElementById('localVideo').srcObject = stream;
// 创建 RTCPeerConnection
const peerConnection = new RTCPeerConnection({
iceServers: [
{ urls: 'stun:stun.l.google.com:19302' },
{ urls: 'stun:stun1.l.google.com:19302' }
]
});
// 添加本地流
stream.getTracks().forEach(track => {
peerConnection.addTrack(track, stream);
});
// 处理远程流
peerConnection.ontrack = (event) => {
document.getElementById('remoteVideo').srcObject = event.streams[0];
};
// 处理 ICE 候选
peerConnection.onicecandidate = (event) => {
if (event.candidate) {
// 通过信令服务器发送 ICE 候选
signalingServer.send({
type: 'ice-candidate',
candidate: event.candidate
});
}
};
// 创建 Offer
return peerConnection.createOffer();
})
.then(offer => {
// 设置本地描述
return peerConnection.setLocalDescription(offer);
})
.then(() => {
// 通过信令服务器发送 Offer
signalingServer.send({
type: 'offer',
sdp: peerConnection.localDescription
});
})
.catch(error => {
console.error('WebRTC 错误:', error);
});
信令服务器
const WebSocket = require('ws');
const wss = new WebSocket.Server({ port: 8080 });
const rooms = new Map();
wss.on('connection', (ws) => {
let currentRoom = null;
ws.on('message', (data) => {
const message = JSON.parse(data);
switch (message.type) {
case 'join':
currentRoom = message.room;
if (!rooms.has(currentRoom)) {
rooms.set(currentRoom, new Set());
}
rooms.get(currentRoom).add(ws);
break;
case 'offer':
case 'answer':
case 'ice-candidate':
// 广播给房间内其他客户端
const room = rooms.get(currentRoom);
if (room) {
room.forEach(client => {
if (client !== ws && client.readyState === WebSocket.OPEN) {
client.send(JSON.stringify(message));
}
});
}
break;
case 'leave':
if (currentRoom && rooms.has(currentRoom)) {
rooms.get(currentRoom).delete(ws);
if (rooms.get(currentRoom).size === 0) {
rooms.delete(currentRoom);
}
}
break;
}
});
ws.on('close', () => {
if (currentRoom && rooms.has(currentRoom)) {
rooms.get(currentRoom).delete(ws);
if (rooms.get(currentRoom).size === 0) {
rooms.delete(currentRoom);
}
}
});
});
踩过的坑
坑一:NAT 穿透失败
内网之间无法建立 WebRTC 连接。
解决:使用 STUN/TURN 服务器。
const peerConnection = new RTCPeerConnection({
iceServers: [
{ urls: 'stun:stun.l.google.com:19302' },
{ urls: 'stun:stun1.l.google.com:19302' },
{
urls: 'turn:turn.example.com:3478',
username: 'username',
credential: 'password'
}
]
});
坑二:消息丢失
WebSocket 连接断开时,消息丢失。
解决:使用消息队列,确保消息不丢失。
const messageQueue = [];
function sendMessage(message) {
if (socket.readyState === WebSocket.OPEN) {
socket.send(JSON.stringify(message));
} else {
// 消息入队
messageQueue.push(message);
}
}
socket.onopen = () => {
// 发送队列中的消息
while (messageQueue.length > 0) {
const message = messageQueue.shift();
socket.send(JSON.stringify(message));
}
};
坑三:性能问题
实时通信客户端太多,服务器扛不住。
解决:
- 使用消息队列(如 RabbitMQ、Kafka)
- 使用集群和负载均衡
- 使用 CDN 分发媒体流
# 使用 Redis Pub/Sub 做消息分发
redis-cli PUBLISH room:1 '{"type":"message","text":"hello"}'
选型建议
简单即时通讯
用 WebSocket 就够了。
视频会议
用 WebRTC + 信令服务器。
大规模实时通讯
WebSocket + 消息队列 + 集群。
低延迟游戏
用 WebSocket 或 UDP。
写在最后
实时通信这东西,不是技术问题,是架构问题。
WebSocket 适合:
- 聊天
- 通知
- 简单的实时更新
WebRTC 适合:
- 视频通话
- 文件传输
- 低延迟数据传输
选型之前先评估:
- 延迟要求
- 同时在线人数
- 是否需要媒体传输
- 网络环境
不是所有实时场景都需要 WebRTC,有时候 WebSocket 加优化就够用。
这次实时通信改造花了三周,从轮询到 WebSocket,再到 WebRTC。改造完成后,消息延迟从 5 秒降到 100 毫秒,用户体验提升明显。
可用性说明:本文发布于 2019 年 5 月,距今已超过五年。文中涉及的软件版本、接口、下载地址、命令参数和操作界面可能已经发生变化,部分方案在当前环境下可能失效。请结合官方最新文档核对后再操作,生产环境使用前务必先行验证。
版权声明: 本文首发于 指尖魔法屋-从WebSocket走到WebRTC:实时通信笔记(https://blog.thinkmoon.cn/post/65-realtime-communication-websocket-webrtc-practice/) 转载或引用必须申明原指尖魔法屋来源及源地址!
评论
使用 GitHub 账号登录后即可留言,支持 Markdown。