系统设计:即时通讯(IM)系统

设计支持千万在线用户的即时通讯系统,详解WebSocket长连接管理、消息可靠性保证、消息ID设计、读扩散与写扩散模型、以及多端同步策略。

系统设计:即时通讯(IM)系统

设计类似微信/钉钉的即时通讯系统,核心挑战在于:海量长连接管理、消息可靠投递、消息顺序保证与多端同步。

1. 需求分析

功能需求

  • 单聊、群聊
  • 消息类型:文本、图片、语音、文件
  • 已读回执
  • 消息撤回
  • 多端同步(手机 + PC + Web)

非功能需求

  • 日活:1 亿
  • 峰值在线:2000 万
  • 消息延迟:< 200ms(P99)
  • 消息不丢失、不重复、不乱序

2. 系统架构

客户端 → DNS → CDN(静态资源)
           ↓
        接入层(LVS + Nginx)
           ↓
    ┌──────┴──────┐
    ↓             ↓
 网关服务      路由服务
(WebSocket)(用户→网关映射)
    ↓             ↓
消息服务 ← → 用户服务
    ↓             ↓
消息队列    数据库集群
    ↓
存储服务(消息持久化)

3. 核心设计

3.1 WebSocket 长连接管理

单台服务器维护的长连接数有限(通常 10 万级),需要水平扩展:

# 网关层维护用户到网关的映射
class GatewayManager:
    def __init__(self):
        self.user_gateway = {}  # user_id → gateway_id
        self.gateway_connections = {}  # gateway_id → 连接数

    def route_to(self, user_id):
        gateway_id = self.user_gateway.get(user_id)
        return gateway_id

    def connect(self, user_id, gateway_id):
        self.user_gateway[user_id] = gateway_id
        self.gateway_connections[gateway_id] += 1

路由服务:用户A发消息给用户B,需查询B在哪个网关,定向推送。

3.2 消息 ID 设计

全局有序的消息 ID 是消息顺序和去重的基础。

方案:Snowflake 变种

| 1 bit | 41 bit 时间戳 | 10 bit 机器ID | 12 bit 序列号 |

优点:趋势递增,支持每秒 4096 × 1024 = 400 万 ID。

3.3 消息可靠性保证

QoS 机制:

  1. 客户端发送 → 携带消息 ID
  2. 服务端 ACK → 服务端收到返回 ACK
  3. 客户端重发 → 未收到 ACK 则定时重发(幂等)
  4. 服务端去重 → 按消息 ID 去重
class MessageService:
    def send_message(self, msg):
        # 1. 生成消息 ID
        msg_id = generate_msg_id()
        # 2. 存储消息(至少一次写入)
        self.store_message(msg)
        # 3. 推送接收方
        self.push_to_recipient(msg)
        return {"msg_id": msg_id}

    def handle_ack(self, msg_id):
        # 标记消息已送达
        self.mark_delivered(msg_id)

3.4 读扩散 vs 写扩散

维度读扩散(如 微信)写扩散(如 微博)
存储一份消息存储每收件人存储一份
读取拉取会话消息读取收件箱
写入简单复杂(群大时O(N))
适合单聊、小群大群、广播

混合方案:小群(< 200 人)用写扩散,大群用读扩散。

3.5 多端同步

每个设备维护消息同步位点(sync_seq):

class DeviceSync:
    def sync_messages(self, user_id, device_id, last_seq):
        # 获取该设备上次同步后的消息
        new_messages = self.get_messages_after(user_id, last_seq)
        new_seq = new_messages[-1].seq if new_messages else last_seq
        return {
            "messages": new_messages,
            "new_seq": new_seq
        }

4. 存储设计

消息存储

CREATE TABLE messages (
    msg_id BIGINT PRIMARY KEY,
    sender_id BIGINT NOT NULL,
    receiver_id BIGINT,
    group_id BIGINT,
    content TEXT,
    msg_type TINYINT,  -- 1:文本 2:图片...
    created_at TIMESTAMP,
    INDEX idx_conversation (sender_id, receiver_id, created_at)
);

近期消息缓存

最近 N 天的消息缓存在 Redis,减少数据库读取:

# 获取消息
messages = redis.get(f"recent_msgs:{conversation_id}")
if not messages:
    messages = db.query(conversation_id, limit=100)
    redis.setex(f"recent_msgs:{conversation_id}", 86400, messages)

5. 群聊特殊处理

超大群策略(> 2000 人)

  • 写扩散成本过高(每条消息复制 2000 份)
  • 改用读扩散:消息只存一份,成员拉取时读取
  • 在线成员实时推送(WebSocket),离线成员走推送服务

6. 面试常见问题

Q: 如何保证消息有序?

  1. 消息 ID 全局递增
  2. 单聊:按时间顺序展示
  3. 群聊:依赖消息 ID 排序,网络延迟可能导致乱序接收,客户端按 ID 排序展示

Q: 如何实现消息撤回?
存储撤回指令消息,客户端渲染时判断:如果消息被撤回则显示「消息已撤回」。

Q: 离线消息怎么处理?

  • 推送系统(APNs、FCM、厂商推送)
  • 用户上线后主动拉取(sync 机制)

Q: WebSocket 连接断了怎么恢复?
客户端定时发送心跳,超时时服务端清理连接映射。客户端断网重连后,携带上次 sync_seq 恢复消息。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页