系统设计:通知系统

从零设计一个亿级用户的通知系统:推送(APNs/FCM)、站内信、邮件多渠道分发,消息模板管理、去重与限流、可靠性投递与重试、幂等消费,包含架构图与核心代码。

系统设计:通知系统

通知系统是几乎所有 App 的标配:点赞、评论、私信、活动提醒都要触达用户。核心挑战:海量投递、多渠道、乱序容忍度低、且必须避免打扰用户。

1. 需求分析

场景

  • 用户量:2 亿注册,日活 5000 万
  • 消息量:高峰 50 万 QPS(大促、直播秒杀时刻)
  • 触达渠道:App 推送(APNs/FCM)、站内信、短信、邮件、Webhook
  • 投递要求:延迟可接受分钟级(非实时通信),但不能丢

核心问题

  1. 多渠道分发:一条业务消息如何路由到用户的多个设备/渠道?
  2. 去重与限流:同一用户瞬间收到 N 条相似推送如何合并?怎么防止推送轰炸?
  3. 可靠性:第三方通道(APNs)不稳定,失败怎么重试?保证不重不漏?
  4. 用户偏好:免打扰时段、渠道开关(只想要邮件不想要短信)

与 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)

扩展方向

  1. 多语言/多地域:模板按 locale 渲染,短信按运营商分片区。
  2. 渠道降级:高优通知若推送失败,自动降级补发短信(如支付验证码)。
  3. 分析追踪:记录送达/点击/转化漏斗,反哺渠道选型与频控策略。
  4. 连接复用:对同一用户的推送合并成一条 APNs 批量请求,降低限流命中率。

兜底方案

  • 推送队列堆积告警:Kafka 消费延迟监控,消费能力不足时丢弃营销低优先级保验证码等高优消息。
  • 数据一致:notification_log 与渠道回执对账,定时任务扫「超时未送达」重投。

渠道选型对比

渠道实时性成本到达率适用场景
App 推送(APNs/FCM)秒级低中(依赖通知权限)通用触达、活动提醒
站内信分钟级极低高(App 内必达)系统通知、账单、隐私类
短信秒级高高验证码、支付/风控告警
邮件分钟级低中(易进垃圾箱)周报、营销、订阅
Webhook秒级低中(依赖接收方)开发者事件回调

选型逻辑:优先「免费且到达率高」的站内信打底,关键交易用短信补强,营销走推送 + 邮件双渠道。高优验证码与低优营销必须在队列与频控上隔离。

数据规模演进

  1. 单机版:同步调用第三方通道,服务内线程池并发。
  2. 队列版:引入 Kafka,通知服务只落库 + 入队,消费端异步投递,抗峰值。
  3. 分布式版:按用户 ID 分片存储,渠道 worker 独立扩容,多机房就近投递。
  4. 智能频控版:基于用户打开率动态调频,学习用户偏好(免打扰时段、渠道顺序),减少流失。

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,面试官通常会顺着这三条追问,提前准备好对应的兜底方案。


相关文章:

继续阅读

探索更多技术文章

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

全部文章 返回首页