从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 是全双工通信协议,建立连接后可以双向通信。

sequenceDiagram participant Client as 客户端 participant Server as 服务器 Client->>Server: HTTP Upgrade 请求 Server->>Client: HTTP 101 Switching Protocols Client-->>Server: WebSocket 连接建立 Client->>Server: 消息 Server->>Client: 消息 Client->>Server: 消息 Server-->>Client: Close

建立 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 用于浏览器之间的点对点通信,不需要服务器中转。

sequenceDiagram participant A as 浏览器A participant Server as 信令服务器 participant B as 浏览器B A->>Server: Offer (SDP) Server->>B: Offer (SDP) B->>Server: Answer (SDP) Server->>A: Answer (SDP) A->>Server: ICE Candidate Server->>B: ICE Candidate B->>Server: ICE Candidate Server->>A: ICE Candidate A-->>B: P2P 连接建立

建立 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/) 转载或引用必须申明原指尖魔法屋来源及源地址!