消息中心案例

wen java案例 1

本文目录导读:

消息中心案例

  1. 业务场景与需求分析(案例背景)
  2. 整体技术架构分层
  3. 核心难点与解决方案(重点)
  4. 核心数据表设计(简版)
  5. 关键流程时序图(文字版)
  6. 可靠性保障(案例分析)
  7. 总结与思考扩展

消息中心是几乎所有互联网产品的“标配”,也是后端架构中非常经典的案例,它涉及推送通道消息存储已读/未读状态实时性以及高并发削峰

下面我从业务建模技术架构核心流程典型方案四个维度,以“类钉钉/企业微信”的企业级消息中心为案例来拆解。


业务场景与需求分析(案例背景)

假设我们要为某SaaS公司设计一个全渠道消息中心,需支持:

  1. 站内信(Web端小红点、App端通知列表)。
  2. App Push(安卓厂商通道 + iOS APNs)。
  3. 短信/邮件(低频、重要告警)。
  4. WebSocket实时推送(客服聊天、协同编辑提醒)。

核心痛点(非功能需求):

  • 高并发写入:双11大促时,运营群发消息可能瞬间产生千万级消息记录。
  • 读扩散与写扩散选择:如果10万人收到同一条系统通知,是存1条还是存10万条?
  • 多端已读同步:Web端读了,App端小红点要消失。

整体技术架构分层

这是一张经典的“消息中心”架构图(文字版):

[触发源] -> [接入层] -> [消息处理层] -> [存储层] -> [推送层] -> [客户端]
  运营后台       API网关        消息分类/过滤      MySQL/Redis     推送服务        Web/App
  业务系统      (鉴权/限流)     优先级/去重       ES/消息队列     (Push/WS)       
  定时任务                      持久化策略                         |
                                 |                               v
                                 +------------------------> [用户触达层]
                                                            (短信/邮件)

核心难点与解决方案(重点)

这是面试或架构设计中必问的点,也是这个案例的精华部分。

存储模型:读扩散 vs 写扩散(关键抉择)

这是消息中心最核心的设计决策。

  • 写扩散(推送模型):用户A发消息给用户B,立即在B的收件箱表里插入一条记录。

    • 优点:用户查询收件箱极快(只查自己的),已读状态天然隔离。
    • 缺点:如果是全员群发(例如10万人),需插入10万条记录,DB压力巨大。
    • 场景单聊、群聊(小群)、点对点通知
  • 读扩散(拉取模型):只存一份消息内容,用户读时去“他的会话列表”里拉取。

    • 优点:写入极快(只存1条)。
    • 缺点:需要维护“用户-消息”的游标或索引,且需额外存储已读状态(Redis Bitmap)。
    • 场景系统通知、运营公告、超大群

本案例采用的混合模式(业界标准)

对于“系统通知/运营活动”:采用读扩散,只在notice表存1条内容,在user_notice_mark表(或Redis)里记录last_read_id(已读位置),用户拉取时,只拉取ID大于last_read_id的公告,小红点逻辑只需比对last_read_id和当前最大ID即可,这样万级用户广播只需要写1行。

对于“业务交互消息”(审批、@提醒):采用写扩散,因为这类消息量小且必须精准送达。

已读/未读状态的高效存储

在写扩散模型下,已读状态是每条记录的一个字段,但在读扩散模型下(一篇长文有10万人看过,不可能在文章表上加10万个字段),采用 Redis Bitmap(位图)

  • Keyread_status:{type}:{notice_id}read_status:notice:12345)。
  • Offset:用户ID(或用户的递增序号)。
  • Value:0或1。
  • 查询当前用户是否已读:GETBIT read_status:notice:12345 {userId}
  • 统计阅读人数:BITCOUNT read_status:notice:12345,性能极高且节省内存(10万用户只需12.5KB)。

消息推送的“削峰填谷”(防轰炸)

如果业务系统直接调消息中心接口,10万用户同时下单会导致瞬间10万条推送请求打到App。

引入“优先级队列”与“批量聚合”

  1. 所有推送请求先进 MQ(RocketMQ/Kafka),不直接调第三方推送(APNs/华为)。
  2. 推送Worker 消费消息,按用户维度设备维度聚合。
  3. 去重合并:如果1秒内用户收到了3条“点赞”通知,合并为1条“你收到3条新点赞”。
  4. 流量控制:针对短信/邮件通道,通过令牌桶算法限制发送速率,防止被运营商风控。

多端同步(Web/App/小程序)的实时性

  • 方案:通过 WebSocket(长连接) 网关维持与客户端的连接。
  • 消息产生后,不直接推送内容,而是推送一个 “新消息通知信号”(包含message_idcursor)。
  • 客户端收到信号后,通过HTTP API拉取具体消息内容。

