人人都会AI编程

22.4 实时业务场景:即时通讯、消息推送、协同编辑

更新时间:2026-07-11

前面三节已经把实时通信的技术底座——WebSocket、Socket.IO 和 SSE——拆解清楚了。但技术最终要落到具体的业务里才能产生价值。本节我们就聚焦三个最典型的实时业务场景:即时通讯(IM)、消息推送和协同编辑,看看在真实项目中如何运用这些技术、会遇到哪些坑以及怎样做出合适的技术决策。

22.4.1 即时通讯(IM)

即时通讯是实时技术最经典的应用,涵盖单聊、群聊、消息回执、离线消息、多端同步等全套需求。这个场景有几个硬指标:

  • 低延迟:消息从发送到接收应在几百毫秒内完成。
  • 高可靠性:消息不能丢、不能乱序,弱网环境下要有补偿机制。
  • 双向通信:用户既要发消息,也要实时收到对方的消息。
  • 状态同步:在线状态、输入状态、已读状态等都需要频繁更新。

技术选型:WebSocket + Socket.IO

无论从哪个角度看,IM 都必须选择全双工的 WebSocket,而 Socket.IO 提供的房间、自动重连、心跳、事件封装可以极大减少重复开发量。SSE 是单向推送,无法满足用户发送消息的需求,不必考虑。

如果团队需要自己封装 WebSocket(比如追求极致性能),可以用 ws 库,但要做好重连、心跳和房间路由的编程工夫。

整体架构

一个生产可用的 IM 系统通常会拆为以下几层:

  • 接入层:Node.js 集群,每个进程维护一批 WebSocket 连接,负责消息收发、协议解析。
  • 路由层:当消息需要跨进程(比如接收方不在同一台机器)时,用 Redis Pub/Sub 或 RabbitMQ 做消息总线,把事件广播给集群中所有节点,再由持有目标连接的进程推送到客户端。
  • 存储层:用 MySQL/MongoDB 持久化消息,用 Redis 存储在线用户的连接节点映射(user:123 -> ws-server-3),保证消息能精准路由。
  • 离线消息:接收方不在线时,消息直接落库并标记未读,用户下次上线时先拉取未读队列再进入正常收发。

关键实现点

消息 ID 与幂等性
每条消息在服务端生成一个全局唯一的 ID(可用雪花算法或 UUID),客户端重连时携带最后收到的一条 ID,服务端补偿推送这段时间内丢失的消息,并过滤重复投递。

已读回执
单聊的已读回执相对简单:当打开聊天窗口时,客户端发送一个 ACK 消息,包含最后一条已读消息的 ID。群聊则需要记录每个成员已读到的消息序号,收发压力更大。

多端同步
同一个用户可能在手机和 PC 同时在线。可以将同一用户的所有连接加入一个专属房间(如 room:user:123),任何给该用户的消息都发给这个房间,这样所有设备都能收到。但要注意,已读回执和在线状态需要更细粒度处理:比如只让活动端回复在线,其他端标记为后台。

心跳与重连
Socket.IO 自带心跳,但如果用原生 WebSocket,必须自己实现 ping/pong,否则经过代理或长时间空闲连接会被切断。重连时需要客户端携带上次收到的最大消息 ID,服务端据此推增量消息。

示例代码架构(Socket.IO)

// 用户上线
io.on('connection', (socket) => {
  const userId = socket.handshake.query.userId;
  // 将用户所有连接加入个人房间
  socket.join(`user:${userId}`);
  
  // 记录在线状态
  redis.hset('online_users', userId, process.pid);
  
  // 接收私聊消息
  socket.on('private_message', (data) => {
    const msg = { id: generateMsgId(), from: userId, to: data.to, content: data.content, time: Date.now() };
    // 持久化消息
    saveMessage(msg);
    // 尝试推送给接收者(本地)
    io.to(`user:${data.to}`).emit('private_message', msg);
    // 若接收者不在本进程,通过 Redis 广播
    redis.publish('message_channel', JSON.stringify({ type: 'private', msg }));
  });
  
  // 处理离线
  socket.on('disconnect', () => {
    redis.hdel('online_users', userId);
  });
});

