消息系统的推拉模型对比:写扩散与读扩散的性能分析

消息系统的推拉模型对比:写扩散与读扩散的性能分析

消息系统的推拉模型对比:写扩散与读扩散的性能分析

一、一条微博发出去,要推给 500 万个粉丝,瞬间写入 500 万条记录

社交平台的消息投递有两种基本模型。写扩散(Push):用户发布内容时,立即把这条内容写入所有粉丝的收件箱。粉丝查看时直接从自己的收件箱读取。读扩散(Pull):用户发布内容时只写入自己的发件箱。粉丝查看时,先去关注列表找到所有关注的人,再从他们的发件箱中拉取内容合并排序。

两种模型各有鲜明的优缺点。写扩散写入成本高但读取成本低,读扩散写入成本低但读取成本高。但问题不是这么简单的二选一——在一个有百万粉丝的大 V 和只有 50 个粉丝的普通用户并存的平台上,统一的 Push 或 Pull 模型都会在某些场景下严重失效。这就引出了推拉结合模型。

二、写扩散的优缺点

优势

  • 读取极快。粉丝看 Feed 只需要一次查询自己收件箱的操作。
  • 实现简单。不需要在读取时做多源合并和排序。

劣势

  • 大 V 发布时写入量爆炸。1000 万粉丝意味着一次发布需要 1000 万次写入操作(即使分批,也是一次巨大的 IO 冲击)。
  • 存储浪费。一个粉丝一年都不登录,但他的收件箱里存了大量从未被阅读的内容。
  • 删除和修改的成本高。如果大 V 删除了某条内容,需要从所有已推送的粉丝收件箱中删除或标记。

优化手段

  • 分批异步推送。大 V 的内容不立即全部推送,而是分成多批慢慢推,利用消息队列缓冲写入压力。
  • 只存 ID。收件箱只存内容 ID 和排序分数,不存完整内容。读取时根据 ID 批量获取。

三、读扩散的优缺点

优势

  • 写入成本极低。发布内容只写一次到自己的发件箱。
  • 无存储浪费。只有被请求时才读取。
  • 删除修改方便。直接在自己的发件箱中操作即可。

劣势

  • 读取时需要聚合多个关注用户的发件箱。如果一个用户关注了 2000 个人,每次刷新 Feed 需要 2000 次查询(或一次复杂的多表查询)。
  • 实现复杂。需要做多源内容的去重、排序、过滤。

优化手段

  • 缓存活跃用户的发件箱。把最近活跃的用户的最近内容缓存在 Redis 中,读取时只需访问缓存命中的用户发件箱。
  • 对关注列表做异步维护。在后台维护每个用户的关注列表的快照,读取时直接用快照。

四、推拉结合:实际生产的务实选择

主流社交平台大多采用推拉结合模式。核心策略是根据用户的粉丝数量和活跃度进行分类:

大 V 的粉丝:读扩散。因为大 V 发布频繁且粉丝数量大,推送成本过高。
普通用户的粉丝:写扩散。因为粉丝数量少,推送成本低。
活跃粉丝:写扩散。因为有高频的读取需求,推送后读取体验更好。
非活跃粉丝:读扩散。因为可能很久不登录,推送的存储是被浪费的。

/** * 推拉结合模式的消息投递服务 * * 分治策略: * - 小 V(粉丝 < 阈值):推送模式 → 写入粉丝收件箱 * - 大 V(粉丝 >= 阈值):拉取模式 → 写入自己发件箱,粉丝拉取 * * 阈值的选择需要根据系统容量做调优 * 典型值:粉丝数 < 5000 走推送,>= 5000 走拉取 */ @Service public class HybridMessageService { // 推送/拉取的粉丝数分界阈值 // 这个值的设定基于压测:当粉丝超过 5000 时,推送延迟开始显著增加 private static final int PUSH_THRESHOLD = 5000; @Resource private RedisTemplate<String, Object> redis; /** * 用户发布内容 */ public void publish(String userId, Message message) { // 写入自己的发件箱(无论 Push 还是 Pull 都要做) // 发件箱存储格式:Sorted Set, score = 发布时间戳 String outboxKey = "outbox:" + userId; redis.opsForZSet().add(outboxKey, message.getId(), message.getTimestamp()); // 获取粉丝数量 long followerCount = getFollowerCount(userId); if (followerCount < PUSH_THRESHOLD) { // 小 V:写扩散 pushToFollowers(userId, message); } else { // 大 V:读扩散 // 只对活跃粉丝做推送(最近 7 天登录的) Set<String> activeFollowers = getActiveFollowers(userId, 7); if (activeFollowers.size() < PUSH_THRESHOLD) { // 活跃粉丝不多,直接全部推送 pushToUsers(activeFollowers, message); } // 非活跃粉丝 → 在他们下次登录时用 Pull 模式拉取 } } /** * 粉丝获取 Feed */ public List<Message> getFeed(String userId, int page, int size) { List<Message> feed = new ArrayList<>(); // 第一步:从自己的收件箱获取推送过来的内容 // 这些是 Push 模式写入的(来自小 V 和活跃大 V 关注者的推送) String inboxKey = "inbox:" + userId; Set<Object> pushedMessages = redis.opsForZSet() .reverseRange(inboxKey, page * size, (page + 1) * size - 1); // 第二步:拉取大 V 的最新内容 // 从关注列表中找到大 V(粉丝 >= 阈值的用户) Set<String> followings = getFollowings(userId); List<String> bigVs = filterBigVs(followings); // 从大 V 的发件箱中拉取最新内容 for (String bigV : bigVs) { String outboxKey = "outbox:" + bigV; Set<Object> recent = redis.opsForZSet() .reverseRange(outboxKey, 0, 19); // 合并到 feed 中 feed.addAll(loadMessages(recent)); } // 第三步:合并 + 排序 feed.addAll(loadMessages(pushedMessages)); feed.sort((a, b) -> Long.compare( b.getTimestamp(), a.getTimestamp())); return feed.subList(0, Math.min(size, feed.size())); } private long getFollowerCount(String userId) { /* Redis 获取 */ return 0; } private Set<String> getActiveFollowers(String userId, int days) { /* */ return null; } private void pushToFollowers(String userId, Message msg) { /* */ } private void pushToUsers(Set<String> users, Message msg) { /* */ } private Set<String> getFollowings(String userId) { /* */ return null; } private List<String> filterBigVs(Set<String> followings) { /* */ return null; } private List<Message> loadMessages(Collection<Object> ids) { /* */ return null; } }

五、总结

消息系统的推拉模型选择不是非黑即白的。纯粹 Push 在大 V 场景下写入爆炸,纯粹 Pull 在关注量大时读取缓慢。推拉结合按粉丝数量和活跃度做分治,是生产环境中务实的选择。核心决策参数是推送阈值——这个值需要根据系统的写入吞吐量和读取延迟的压测数据来确定。阈值太低,太多走 Pull,读取变慢;阈值太高,太多走 Push,写入压力和存储成本变大。找到这个平衡点,就是消息投递系统设计的工程功力所在。