WhatsApp 发送频率控制的令牌桶算法实现

📅 2026/7/24 12:25:47 👁️ 阅读次数 📝 编程学习
WhatsApp 发送频率控制的令牌桶算法实现

WhatsApp 发送频率控制的令牌桶算法实现

目录

  1. 为什么固定间隔的限速不够用
  2. 令牌桶的核心思想
  3. 基础实现:单线程令牌桶
  4. 进阶:多节点独立桶 + 全局配额
  5. 与调度器的集成
  6. 生产环境的落地经验
  7. 小结

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)、分布式令牌桶(跨进程/跨机器共享配额)、或者令牌借用与偿还机制,都可以在这个基础上扩展。