// 订阅 Redis 频道,接收其他节点的消息转发
redis.subscribe('message_channel');
redis.on('message', (channel, message) => {
  const data = JSON.parse(message);
  if (data.type === 'private') {
    io.to(`user:${data.msg.to}`).emit('private_message', data.msg);
  }
});

22.4.2 消息推送

消息推送与 IM 的核心区别在于:驱动方是服务器而非客户端,且大多数场景下只需要服务器向客户端单向通知,如:

  • 订单状态变更、物流更新
  • 系统公告、后台配置下发的实时开关
  • 实时行情、股票价格推送

这类场景对双向互相通信的需求不强,所以技术选型更为灵活。

SSE 与 WebSocket 如何选?

| 特性 | SSE (Server-Sent Events) | WebSocket |
|------|--------------------------|-----------|
| 通信方向 | 服务器 -> 客户端 | 双向 |
| 协议 | HTTP | 独立 ws/wss 协议 |
| 浏览器兼容 | 除 IE 外均支持 | 全部现代浏览器 |
| 重连 | 自动重连 (EventSource API) | 需手动实现或借助库 |
| 同域连接数限制 | 默认 6 个(HTTP/1.1) | 无限制 |

建议:如果业务只有服务端主动推送,且推送频率适中(每秒一条以内),SSE 是更简单、更省资源的方案。它走普通的 HTTP,能复用现有的负载均衡和鉴权体系。但若需要客户端频繁地向服务端发数据包(如用户操作反馈),那还是用 WebSocket 更直接。

Node.js 实现 SSE 推送

核心就是设置响应头 Content-Type: text/event-stream 并保持连接打开。

app.get('/stream', (req, res) => {
  res.writeHead(200, {
    'Content-Type': 'text/event-stream',
    'Cache-Control': 'no-cache',
    'Connection': 'keep-alive',
  });

  // 每个客户端分配一个 ID
  const clientId = Date.now();
  // 加入连接池
  clients.set(clientId, res);

  req.on('close', () => {
    clients.delete(clientId);
  });
});

// 在业务某处推送事件
function pushEvent(data) {
  clients.forEach((res, clientId) => {
    res.write(`data: ${JSON.stringify(data)}\n\n`);
  });
}

多进程部署时,同样可以借助 Redis Pub/Sub,但注意只推送“事件通知”而非完整数据,由客户端收到事件后再拉取最新详情,可以降低服务器出口带宽和复杂度。

推送的可靠性保障

推送消息不一定保证到达(网络抖动、连接断开),重要业务需要配合“拉取确认”机制:

  1. 推送一条消息,同时附带一个递增的消息序列号。
  2. 客户端收到后发送确认请求(或携带最后的序列号)。
  3. 服务端记录待确认列表,定期重推未确认的消息,直到 ACK 或超时丢弃。

22.4.3 协同编辑

协同编辑是最具挑战性的实时场景之一,代表的体验是 Google Docs、在线白板、多人可编辑表格等。它不仅要解决消息的可靠传递,更大难点在于多个用户同时对同一份文档进行修改时的冲突解决

核心挑战

  • 低延迟操作同步:用户的每一次插入、删除都应该在 100ms 内反映到其他人屏幕上。
  • 一致性保证:所有参与编辑者最终看到一致的内容,无论网络延迟如何。
  • 光标与选区同步:不仅要同步内容,还要同步每个参与者的光标位置和选中范围。
  • 离线编辑支持:用户可能短暂断网,恢复连接后能合并离线修改。

技术选型:OT 与 CRDT

