消息系统的推拉模型对比写扩散与读扩散的性能分析一、一条微博发出去要推给 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 RedisTemplateString, 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 天登录的 SetString activeFollowers getActiveFollowers(userId, 7); if (activeFollowers.size() PUSH_THRESHOLD) { // 活跃粉丝不多直接全部推送 pushToUsers(activeFollowers, message); } // 非活跃粉丝 → 在他们下次登录时用 Pull 模式拉取 } } /** * 粉丝获取 Feed */ public ListMessage getFeed(String userId, int page, int size) { ListMessage feed new ArrayList(); // 第一步从自己的收件箱获取推送过来的内容 // 这些是 Push 模式写入的来自小 V 和活跃大 V 关注者的推送 String inboxKey inbox: userId; SetObject pushedMessages redis.opsForZSet() .reverseRange(inboxKey, page * size, (page 1) * size - 1); // 第二步拉取大 V 的最新内容 // 从关注列表中找到大 V粉丝 阈值的用户 SetString followings getFollowings(userId); ListString bigVs filterBigVs(followings); // 从大 V 的发件箱中拉取最新内容 for (String bigV : bigVs) { String outboxKey outbox: bigV; SetObject 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 SetString getActiveFollowers(String userId, int days) { /* */ return null; } private void pushToFollowers(String userId, Message msg) { /* */ } private void pushToUsers(SetString users, Message msg) { /* */ } private SetString getFollowings(String userId) { /* */ return null; } private ListString filterBigVs(SetString followings) { /* */ return null; } private ListMessage loadMessages(CollectionObject ids) { /* */ return null; } }五、总结消息系统的推拉模型选择不是非黑即白的。纯粹 Push 在大 V 场景下写入爆炸纯粹 Pull 在关注量大时读取缓慢。推拉结合按粉丝数量和活跃度做分治是生产环境中务实的选择。核心决策参数是推送阈值——这个值需要根据系统的写入吞吐量和读取延迟的压测数据来确定。阈值太低太多走 Pull读取变慢阈值太高太多走 Push写入压力和存储成本变大。找到这个平衡点就是消息投递系统设计的工程功力所在。
消息系统的推拉模型对比:写扩散与读扩散的性能分析
消息系统的推拉模型对比写扩散与读扩散的性能分析一、一条微博发出去要推给 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 RedisTemplateString, 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 天登录的 SetString activeFollowers getActiveFollowers(userId, 7); if (activeFollowers.size() PUSH_THRESHOLD) { // 活跃粉丝不多直接全部推送 pushToUsers(activeFollowers, message); } // 非活跃粉丝 → 在他们下次登录时用 Pull 模式拉取 } } /** * 粉丝获取 Feed */ public ListMessage getFeed(String userId, int page, int size) { ListMessage feed new ArrayList(); // 第一步从自己的收件箱获取推送过来的内容 // 这些是 Push 模式写入的来自小 V 和活跃大 V 关注者的推送 String inboxKey inbox: userId; SetObject pushedMessages redis.opsForZSet() .reverseRange(inboxKey, page * size, (page 1) * size - 1); // 第二步拉取大 V 的最新内容 // 从关注列表中找到大 V粉丝 阈值的用户 SetString followings getFollowings(userId); ListString bigVs filterBigVs(followings); // 从大 V 的发件箱中拉取最新内容 for (String bigV : bigVs) { String outboxKey outbox: bigV; SetObject 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 SetString getActiveFollowers(String userId, int days) { /* */ return null; } private void pushToFollowers(String userId, Message msg) { /* */ } private void pushToUsers(SetString users, Message msg) { /* */ } private SetString getFollowings(String userId) { /* */ return null; } private ListString filterBigVs(SetString followings) { /* */ return null; } private ListMessage loadMessages(CollectionObject ids) { /* */ return null; } }五、总结消息系统的推拉模型选择不是非黑即白的。纯粹 Push 在大 V 场景下写入爆炸纯粹 Pull 在关注量大时读取缓慢。推拉结合按粉丝数量和活跃度做分治是生产环境中务实的选择。核心决策参数是推送阈值——这个值需要根据系统的写入吞吐量和读取延迟的压测数据来确定。阈值太低太多走 Pull读取变慢阈值太高太多走 Push写入压力和存储成本变大。找到这个平衡点就是消息投递系统设计的工程功力所在。