系统设计:通知系统
通知系统是几乎所有 App 的标配:点赞、评论、私信、活动提醒都要触达用户。核心挑战:海量投递、多渠道、乱序容忍度低、且必须避免打扰用户。
1. 需求分析
场景
- 用户量:2 亿注册,日活 5000 万
- 消息量:高峰 50 万 QPS(大促、直播秒杀时刻)
- 触达渠道:App 推送(APNs/FCM)、站内信、短信、邮件、Webhook
- 投递要求:延迟可接受分钟级(非实时通信),但不能丢
核心问题
- 多渠道分发:一条业务消息如何路由到用户的多个设备/渠道?
- 去重与限流:同一用户瞬间收到 N 条相似推送如何合并?怎么防止推送轰炸?
- 可靠性:第三方通道(APNs)不稳定,失败怎么重试?保证不重不漏?
- 用户偏好:免打扰时段、渠道开关(只想要邮件不想要短信)
与 IM 的区别
通知系统不是实时通信(不做 WebSocket 长连接),而是「尽力而为、可靠投递」的异步触达系统,天然适合消息队列。实时推送可参考 IM 即时通讯系统。
2. 系统架构
业务方(点赞/评论/订单)
↓ HTTP 调用
通知 API 网关(鉴权、参数校验)
↓
通知服务(异步落库 + 入队)
↓
消息队列(Kafka,按渠道 topic 分区)
↓
渠道分发器(消费,按用户偏好路由)
┌────┬─────┬────┬──────┐
↓ ↓ ↓ ↓ ↓
推送 worker 站内信 短信 邮件 Webhook
↓
第三方推送(APNs / FCM)
↓
用户设备
3. 关键技术
3.1 消息模型与模板
消息 = 模板 + 参数,模板带渠道语言描述,方便多渠道渲染。
# 模板示例
template = {
"id": "order_shipped",
"title": "订单已发货",
"channels": {
"push": {"title": "订单已发货 🚚", "body": "您的订单 {order_id} 已发出"},
"email": {"subject": "您的订单已发货", "html": "<p>订单号:{order_id}</p>"},
"sms": {"text": "您的订单 {order_id} 已发货,请查收"}
}
}
def render(template_id, params):
t = get_template(template_id)
return {ch: t["channels"][ch].format(**params) for ch in t["channels"]}
为什么用模板? 模板把「业务方发什么」和「每个渠道怎么渲染」解耦,改文案不用改业务代码,且便于审核合规。
3.2 渠道分发器(路由)
分发器消费队列,按用户偏好 + 渠道可用性 + 优先级路由:
def route(notification):
user = get_user_preferences(notification.user_id) # 渠道开关、免打扰
if user.dnd and is_in_dnd_window(): # 免打扰时段
queue_to_delayed(notification, until=user.dnd_until)
return
for channel in ordered_channels(notification, user):
if channel.enabled(user):
send_to_channel_worker(channel, notification)
优先级队列:验证码/支付结果 → 高优先级(立即投递);营销活动 → 低优先级(错峰)。
3.3 推送通道:APNs / FCM 接入
移动推送必须走苹果 APNs 或谷歌 FCM,服务端持有设备 token。
import requests
def send_push(token, payload, channel="apns"):
if channel == "apns":
resp = requests.post(
"https://api.push.apple.com/3/device/" + token,
json={"aps": {"alert": payload["title"], "sound": "default"}},
headers={"authorization": "Bearer " + get_jwt()},
)
else: # fcm
resp = requests.post(
"https://fcm.googleapis.com/fcm/send",
json={"to": token, "notification": payload},
headers={"authorization": "key=" + FCM_KEY},
)
return resp.status_code
关键点:
- 设备 token 会失效(卸载/换机),返回
410 Gone时从设备表删除该 token。 - APNs 对请求有速率限制(按服务端 key),推送 worker 要并发受限 + 批量接口。
- token 管理与设备注册表放在 Redis/MySQL,推送时批量拉取。
3.4 限流与去重
推送轰炸是用户体验杀手,三层防护:
# 1. 单用户限流:同一用户每分钟最多 N 条(Redis 计数)
def rate_limit(user_id, channel, limit=20, window=60):
key = f"notify:rate:{user_id}:{channel}"
count = redis.incr(key)
if count == 1:
redis.expire(key, window)
return count <= limit
# 2. 内容去重:相同内容的通知在时间窗口内合并
def dedup(user_id, content_hash, window=300):
key = f"notify:dedup:{user_id}:{content_hash}"
return redis.set(key, 1, nx=True, ex=window) is True
# 3. 合并(App 端聚合):同类型多条推送折叠为一条
# 如「你关注的主播开播了 ×5」→「5 位主播开播了」
静默推送 vs 展示推送:可以在推送消息里带 data 字段让客户端静默处理,由 App 本地聚合后再弹一条,大幅降低打扰。
3.5 可靠性投递与重试
第三方通道(APNs/FCM/运营商)不可控,必须设计失败重试 + 幂等消费:
# 消费端幂等:用全局消息 ID 去重,保证「不重」
def process_message(msg):
if not redis.set(f"consumed:{msg.id}", 1, nx=True, ex=86400):
return # 已消费过,跳过
try:
send_via_channel(msg)
update_status(msg.id, "SENT")
except ChannelUnavailable:
# 1. 重试队列:指数退避(1s/2s/4s... 最多 8 次)
retry_delay(msg, backoff(times(msg.retries)))
# 2. 超过最大重试 → 死信队列 + 告警人工介入
if msg.retries >= 8:
send_to_dead_letter(msg)
投递语义:
- 至多一次:适合营销消息(丢了也无所谓)。
- 至少一次:适合交易通知(可重复但绝不能丢),靠幂等消费保证业务正确。
持久化:每条通知先落库(notification 表带状态机:PENDING → SENT → FAILED/DEAD),定时任务扫描补偿未投递的消息。
3.6 站内信与邮件
- 站内信:写扩散到收件箱表
(user_id, msg_id, read_flag),用户拉取时查询;也可用 Redis 存未读数。 - 邮件:交给专用邮件服务(SES/自建),模板渲染 + 退信处理(
bounce标记,防止持续向无效邮箱投递)。
4. 数据模型与扩展
核心表
users (user_id, device_token, channels_pref, dnd_window)
notification (id, user_id, template_id, params_json, status, created_at)
notification_log (id, notification_id, channel, status, retries, updated_at)
扩展方向
- 多语言/多地域:模板按
locale渲染,短信按运营商分片区。 - 渠道降级:高优通知若推送失败,自动降级补发短信(如支付验证码)。
- 分析追踪:记录送达/点击/转化漏斗,反哺渠道选型与频控策略。
- 连接复用:对同一用户的推送合并成一条 APNs 批量请求,降低限流命中率。
兜底方案
- 推送队列堆积告警:Kafka 消费延迟监控,消费能力不足时丢弃营销低优先级保验证码等高优消息。
- 数据一致:
notification_log与渠道回执对账,定时任务扫「超时未送达」重投。
渠道选型对比
| 渠道 | 实时性 | 成本 | 到达率 | 适用场景 |
|---|---|---|---|---|
| App 推送(APNs/FCM) | 秒级 | 低 | 中(依赖通知权限) | 通用触达、活动提醒 |
| 站内信 | 分钟级 | 极低 | 高(App 内必达) | 系统通知、账单、隐私类 |
| 短信 | 秒级 | 高 | 高 | 验证码、支付/风控告警 |
| 邮件 | 分钟级 | 低 | 中(易进垃圾箱) | 周报、营销、订阅 |
| Webhook | 秒级 | 低 | 中(依赖接收方) | 开发者事件回调 |
选型逻辑:优先「免费且到达率高」的站内信打底,关键交易用短信补强,营销走推送 + 邮件双渠道。高优验证码与低优营销必须在队列与频控上隔离。
数据规模演进
- 单机版:同步调用第三方通道,服务内线程池并发。
- 队列版:引入 Kafka,通知服务只落库 + 入队,消费端异步投递,抗峰值。
- 分布式版:按用户 ID 分片存储,渠道 worker 独立扩容,多机房就近投递。
- 智能频控版:基于用户打开率动态调频,学习用户偏好(免打扰时段、渠道顺序),减少流失。
5. 面试常见问题
Q: 通知系统怎么保证「不重复投递」?
消费端幂等:全局消息 ID + Redis SETNX 去重;数据库唯一索引 (notification_id, channel) 双保险。
Q: 用户瞬间收到大量推送怎么办?
三层:单用户频控(Redis 计数)、内容去重合并、客户端静默聚合。营销类还要做全站错峰。
Q: APNs token 失效怎么处理?
推送返回 410/Unregistered 时删除设备 token,避免反复向无效设备发送。
Q: 通知系统和 IM 的区别?
IM 是实时长连接(WebSocket),要求秒级送达与消息顺序;通知系统是异步触达(分钟级可接受),核心是渠道分发、频控与可靠性,二者常共用推送网关但业务模型不同。
Q: 如果消息队列挂了怎么办?
API 网关同步兜底:通知先落库为 PENDING,队列恢复后由定时任务扫描补投;同时 Kafka 多副本 + 消费组容灾。
Q: 定时通知(如「明天 10 点提醒」)怎么实现?
延迟队列:Kafka + 时间轮 / Redis ZSet 按触发时间排序,到点再投递到业务队列;或直接落库后由调度任务扫描到期消息入队。
Q: 业务方接入的成本如何控制?
提供统一 SDK + 消息模板中心:业务方只传 template_id + params + user_id,渠道、文案、频控全部由通知平台接管,避免每个业务各搞一套推送。
总结
通知系统的本质是把「业务事件」转成「多渠道、不丢、不打扰」的用户触达。答题时抓住三条主线即可拿高分:渠道分发(模板 + 路由 + 优先级)、防打扰(频控 + 去重 + 聚合)、可靠性(落库 + 幂等 + 重试 + 对账)。先画架构图,再逐条展开 trade-off,面试官通常会顺着这三条追问,提前准备好对应的兜底方案。
相关文章:
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。