系统设计:即时通讯(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 机制:
- 客户端发送 → 携带消息 ID
- 服务端 ACK → 服务端收到返回 ACK
- 客户端重发 → 未收到 ACK 则定时重发(幂等)
- 服务端去重 → 按消息 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: 如何保证消息有序?
- 消息 ID 全局递增
- 单聊:按时间顺序展示
- 群聊:依赖消息 ID 排序,网络延迟可能导致乱序接收,客户端按 ID 排序展示
Q: 如何实现消息撤回?
存储撤回指令消息,客户端渲染时判断:如果消息被撤回则显示「消息已撤回」。
Q: 离线消息怎么处理?
- 推送系统(APNs、FCM、厂商推送)
- 用户上线后主动拉取(sync 机制)
Q: WebSocket 连接断了怎么恢复?
客户端定时发送心跳,超时时服务端清理连接映射。客户端断网重连后,携带上次 sync_seq 恢复消息。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。