目前工业界主要使用两种算法:

  • OT (Operational Transformation):操作变换,通过转换位置偏移来合并并发操作。经典的实现有 Google Docs 早期的引擎、ShareDB。OT 一般需要一个中心服务器进行转换和协调,难以实现纯去中心化。
  • CRDT (Conflict-free Replicated Data Types):无冲突复制数据类型,数据结构本身保证无论操作以什么顺序到达,最终结果都一致。代表库有 Yjs、Automerge。CRDT 可以中心化或去中心化,不要求单一权威服务器。

推荐方案:对于 Node.js 开发者,直接使用成熟的库可以大幅降低门槛。

  • 如果只需要文本协同编辑,Yjs + y-websocket 是一条捷径。Yjs 内置了 Y.Text、Y.Map 等数据类型,集成 prosemirror、quill 等富文本编辑器也很方便。
  • 如果需要一个中心化架构并希望更精细的权限控制,ShareDB(支持 OT)配合 JSON 文档很合适。

Node.js 在协同编辑中的角色

无论用哪种算法,Node.js 服务器主要做三件事:

  1. 接收来自客户端的操作(如插入字符、删除范围)。
  2. 对操作进行转换/合并(OT 模式)或应用 CRDT 要求(通常由库完成)。
  3. 将转换后的操作广播给同一文档的其他客户端。

因为所有连接都汇集到 Node.js 进程,事件循环的高并发能力可以维持成百上千个文档的实时同步,CPU 开销也相对可控——大部分运算由库内部的轻量级操作完成,而不是 JavaScript 处理大循环。

Yjs 示例架构

// server.js (使用 y-websocket)
const WebSocket = require('ws');
const { setupWSConnection } = require('y-websocket/bin/utils');

const wss = new WebSocket.Server({ port: 1234 });
wss.on('connection', (ws, req) => {
  setupWSConnection(ws, req, { gcEnabled: true });
});

客户端(如浏览器中使用 Yjs):

import * as Y from 'yjs';
import { WebsocketProvider } from 'y-websocket';

const ydoc = new Y.Doc();
const provider = new WebsocketProvider('ws://localhost:1234', 'room-doc-1', ydoc);
const ytext = ydoc.getText('content');

// 绑定到编辑器,如 Quill
const binding = new QuillBinding(ytext, quillEditor);

连接建立后,所有编辑操作会自动同步。Yjs 内部使用 CRDT,不需要中心仲裁,服务器只负责消息转发。

性能与扩展

协同编辑中用户操作非常频繁(打一个字就一个操作),所以不建议让每个小的编辑操作都经过数据库持久化。更好的做法是:

  • 定期保存文档快照到数据库(如每 50 个操作或每 5 秒)。
  • 利用 Redis 做操作日志缓存,需要回放历史时从最近快照开始重放操作。
  • 多进程扩展:利用 Redis Pub/Sub 或 MQ 在进程间广播操作,确保所有持有同一文档的进程都能收到转发。

常见坑点

  1. 操作太大导致消息积压:大量、高频的操作可能导致消息队列拥塞,需要做操作合并和限速。
  2. 网络延迟导致的无序操作:CRDT 天然解决此问题;OT 则依赖服务端按顺序处理操作,必须保证消息顺序(可用时间戳+版本号)。
  3. 内存占用:长时间编辑的文档其操作历史可能在内存中堆积,需要定期清理或持久化释放。
  4. 光标跳变:光标同步需要单独的低延迟通道,避免与内容同步竞争带宽。

小结:无论是即时通讯、消息推送还是协同编辑,Node.js 的事件驱动和非阻塞 I/O 都使其成为构建这类实时服务的第一梯队选择。最重要的不是选用哪个库,而是根据业务特性正确评估通信方向、可靠性和一致性需求,然后针对性地设计分布式架构和异常恢复机制。在下一章,我们将进入微服务与分布式领域,看看 Node.js 怎样承担更复杂的系统职责。