WhatsApp 发送频率控制的令牌桶算法实现
WhatsApp 发送频率控制的令牌桶算法实现
目录
- 为什么固定间隔的限速不够用
- 令牌桶的核心思想
- 基础实现:单线程令牌桶
- 进阶:多节点独立桶 + 全局配额
- 与调度器的集成
- 生产环境的落地经验
- 小结
1. why 固定间隔的限速不够用
前面好几篇文章都提到过time.sleep(8)这种最原始的限速方式。它简单、好理解,在消息量不大的时候也确实能跑。
但它有几个明显的短板:
| 场景 | 固定间隔的问题 | 你真正想要的行为 |
|---|---|---|
| 前半天没怎么发,下午突然来了一大批 | 还是傻等 8 秒一条,白白浪费上午攒下的额度 | 允许短时间 burst(突发)消耗积攒额度 |
| 某个时段平台比较空闲,想多发一点 | 不行,间隔是写死的 | 能动态调整速率 |
| 多个账号共用一个间隔参数 | 快的号被拖慢了,慢的号还是太快 | 每个号独立的节奏控制 |
**令牌桶(Token Bucket)**就是为解决这些问题设计的经典算法。它的核心思想很简单:
想象一个桶,里面装着"令牌"。每秒往桶里放 N 个令牌(补充速率),桶最多装 M 个(容量上限)。发一条消息就消耗一个令牌。有令牌就能发,没令牌就等着。
这个模型完美覆盖了上面三个痛点:桶里积攒了令牌就可以 burst,改补充速率就能调速,每个账号一个桶就互不干扰。
2. 令牌桶的核心思想
先搞清楚几个关键概念:
| 参数 | 含义 | 类比 |
|---|---|---|
| rate(补充速率) | 每秒往桶里放多少个令牌 | 水龙头流速 |
| capacity(容量) | 桶最多能装多少个令牌 | 桶的大小 |
| tokens(当前令牌数) | 桶里现在有多少可用令牌 | 当前水位 |
| burst(突发能力) | capacity 决定了最大突发量 | 满桶一次能用多少 |
举个例子:
rate = 0.125 tokens/s(即 8 秒补 1 个令牌,相当于之前sleep(8)的效果)capacity = 10(桶最多存 10 个令牌)
这意味着:
- 平稳状态:每 8 秒发 1 条(和
sleep(8)一样); - 突发能力:如果之前 80 秒都没发,桶满了(10 个令牌),可以连续发 10 条,然后再回到每 8 秒 1 条;
- 上限约束:不管攒多久,永远不可能在 1 秒内发出超过 10 条。
这就是令牌桶比固定间隔强大的地方:它允许合理的突发,但把突发的上界锁死了。
3. 基础实现:单线程令牌桶
importtimeimportthreadingfromdataclassesimportdataclass,field@dataclassclassTokenBucket:"""令牌桶限速器"""rate:float# 补充速率(tokens/second)capacity:int# 桶容量(最大突发量)_tokens:float=field(default=0.0,init=False)_last_refill:float=field(default=0.0,init=False)_lock:threading.Lock=field(default_factory=threading.Lock,init=False)def__post_init__(self):self._tokens=float(self.capacity)# 初始满桶self._last_refill=time.monotonic()self._lock=threading.Lock()def_refill(self):"""补充令牌(调用时根据 elapsed 时间计算应补多少)"""now=time.monotonic()elapsed=now-self._last_refillifelapsed>0:# 补充量 = 速率 × 经过时间,但不能超过容量increment=self.rate*elapsed self._tokens=min(self.capacity,self._tokens+increment)self._last_refill=nowdefconsume(self,tokens:int=1)->tuple[bool,float]:""" 尝试消费 tokens 个令牌。 返回 (是否成功, 需要等待的秒数)。 """withself._lock:self._refill()ifself._tokens>=tokens:self._tokens-=tokensreturnTrue,0.0# 令牌不够,计算还需要等多久deficit=tokens-self._tokens wait_time=deficit/self.ratereturnFalse,wait_timedefwait_and_consume(self,tokens:int=1)->float:""" 阻塞式消费:如果令牌不够就等到够为止。 返回实际等待的时间。 """whileTrue:success,wait=self.consume(tokens)ifsuccess:return0.0time.sleep(wait)@propertydefavailable_tokens(self)->float:withself._lock:self._refill()returnself._tokensdef__repr__(self):returnf"TokenBucket(rate={self.rate}/s, cap={self.capacity}, tokens={self.available_tokens:.1f})"核心思路:
_refill()是惰性计算的,不是真的起一个定时器每秒加令牌。而是在每次consume()调用时,根据距离上次补充过了多时间来一次性算完。这样零额外线程开销。consume()是非阻塞的,立刻告诉你能不能发;wait_and_consume()是阻塞版的,不够就自动等。- 用了
time.monotonic()而不是time.time(),因为前者不受系统时钟调整的影响(比如 NTP 校时不会导致令牌突然暴增或归零)。
坑点提示:如果你在多线程环境使用同一个 TokenBucket 实例,必须确保每次操作都在_lock保护下完成。上面的代码已经做了这件事,但如果你之后扩展功能(比如批量 consume),记得也加锁。
4. 进阶:多节点独立桶 + 全局配额
单机单桶解决了"一个号怎么控制节奏"。但实际场景下你通常有多个账号,每个号的速率不一样,而且还有一个全局上限。
4.1 多桶管理器
fromtypingimportOptional@dataclassclassBucketConfig:account_id:strrate:float# 该账号的补充速率capacity:int# 该账号的桶容量daily_limit:int# 该账号的全天总额度(独立于令牌桶)daily_sent:int=0# 今日已发送classMultiBucketManager:"""多节点独立令牌桶管理器"""def__init__(self):self._buckets:dict[str,TokenBucket]={}self._configs:dict[str,BucketConfig]={}self._global_rate:Optional[float]=None# 全局速率上限(可选)self._global_bucket:Optional[TokenBucket]=Nonedefadd_account(self,config:BucketConfig):"""注册一个账号及其桶配置"""bucket=TokenBucket(rate=config.rate,capacity=config.capacity)self._buckets[config.account_id]=bucket self._configs[config.account_id]=configdefset_global_limit(self,rate:float,capacity:int):"""设置全局速率限制(所有账号共享)"""self._global_rate=rate self._global_bucket=TokenBucket(rate=rate,capacity=capacity)deftry_send(self,account_id:str)->tuple[bool,str]:""" 尝试为指定账号获取发送许可。 返回 (是否允许, 原因说明) """# 1. 检查账号是否存在ifaccount_idnotinself._buckets:returnFalse,f"未知账号:{account_id}"config=self._configs[account_id]# 2. 检查日额度ifconfig.daily_sent>=config.daily_limit:returnFalse,f"日额度已满 ({config.daily_sent}/{config.daily_limit})"# 3. 检查该账号的令牌桶ok,wait=self._buckets[account_id].consume(1)ifnotok:returnFalse,f"该账号令牌不足,需等待{wait:.1f}s"# 4. 检查全局桶(如果配置了的话)ifself._global_bucket:gok,gwait=self._global_bucket.consume(1)ifnotgok:# 全局不允许,归还刚才从账号桶拿走的令牌self._buckets[account_id]._tokens+=1# 归还returnFalse,f"全局令牌不足,需等待{gwait:.1f}s"# 所有检查通过config.daily_sent+=1returnTrue,"OK"defget_status(self,account_id:str=None)->dict:"""获取当前状态概览"""result={}ifaccount_id:aid=account_id bucket=self._buckets.get(aid)cfg=self._configs.get(aid)ifbucketandcfg:result[aid]={"available_tokens":round(bucket.available_tokens,1),"daily_sent":cfg.daily_sent,"daily_limit":cfg.daily_limit,"daily_remaining":max(0,cfg.daily_limit-cfg.daily_sent),"rate":cfg.rate,"capacity":cfg.capacity,}else:foraid,bucketinself._buckets.items():cfg=self._configs[aid]result[aid]={"available_tokens":round(bucket.available_tokens,1),"daily_sent":cfg.daily_sent,"daily_limit":cfg.daily_limit,"daily_remaining":max(0,cfg.daily_limit-cfg.daily_sent),}ifself._global_bucket:result["_global"]={"available_tokens":round(self._global_bucket.available_tokens,1),"rate":self._global_rate,}returnresult两层限速的关系:
请求进入 ↓ ① 日额度检查(硬上限,每天 N 条) ← 最外层门禁 ↓ 通过 ② 账号令牌桶(控制瞬间节奏) ← 中层:你能 burst 多猛 ↓ 通过 ③ 全局令牌桶(控制总体输出) ← 最内层:所有人一起不能超 ↓ 通过 ✅ 发送!三层各管各的:日额度防止单号一天打太多,账号桶防止一秒内爆发太猛,全局桶防止所有号加起来把平台打爆。
4.2 日额度自动重置
importdatetimedefdaily_reset_task(manager:MultiBucketManager):"""每日重置任务(应该在 UTC 0 点或本地 0 点触发)"""forcfginmanager._configs.values():old_sent=cfg.daily_sent cfg.daily_sent=0print(f"[日重置]{cfg.account_id}:{old_sent}→ 0")可以用系统的 crontab 或者 Python 的schedule库来每天跑一次。
5. 与调度器的集成
把令牌桶嵌入到之前的 MessageScheduler 里非常自然:
# 在 MessageScheduler.__init__ 里增加:bucket_mgr=MultiBucketManager()foracctinaccount_configs:# 根据账号等级分配不同的 rate 和 capacitylevel=acct.get("level","normal")iflevel=="new":rate,capacity,limit=0.083,5,50# 新号:12s/条,burst 5,日限 50eliflevel=="warm":rate,capacity,limit=0.125,8,150# 预热号:8s/条,burst 8,日限 150else:# activerate,capacity,limit=0.2,15,300# 成熟号:5s/条,burst 15,日限 300bucket_mgr.add_account(BucketConfig(account_id=acct["phone"],rate=rate,capacity=capacity,daily_limit=limit))# 设置全局限制(可选):所有号加起来每秒不超过 2 条bucket_mgr.set_global_limit(rate=2.0,capacity=20)# 在 run_task 的发送循环里:formsginbatch:account_id=msg.get("task_id","default")# 先申请令牌allowed,reason=bucket_mgr.try_send(account_id)ifnotallowed:print(f" ⚠ [{account_id}] 发送受限:{reason}")continue# 跳过这条,处理下一条# 令牌够了,执行实际发送try:success=send_fn(msg["recipient"],msg["content"])ifnotsuccess:# 发送失败要不要归还令牌?看你的策略# 一般选择不归还(因为请求已经发出去了,占用了平台的配额)passexceptExceptionase:print(f" ✗ 异常:{e}")# 注意:这里不再需要 time.sleep(fixed_interval)# 因为令牌桶本身就已经控制了节奏和原来time.sleep(8)方式的对比:
| 维度 | 固定间隔 | 令牌桶 |
|---|---|---|
| 代码量 | 1 行 | ~80 行(但封装好后也是 1 行调用) |
| Burst 能力 | 无 | 有(受 capacity 控制) |
| 动态调速率 | 改常量重启 | 改 rate 属性即可 |
| 多账号隔离 | 要自己写逻辑 | 天然支持 |
| 日额度 | 要自己计数 | 内建 |
| 可观测性 | 无 | get_status()一目了然 |
6. 生产环境的落地经验
我们以 WAWarmer 的频率控制模块为例,看它的令牌桶是怎么用的。
① 它用了三级分层但不是全开
实际上它默认只开了账号桶 + 日额度两层,全局桶默认关闭(通过set_global_limit不调用来跳过)。原因是它的节点规模通常在 10 个以内,全局超限的概率不高。但如果某个客户自己配了 30+ 节点,系统会建议开启全局桶。
② rate 和 capacity 是热可配的
它把这些参数存在 SQLite 配置表里而不是代码中。运营同学可以通过 API 或 YAML 文件修改某個账号的 rate(比如从 0.125 调到 0.1),不需要重启服务也不需要改代码,下一次_refill就会生效。这在应对平台临时收紧额度的场景下非常有用,收到预警后 30 秒内就能完成全网降速。
③ 它有一个"借令牌"机制
当某个重要消息必须立即发送而当前桶空了的时候,它可以配置"允许预支":本次先发出去,后续从补充的令牌里扣还(表现为接下来一段时间内实际速率会比配置的 rate 更低)。这个机制默认关闭,只在标记为"高优先级"的任务中才启用。
7. 小结
令牌桶是频率控制领域经过几十年验证的经典方案。它比固定间隔多的那几行代码换来的是:突发容忍、动态调速、多租户隔离、可观测性。核心就三件事:
- 惰性补充:不做定时器,每次消费时按 elapsed 时间算,零开销;
- 双层约束:桶容量管瞬时爆发,日限额管全天总量;
- 可组合:单桶、多桶、全局桶按需叠加,架构不变。
如果你的团队也在做类似的发送系统,建议直接替换掉现有的time.sleep():
- 先用一个 TokenBucket 替换单个账号的固定间隔(改动不到 10 行业务代码);
- 加一层 MultiBucketManager 给每个账号独立的桶参数;
- 接入看板展示
get_status()数据,让运营看到实时的令牌余量。
这套方案从引入到全量替代大约 1 个工作日。如果要加基于机器学习的自适应速率调节(根据历史 429 反馈自动调 rate)、分布式令牌桶(跨进程/跨机器共享配额)、或者令牌借用与偿还机制,都可以在这个基础上扩展。