8.3 实战:排行榜计数与消息队列 本节摘要:两个端到端方案收束全书:活动排行榜用跳表加过期加分片计数器实现;订单事件流用 Stream 加消费组加 XACK 实现可恢复的消费。每个设计决策都能回溯到前几章的一个结构结论——这是数据结构透视方法的最终验收。 方案一:活动排行榜 需求:一场 7 天活动,千万用户积分排行,实时查询前 100 名与"我的排名",每天结算分区榜。 每个选型背后都站着第 2 章:ZRANGE 与 ZRANK 是跳表的范围查询与跨度累计,O(logN) 让千万成员下仍是毫秒级;分区榜按天分键,让"周期清榜"变成过期字典里的一行台账(第 5 章),不需要任何删除任务。
本节摘要:两个端到端方案收束全书:活动排行榜用跳表加过期加分片计数器实现;订单事件流用 Stream 加消费组加 XACK 实现可恢复的消费。每个设计决策都能回溯到前几章的一个结构结论——这是数据结构透视方法的最终验收。
需求:一场 7 天活动,千万用户积分排行,实时查询前 100 名与"我的排名",每天结算分区榜。
# 写入:积分变动,双写总分榜与当日分区榜 ZINCRBY rank:total 5 user:88 ZINCRBY rank:day:3 5 user:88 # 查询:前100名(分数带成员,反转取区间) ZREVRANGE rank:total 0 99 WITHSCORES # 我的排名与分数 ZREVRANK rank:total user:88 ZSCORE rank:total user:88 # 同分按先到先排:分数拼接时间戳余量 ZADD rank:total 1000005 user:99 # 100万为分数基,尾数为时序 # 每天凌晨结算后给分区榜设过期,7天后自动回收 EXPIRE rank:day:3 604800
每个选型背后都站着第 2 章:ZRANGE 与 ZRANK 是跳表的范围查询与跨度累计,O(logN) 让千万成员下仍是毫秒级;分区榜按天分键,让"周期清榜"变成过期字典里的一行台账(第 5 章),不需要任何删除任务。
读多写多的榜单要加缓存层:前 100 名结果缓存在 String,500 毫秒过期,写入路径顺手刷新——把 ZREVRANGE 的调用量砍掉九成,配合第 5 章的抖动 TTL 防集体失效。
容量与性能的预估账也一并摆出来,评审时就说这三个数。写入侧:ZINCRBY 是跳表加 dict 双结构各动一次,千万成员的榜上单次仍约几十微秒,写瓶颈通常先出现在网络往返而不是结构;读取侧:前 100 名的 ZREVRANGE 是对数级定位加百次顺序读,毫秒以内,加结果缓存后实际打到跳表的 QPS 可以压到个位数;内存侧:skiplist 编码每成员粗算数百字节,千万成员合计数 GB 级——这个数决定方案评审能不能过,提前按"预计参与人数乘留存率"算两遍。变式:预算卡死时按活跃分层,只给活跃用户建榜、长尾用户只记分不入榜,内存立省一大截。
需求:订单状态变更产生事件,三个下游(通知、统计、风控)各自独立消费,处理失败可重试,服务重启不丢任务。
上线前的压测脚本也是方案的一部分,两组核心断言:秒杀窗口模拟一万并发 ZINCRBY,验收写入无报错、榜首查询延迟在毫秒位;事件流灌入十万事件、中途杀掉一个消费者,验收接管后无丢失(组内积压清零、业务侧对账无缺口)。把这两个脚本随方案一起归档,下一次活动复用时,验收成本只剩十来分钟。
# 生产端:订单状态机每次跃迁追加一条 XADD order:events * oid 1001 from PAID to SHIPPED ts 1724301000 # 三个下游各自建组,都从头消费 XGROUP CREATE order:events notify 0 XGROUP CREATE order:events stats 0 XGROUP CREATE order:events risk 0 # 通知组消费者:取一条处理确认一条 XREADGROUP GROUP notify worker-1 COUNT 10 BLOCK 2000 STREAMS order:events > # 处理成功 XACK order:events notify 1724301000-5 # 处理失败:不ACK,留给接管流程

兜底与扩容:消费者崩溃后,其 PEL 中未确认消息由监控脚本 XAUTOCLAIM 转移给存活 worker 重投;下游处理要幂等(以事件 id 去重),因为"至少一次"投递必然伴随偶发重复;日志按 MAXLEN 100 万裁剪,内存占用有界。
幂等去重的落地虽然小,但值得写完整——它是事件流方案的保险丝:
def handle(event_id, payload): # 去重键:处理成功标记,TTL盖过最大重试周期即可 if not r.set(f"dedup:{event_id}", 1, nx=True, ex=86400): return "duplicate" # 已处理过,直接确认 try: do_business(payload) # 真正的下游动作 except Exception: r.delete(f"dedup:{event_id}") # 失败要允许重试,标记回滚 raise return "ok"
解读两处细节:去重标记与业务动作不是分布式事务,极端时序下仍可能出现"标记写了、业务没跑成",所以标记先写、失败回滚的模式只能把重复压到极小,配合下游自身的唯一约束才是闭环;TTL 取一天足够覆盖接管重投的窗口,过期自动清理不占内存。这个三十行的组合,是所有"至少一次"消息方案的标配前置件。
| 关注点 | 排行榜方案 | 事件流方案 |
|---|---|---|
| 结构支点 | 跳表加过期字典 | RADIX 日志加 PEL |
| 峰值表现 | 写 O(logN),读走结果缓存 | 追加 O(1),消费可横向扩 |
| 故障路径 | 缓存失效回源重建 | 消费者挂由接管流程兜底 |
| 容量边界 | 单键百万级成员内 | XTRIM 约束下的有限历史 |
事件流方案的容量账同样三句话:写入侧 XADD 是对数级的日志定位加追加,单次微秒级,日百万事件毫无压力;存储侧每事件数百字节,一百万条保留量对应数百 MB,XTRIM 的 MAXLEN 就是内存闸门;消费侧三个组各自独立游标,组数只影响元数据多少,不复制消息本体——"一份数据多组消费"的内存经济学,正是它对"每组一条队列"方案的降维优势。
💡 关键直觉:方案评审时把"结构支点"一栏填出来,风险就藏不住了——支点不稳(大键、无界增长、无确认),方案一定会在生产现形。