核心数据表设计(简版)

这是消息中心的“地基”。

表1:message(消息内容表 - 全局唯一)

-- 为了读扩散存的实体内容
CREATE TABLE `message` (
  `id` bigint(20) NOT NULL AUTO_INCREMENT,
  `biz_type` varchar(32) NOT NULL COMMENT '业务类型:SYSTEM/APPROVE/CHAT', varchar(200) DEFAULT NULL,
  `content` text COMMENT '消息体(JSON或富文本)',
  `sender_id` bigint(20) DEFAULT NULL,
  `send_time` datetime NOT NULL,
  `extra` json DEFAULT NULL COMMENT '扩展字段(跳转链接、图片)',
  PRIMARY KEY (`id`),
  KEY `idx_biz_type_sendtime` (`biz_type`, `send_time`)
) ENGINE=InnoDB COMMENT='消息内容表';

表2:user_message_box(收件箱表 - 写扩散的产物)

-- 数据量大,必须分库分表(按user_id哈希)
CREATE TABLE `user_message_box` (
  `id` bigint(20) NOT NULL AUTO_INCREMENT,
  `user_id` bigint(20) NOT NULL,
  `message_id` bigint(20) NOT NULL,
  `status` tinyint(4) NOT NULL DEFAULT '0' COMMENT '0未读 1已读 2删除',
  `read_time` datetime DEFAULT NULL,
  PRIMARY KEY (`id`),
  UNIQUE KEY `uk_user_msg` (`user_id`, `message_id`), -- 唯一约束
  KEY `idx_user_status` (`user_id`, `status`) -- 查询未读列表
) ENGINE=InnoDB COMMENT='用户收件箱表';

对于那些超大群发通知,无需写入user_message_box,只需在Redis维护last_read_id


关键流程时序图(文字版)

场景:运营发一条全员公告(读扩散)

  1. Admin 点击群发 -> 接入层校验权限。
  2. 服务端写入 message 表,状态为ACTIVE
  3. 服务端将 messageIdtarget 放入MQ。
  4. 调度Worker 消费MQ -> 调用 WebSocket Server,广播给所有在线用户(只发信号:你有一条新公告)。
  5. App 收到信号 -> 调用 GET /api/v1/notices?lastId=xxx
  6. 服务端比对传参lastIdmessage表最大ID -> 返回最新公告内容。
  7. App展示小红点,用户点击查看 -> 调用 PUT /read
  8. 服务端在Redis Bitmap中设置该用户位置为1,小红点消失。

场景:用户A给用户B发送一条业务提醒(写扩散)

  1. 用户A触发审批动作。
  2. 业务系统调用消息中心SDK。
  3. MQ 异步消费落库:
    • message 表插入一条内容。
    • user_message_box 表插入一条user_id=B的记录。
  4. 推送服务查找B的设备Token,通过APNs/厂商通道推送通知栏消息。
  5. 若B在线,通过WebSocket直接下发新消息数量变更。

可靠性保障(案例分析)

这是消息中心不能忽视的部分,也是体现架构完备性的地方。

  • 最终一致性:使用MQ解耦,即使推送服务宕机,消息也会积压在MQ中,恢复后继续推送,不丢失业务数据。
  • 消息去重:为防止网络抖动导致重复推送,message表有biz_token(业务幂等ID),DB层面加唯一索引,消息重复提交时直接报错忽略。
  • 降级方案:如果Redis(Bitmap)崩溃,系统返回“已读状态未知”,不阻塞用户拉取消息列表;如果WebSocket断开,则降级为轮询HTTP接口(长轮询)。

总结与思考扩展

这个案例最核心的价值在于解决了“读”与“写”的平衡问题: 它根据消息的广度(点对点 vs 广播)自适应选择写扩散读扩散,并采用Bitmap解决海量用户的已读存储问题。

扩展思考(面试加分项):

  1. 如果用户量达到亿级user_message_box表如何拆分?(按user_id模16分库,按create_time做冷热分离)。
  2. 如果要做消息撤回(如钉钉的24小时内撤回),在读扩散模型下如何实现?(更新message表状态为RECALLED,客户端拉取时遇到该状态则本地替换文案)。
  3. 如何做通知频控?(同一用户1分钟内最多收到2条APP Push,超过则折叠进通知栏)。

希望这个案例对你有帮助,如果你需要针对某个特定模块(比如WebSocket长连接设计或分库分表策略)深入展开,可以告诉我。

抱歉,评论功能暂时关闭!