8.3 实战:排行榜计数与消息队列


文档摘要

8.3 实战:排行榜计数与消息队列 本节摘要:两个端到端方案收束全书:活动排行榜用跳表加过期加分片计数器实现;订单事件流用 Stream 加消费组加 XACK 实现可恢复的消费。每个设计决策都能回溯到前几章的一个结构结论——这是数据结构透视方法的最终验收。 方案一:活动排行榜 需求:一场 7 天活动,千万用户积分排行,实时查询前 100 名与"我的排名",每天结算分区榜。 每个选型背后都站着第 2 章:ZRANGE 与 ZRANK 是跳表的范围查询与跨度累计,O(logN) 让千万成员下仍是毫秒级;分区榜按天分键,让"周期清榜"变成过期字典里的一行台账(第 5 章),不需要任何删除任务。

8.3 实战:排行榜计数与消息队列

本节摘要:两个端到端方案收束全书:活动排行榜用跳表加过期加分片计数器实现;订单事件流用 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,留给接管流程

Stream 支撑多下游消费

Stream 支撑多下游消费

兜底与扩容:消费者崩溃后,其 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 就是内存闸门;消费侧三个组各自独立游标,组数只影响元数据多少,不复制消息本体——"一份数据多组消费"的内存经济学,正是它对"每组一条队列"方案的降维优势。

💡 关键直觉:方案评审时把"结构支点"一栏填出来,风险就藏不住了——支点不稳(大键、无界增长、无确认),方案一定会在生产现形。

本节要点回顾

  • 排行榜:跳表范围查询加按天分键加过期自动清榜,读路径加结果缓存
  • 事件流:一组日志多组消费,XACK 加 XAUTOCLAIM 构成恢复闭环
  • 幂等是 Stream 方案的必修课,至少一次投递必有重复
  • 每个决策回溯一个结构结论,这是全书方法论的最终验收方式

作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U