异步并发怎么限流?一次 gather 一万个任务会发生什么?
简化版
asyncio.gather(*[fetch(u) for u in urls]) 会同时创建并启动所有任务——一万个 URL 就是一万个协程同时发起连接。后果是四重的:① 瞬间打满本地的文件描述符(每个连接一个 fd,ulimit -n 通常 1024);② 一万个协程对象和缓冲区同时驻留内存;③ 对下游发起 DDoS(对方限流、封 IP,或者被你打挂);④ 事件循环里堆积海量回调,延迟飙升。而且 gather 是先把所有任务创建出来再等,所以「内存峰值」在第一时间就到了。控制并发的标准工具是 asyncio.Semaphore:sem = asyncio.Semaphore(20),在每个任务内部用 async with sem: 包住真正的 IO——注意信号量要包在任务里,而不是包在创建任务的地方,否则并发根本没被限制住。但要分清两个不同的概念:Semaphore 限制的是「同时进行的数量(并发度)」,而不是「每秒多少次(速率)」——如果下游要求「QPS ≤ 100」,你用 Semaphore(100) 但每个请求只要 10ms,实际 QPS 会飙到 10000。限速率要用令牌桶(按时间补充令牌)。更完整的方案是「有界队列 + 固定数量的 worker」:生产者往队列里放任务、N 个 worker 协程消费——它同时实现了并发控制和背压(队列满时生产者会被阻塞,天然向上游传递压力)。核心记忆:gather 是无界并发;Semaphore 管并发数、令牌桶管速率;大批量用队列 + worker 池。
详细版
四种并发控制手段:
| 手段 | 控制什么 | 适合 |
|---|---|---|
asyncio.Semaphore(N) | 同时进行的数量 | 保护下游连接数、限制本地资源 |
| 令牌桶 / 漏桶 | 单位时间的次数(QPS) | 对方按 QPS 限流时 |
| 有界队列 + worker 池 | 并发数 + 背压 | 大批量、流式、长期运行 |
| 分批(chunk) | 每批的并发 | 简单批处理 |
import asyncio, time, random
# ① ★反面:无界并发★
async def bad(urls):
return await asyncio.gather(*(fetch(u) for u in urls)) # ★1 万个同时发起★
# → fd 耗尽 / 内存暴涨 / 下游被打挂 / 事件循环堆积
# ② ★Semaphore 限并发(注意信号量的位置!)★
sem = asyncio.Semaphore(20)
async def fetch_limited(url):
async with sem: # ★① 在任务内部★
return await fetch(url)
async def good(urls):
return await asyncio.gather(*(fetch_limited(u) for u in urls))
# ★1 万个协程仍然被创建,但同时只有 20 个在真正发请求★
# ✗ 常见错误:把信号量放在创建任务的地方
async def wrong(urls):
tasks = []
for u in urls:
async with sem: # ★✗ 这里只限制了"创建任务"的速度★
tasks.append(asyncio.create_task(fetch(u)))
return await asyncio.gather(*tasks) # ★任务一创建就跑,并发完全没限制★
# ③ ★并发数 ≠ 速率★
# Semaphore(100) + 每个请求 10ms → ★实际 QPS = 100/0.01 = 10000★
# 想限制 QPS 必须用令牌桶:
class RateLimiter:
"""令牌桶:每秒补充 rate 个令牌,最多攒 burst 个"""
def __init__(self, rate: float, burst: int | None = None):
self.rate = rate
self.capacity = burst or int(rate)
self.tokens = float(self.capacity)
self.updated = time.monotonic()
self._lock = asyncio.Lock()
async def acquire(self, n: int = 1):
async with self._lock:
while True:
now = time.monotonic()
self.tokens = min(self.capacity,
self.tokens + (now - self.updated) * self.rate)
self.updated = now
if self.tokens >= n:
self.tokens -= n
return
wait = (n - self.tokens) / self.rate # ★算出还要等多久★
await asyncio.sleep(wait)
limiter = RateLimiter(rate=100) # 100 QPS
async def fetch_rate_limited(url):
await limiter.acquire()
return await fetch(url)
# ④ ★有界队列 + worker 池(大批量的标准方案)★
async def worker(q: asyncio.Queue, results: list):
while True:
item = await q.get()
try:
results.append(await fetch(item))
except Exception as e:
results.append(e)
finally:
q.task_done() # ★必须调,否则 join() 永远不返回★
async def run_pool(items, concurrency=20):
q = asyncio.Queue(maxsize=concurrency * 2) # ★有界 → 天然背压★
results = []
workers = [asyncio.create_task(worker(q, results)) for _ in range(concurrency)]
for it in items:
await q.put(it) # ★队列满时在这里等 → 背压★
await q.join() # 等所有任务处理完
for w in workers:
w.cancel() # ★清理 worker★
await asyncio.gather(*workers, return_exceptions=True)
return results
# ⑤ 分批处理(最简单,但批内快批间慢时利用率低)
async def in_batches(items, size=50):
out = []
for i in range(0, len(items), size):
out += await asyncio.gather(*(fetch(x) for x in items[i:i+size]),
return_exceptions=True)
return out
# ⑥ ★流式消费:as_completed + Semaphore(省内存,能提前处理结果)★
async def stream(items, concurrency=20):
sem = asyncio.Semaphore(concurrency)
async def one(x):
async with sem:
return await fetch(x)
for fut in asyncio.as_completed([one(x) for x in items]):
yield await fut # ★谁先完成先处理★
⚠️ 三个最容易做错的地方:① 信号量必须包在「真正做 IO 的那段代码」外面,而不是包在「创建任务」的地方。写成
async with sem: tasks.append(create_task(...))时,信号量只限制了创建任务的速度——任务一旦创建就立刻开始跑,并发完全没被限制。正确写法是把async with sem:放进任务协程的内部,包住await fetch(...)。②Semaphore限的是并发数,不是速率。Semaphore(100)配上每次 10ms 的请求,实际 QPS 是 100 / 0.01 = 10000——如果对方限的是「每秒 100 次」,你照样会被限流甚至封禁。并发数控制的是「同时占用多少资源」,速率控制的是「单位时间做多少次」,两者经常要同时用(Semaphore+ 令牌桶)。③gather即使配了信号量,一万个协程对象仍然会被同时创建——信号量只挡住了「同时执行 IO」的数量,但每个协程的对象、闭包、参数都已经在内存里了。对超大批量(十万级以上)应该用「有界队列 + 固定 worker 池」:任务按需从队列取,内存占用恒定,而且队列满时生产者会被阻塞,天然形成背压向上游传递压力。
完整版教学
一、无界并发的真实后果
await asyncio.gather(*(fetch(u) for u in urls)) # urls 有 10000 个
★ 会发生什么(按时间顺序):
t0 创建 10000 个协程对象 → ★内存瞬间上涨★(每个协程 + 闭包 + 参数)
t0+ 10000 个任务全部进入事件循环的 ready 队列
t1 10000 个 TCP 连接同时发起
→ ★本地 fd 耗尽★:ulimit -n 默认 1024 → OSError: Too many open files
→ 端口耗尽:本地临时端口范围约 28000 个(★大量 TIME_WAIT 时更少★)
t2 下游:10000 个并发请求
→ 对方限流(429)、封 IP、或者★被你打挂★
t3 事件循环里堆积海量回调 → ★调度延迟飙升★
→ 连心跳、超时检测都被拖慢
★ 量化一下内存:
一个协程对象 + 帧 + 参数 ≈ 1~3 KB
一个 aiohttp 请求的缓冲区 ≈ 几十 KB(★响应体大时更多★)
10000 并发 × 50KB = ★500 MB★(还没算响应内容)
→ 100 万个任务?直接 OOM
★ 更隐蔽的问题:错误雪崩
下游被打挂 → 全部超时 → 你的重试逻辑触发 → ★并发翻倍★
→ 彻底压垮下游,且自己也 OOM
★ "但我本地测 100 个没问题啊"
→ 100 和 10000 是两个世界:
fd、端口、内存、下游承受力都是★阶跃式★崩溃,不是线性劣化
→ ★永远不要写没有上限的并发★
★ 正确的心智模型:
┌────────────────────────────────────────────────┐
│ 你要控制的是三个独立的量: │
│ ① ★并发数★ 同时有多少个请求在飞(占资源) │
│ ② ★速率★ 每秒发起多少个(对方的限流规则) │
│ ③ ★内存占用★ 同时有多少任务对象/数据在内存里 │
│ Semaphore 只解决 ①;令牌桶解决 ②;队列+worker 解决 ③│
└────────────────────────────────────────────────┘
无界并发的后果是阶跃式的而不是线性劣化。按时间顺序:先是一万个协程对象同时创建、内存瞬间上涨,接着一万个 TCP 连接同时发起——本地 fd 耗尽(ulimit -n 默认 1024,报 Too many open files)、临时端口耗尽(约 28000 个,大量 TIME_WAIT 时更少),然后下游被限流、封 IP 甚至被打挂,最后事件循环里堆积海量回调、调度延迟飙升(连心跳和超时检测都被拖慢)。量化一下:一个 aiohttp 请求的缓冲区约几十 KB,一万并发就是 500MB,一百万个任务直接 OOM。还有更隐蔽的错误雪崩:下游挂了 → 全部超时 → 重试逻辑触发 → 并发翻倍 → 彻底压垮。关键的心智模型是:你要控制的是三个独立的量——并发数(同时多少个在飞)、速率(每秒发起多少)、内存占用(同时多少任务在内存里),而 Semaphore 只解决第一个。
二、Semaphore:控制并发数
★ 原理:一个计数器 + 等待队列
sem = asyncio.Semaphore(20)
async with sem: # ★计数减 1;减到 0 时后来者在这里挂起★
await do_io()
# 退出时计数加 1,唤醒一个等待者
★ 三个正确用法:
① 包住 IO(★最常见★)
async def task(x):
async with sem:
return await fetch(x)
② 作为装饰器/包装器复用
def limited(sem):
def deco(fn):
async def wrapper(*a, **kw):
async with sem:
return await fn(*a, **kw)
return wrapper
return deco
③ 手动 acquire/release(★需要跨函数持有时★)
await sem.acquire()
try:
...
finally:
sem.release() # ★必须在 finally★
★ 四个常见错误:
✗ ① 信号量包在了"创建任务"的地方(★最常见★)
async with sem:
tasks.append(asyncio.create_task(work())) # ★只限制了创建速度★
✗ ② 在同步函数里用 threading.Semaphore
→ ★asyncio 要用 asyncio.Semaphore★(threading 版会阻塞整个事件循环)
✗ ③ 信号量作用域太大,把不该限制的也包进去了
async with sem:
data = await fetch(x) # 需要限制
result = parse(data) # ★CPU 计算不需要占着信号量★
await save(result) # 这是另一个下游,应该用另一个信号量
✓ 拆开:不同的下游用不同的信号量,CPU 部分不占信号量
✗ ④ 全局单例信号量在多个事件循环间共享
→ ★asyncio 的同步原语绑定在创建它的事件循环上★
→ 3.10+ 移除了 loop 参数,但仍然★不能跨 loop 使用★
✓ 在 async 函数里创建,或用惰性初始化
★ 并发数怎么定:
受限于下游 → 按对方文档/协议(如 API 允许 50 并发)
受限于本地资源 → fd 上限、内存、连接池大小
★受限于延迟★ → 并发数 ≈ 目标 QPS × 平均延迟(利特尔法则)
例:目标 200 QPS、平均延迟 100ms → ★并发 20★
→ ★压测确定★:从小到大扫,看吞吐和 P99 延迟的拐点
★ Semaphore vs 连接池:
aiohttp 的 TCPConnector(limit=100) 本身就是一种并发限制
→ ★两者会叠加★:Semaphore(20) + limit=100 → 实际并发 20
→ 别重复限制到互相打架;通常让连接池限制略大于信号量
asyncio.Semaphore 的原理是「计数器 + 等待队列」,用法就是 async with sem: 包住 IO。四个常见错误要记住:① 把信号量包在「创建任务」的地方(只限制了创建速度,最常见);② 在异步代码里用 threading.Semaphore(会阻塞整个事件循环,必须用 asyncio.Semaphore);③ 作用域太大(把 CPU 计算和另一个下游的调用也包进去了,应该拆成不同的信号量、CPU 部分不占);④ 跨事件循环共享(asyncio 的同步原语绑定在创建它的 loop 上)。并发数怎么定有个实用公式——利特尔法则:并发数 ≈ 目标 QPS × 平均延迟(目标 200 QPS、延迟 100ms → 并发 20),最终还是要压测扫参数、找吞吐和 P99 延迟的拐点。最后注意 Semaphore 和连接池的限制会叠加(aiohttp 的 TCPConnector(limit=100) 本身就是并发限制),别重复限制到互相打架。
三、并发数 vs 速率:两个不同的问题
★ 核心区别:
★并发数(concurrency)★ = 同一时刻有多少个请求"在飞"
→ 控制的是★资源占用★(连接、fd、内存、下游的处理槽位)
→ 工具:Semaphore、连接池、worker 数量
★速率(rate)★ = 单位时间内发起多少次
→ 控制的是★对方的限流规则★(QPS、每分钟配额)
→ 工具:令牌桶、漏桶、滑动窗口
★ 两者的换算(利特尔法则):
★并发数 = 速率 × 平均延迟★
Semaphore(100) + 平均延迟 10ms → 速率 = 100/0.01 = ★10000 QPS★
Semaphore(100) + 平均延迟 1s → 速率 = 100/1 = ★100 QPS★
→ ★同一个并发限制,在不同延迟下对应完全不同的 QPS★
→ 所以"下游按 QPS 限流"时,光用 Semaphore ★控制不住★
★ 令牌桶(token bucket):最常用的限速算法
概念:桶里以固定速率 rate 补充令牌,最多攒 capacity 个;
每次请求消耗 1 个令牌,没令牌就等
★ 特点:★允许突发★(攒够的令牌可以一次性用掉,最多 capacity 个)
实现要点:
- 不需要定时器:★按"距上次的时间差"惰性补充令牌★
- 用 ★time.monotonic()★(不受系统改时间影响)
- 用 asyncio.Lock 保护 tokens 的读-改-写
- 令牌不够时:sleep((需要-现有)/rate) ★而不是轮询★
★ 漏桶(leaky bucket):
请求以固定速率流出 → ★完全平滑,不允许突发★
实现更简单:记录 next_time,每次 sleep 到 next_time 再 += 间隔
class LeakyBucket:
def __init__(self, rate): self.interval = 1/rate; self.next = time.monotonic()
async def acquire(self):
now = time.monotonic()
if self.next > now: await asyncio.sleep(self.next - now)
self.next = max(now, self.next) + self.interval
★ 滑动窗口:
精确匹配"每 60 秒最多 N 次"这类规则(记录最近 N 次的时间戳)
→ 内存占 O(N),但语义最贴合大多数 API 的限流文档
★ 选择:
┌────────────────┬────────────────────────────────┐
│ 对方规则 │ 用什么 │
├────────────────┼────────────────────────────────┤
│ 最多 N 个并发 │ ★Semaphore★ │
│ 每秒 N 次 │ ★令牌桶★(允许突发)/ 漏桶(平滑)│
│ 每分钟 N 次 │ ★滑动窗口★ 或 令牌桶(rate=N/60) │
│ 两者都有 │ ★Semaphore + 令牌桶叠加★ │
└────────────────┴────────────────────────────────┘
★ 分布式限流:
单进程的令牌桶在多副本部署下会 ★N 倍超发★
→ 用 Redis 实现(INCR + EXPIRE,或 Lua 脚本做原子令牌桶)
→ 或者把总配额按副本数平分(简单但利用率低)
并发数和速率是两个不同的问题:并发数控制「资源占用」(同时多少个在飞,用 Semaphore、连接池),速率控制「单位时间的次数」(用令牌桶、漏桶)。两者的换算是利特尔法则:并发数 = 速率 × 平均延迟——所以 Semaphore(100) 在 10ms 延迟下是 10000 QPS、在 1 秒延迟下是 100 QPS,同一个并发限制在不同延迟下对应完全不同的 QPS,这就是「下游按 QPS 限流时光用 Semaphore 控制不住」的原因。令牌桶是最常用的限速算法(按时间差惰性补充令牌、不需要定时器、用 time.monotonic()、允许突发),漏桶完全平滑不允许突发,滑动窗口最贴合「每 60 秒最多 N 次」这类 API 文档的语义。要注意分布式场景:单进程的令牌桶在多副本部署下会 N 倍超发,必须用 Redis 做全局限流或把配额按副本数平分。
四、队列 + worker 池:大批量的标准方案
★ 为什么需要它(Semaphore 解决不了的问题):
gather + Semaphore:★所有协程对象仍然同时存在于内存★
→ 100 万个任务 = 100 万个协程对象 = ★OOM★
队列 + worker:★只有 N 个 worker 协程 + 队列里的 M 个待办项★
→ 内存恒定,与总任务数无关
★ 标准实现:
async def worker(name, q, out):
while True:
item = await q.get()
try:
out.append(await process(item))
except Exception as e:
logging.exception("处理 %s 失败", item)
out.append(e)
finally:
q.task_done() # ★必须调用,否则 join() 永远不返回★
async def run(items, concurrency=20, queue_size=100):
q = asyncio.Queue(maxsize=queue_size) # ★有界!★
out = []
workers = [asyncio.create_task(worker(f"w{i}", q, out))
for i in range(concurrency)]
try:
for it in items: # ★可以是生成器/异步流★
await q.put(it) # ★队列满时在这里等 = 背压★
await q.join() # 等队列清空
finally:
for w in workers: w.cancel()
await asyncio.gather(*workers, return_exceptions=True)
return out
★ ★背压(backpressure)★——这是队列方案的核心价值:
生产者(读文件/拉数据库/收网络流)比消费者快时:
✗ 无界队列 → 内存无限增长 → ★OOM★
✓ 有界队列 → put() 阻塞 → ★生产者自动放慢★ → 压力向上游传递
→ "背压"就是"下游忙不过来时,让上游慢下来"
★ maxsize 怎么定:够 worker 不空转即可,通常 ★concurrency × 2~5★
★ 几个关键细节:
① ★task_done() 必须调用★(放在 finally 里),否则 q.join() 永远挂起
② ★worker 里要吞掉异常★,否则一个 worker 挂了并发度就降了
③ 结束方式两种:
- q.join() 等队列清空(★推荐★)
- 或往队列放 N 个哨兵值 None,worker 见到就 return
④ ★worker 用完要 cancel★,否则 while True 会一直挂着(任务泄漏)
⑤ 结果收集:append 到 list 是安全的(★单线程事件循环★),
但如果需要顺序,要带上索引
★ 三种方案的选择:
┌────────────────────┬──────────────┬──────────────┬────────────┐
│ │ 内存 │ 背压 │ 实现复杂度 │
├────────────────────┼──────────────┼──────────────┼────────────┤
│ gather + Semaphore │ ★O(总任务数)★│ ❌ 无 │ ★最简单★ │
│ 分批 gather │ O(批大小) │ ❌ 无 │ 简单 │
│ ★队列 + worker★ │ ★O(并发数)★ │ ✅ ★有★ │ 中等 │
└────────────────────┴──────────────┴──────────────┴────────────┘
→ 任务数 < 几千:gather + Semaphore 足够
→ 任务数很大 / 流式输入 / 需要背压:★队列 + worker★
★ 3.11+ 用 TaskGroup 管理 worker(更安全):
async with asyncio.TaskGroup() as tg:
for i in range(concurrency):
tg.create_task(worker(i, q, out))
for it in items:
await q.put(it)
await q.join()
raise StopWorkers # 或用哨兵让 worker 自然退出
★ TaskGroup 保证退出时 worker 都被清理
队列 + worker 池解决的是 Semaphore 解决不了的问题:gather + Semaphore 虽然限制了同时执行的 IO 数,但所有协程对象仍然同时存在于内存——100 万个任务就是 100 万个协程对象,直接 OOM;而队列方案只有 N 个 worker 协程 + 队列里的 M 个待办项,内存与总任务数无关。它的核心价值是背压:有界队列在满时会让 put() 阻塞,生产者自动放慢,压力向上游传递——这是无界方案根本做不到的。实现有五个关键细节:task_done() 必须放在 finally 里调用(否则 q.join() 永远挂起)、worker 里要吞掉异常(否则一个 worker 挂了并发度就降了)、用 q.join() 或哨兵值结束、worker 用完要 cancel(否则 while True 一直挂着造成任务泄漏)、结果收集时如果需要顺序要带索引。选择标准很简单:任务数几千以内用 gather + Semaphore,任务数很大或需要背压就用队列 + worker。
五、重试、退避与熔断
★ 重试会让并发失控(★最容易忽略的一点★):
200 个并发 + 每个失败重试 3 次 → 峰值可能是 ★800 个请求★
→ 下游本来就慢/挂了,你的重试把它彻底压垮 → ★重试风暴★
✓ 对策:
- ★重试也要走同一个 Semaphore/限流器★(别绕过)
- ★指数退避 + 抖动★
- ★熔断★:连续失败就快速失败,不再打下游
★ 指数退避 + 抖动(jitter):
async def retry(fn, attempts=3, base=0.5, cap=10):
for i in range(attempts):
try:
return await fn()
except RetryableError:
if i == attempts - 1:
raise
delay = min(cap, base * (2 ** i))
delay = random.uniform(0, delay) # ★全抖动★
await asyncio.sleep(delay)
★ 为什么必须加抖动:
没有抖动 → 所有失败的请求★在同一时刻★重试 → ★同步化的尖峰★
→ 下游刚缓过来又被打挂 → 反复震荡
★ AWS 的经典文章推荐"full jitter":sleep(random(0, min(cap, base*2^i)))
★ 哪些错误该重试(★别无脑重试★):
✓ 网络超时、连接错误、502/503/504、429(★但要看 Retry-After★)
✗ 400/401/403/404(★重试无意义★)
✗ ★非幂等的写操作★(重试可能导致重复下单)
→ 要重试就必须带幂等键
★ 429 的正确处理:读 Retry-After 头,按它等待(★别自己猜★)
★ 熔断器(circuit breaker):
三个状态:
CLOSED(正常) → 失败率超阈值 → OPEN
OPEN(熔断) → ★直接快速失败,不打下游★ → 冷却时间后 → HALF_OPEN
HALF_OPEN(试探)→ 放少量请求过去 → 成功则 CLOSED,失败则回 OPEN
→ 价值:★给下游恢复的时间★,同时让自己快速失败而不是全部超时堆积
★ 超时是限流的一部分(★经常被忘记★):
没有超时 → 慢请求一直占着信号量的槽位 → ★有效并发降为 0★
✓ 每个请求都要有超时:
async with sem:
async with asyncio.timeout(5): # ★3.11+★
return await fetch(url)
# 或 await asyncio.wait_for(fetch(url), timeout=5)
★ 超时值要分层递减:上游 10s > 本层 5s > 下游 3s
★ 完整的"受控并发请求"模板:
class Client:
def __init__(self, concurrency=20, qps=100):
self.sem = asyncio.Semaphore(concurrency) # ★并发★
self.limiter = RateLimiter(qps) # ★速率★
async def get(self, url):
async with self.sem:
await self.limiter.acquire()
async with asyncio.timeout(5): # ★超时★
return await self._retry(url) # ★退避重试★
重试是最容易让并发失控的地方:200 个并发配上每个失败重试 3 次,峰值可能是 800 个请求——下游本来就慢,你的重试把它彻底压垮,形成重试风暴。对策有三条:重试也要走同一个信号量和限流器(别绕过)、指数退避 + 抖动、熔断。抖动是必须的——没有抖动时所有失败请求会在同一时刻重试,形成同步化的尖峰,下游刚缓过来又被打挂、反复震荡(AWS 推荐 full jitter:sleep(random(0, min(cap, base*2^i))))。还要区分哪些错误该重试:网络超时、502/503/504 可以重试,4xx(除 429)重试无意义,非幂等的写操作必须带幂等键才能重试,而 429 应该读 Retry-After 头而不是自己猜。最后一个经常被忘的点:超时是限流的一部分——没有超时的话,慢请求会一直占着信号量的槽位、有效并发降为 0,所以每个请求都要有超时,且超时值要分层递减(上游 10s > 本层 5s > 下游 3s)。
六、实践清单
★ 决策流程:
任务数 < 100 → 直接 gather(★但仍要有超时★)
任务数 100~几千 → ★gather + Semaphore★
任务数很大 / 流式 → ★有界队列 + worker 池★(有背压)
对方按 QPS 限流 → ★叠加令牌桶★
多副本部署 → ★分布式限流(Redis)★
★ 并发数怎么定:
① 利特尔法则估算:★并发 ≈ 目标 QPS × 平均延迟★
② 看下游限制(API 文档、连接池上限、DB 最大连接数)
③ 看本地限制(ulimit -n、内存、CPU)
④ ★压测扫参数★:并发从 5→10→20→50→100,看吞吐拐点和 P99
→ ★吞吐不再上升但 P99 开始恶化的点,就是最优并发★
★ 检查清单:
□ ★没有任何一处是无界并发★(grep 一下裸的 gather)
□ 信号量包在 ★IO 外面★而不是"创建任务"外面
□ ★每个请求都有超时★(否则慢请求会占死信号量)
□ 分清 ★并发数 vs QPS★,按对方规则选工具
□ 重试★走同一个限流器★ + 指数退避 + ★抖动★
□ 429 按 ★Retry-After★ 等待
□ 大批量用 ★有界队列★(背压),队列 size ≈ 并发 × 2~5
□ worker 里 ★吞掉异常★ + ★finally 里 task_done()★
□ 结束时 ★cancel worker★,别泄漏任务
□ 多副本部署时确认限流是★全局的★还是每副本的
□ 监控:★在途请求数、队列长度、限流等待时间、429 次数★
★ 监控要点(出问题时能救命):
- ★队列长度持续增长★ = 消费能力不足 → 加 worker 或降上游速率
- 限流等待时间上升 = 打到了速率上限
- ★在途请求数长期贴着上限★ = 下游变慢了(延迟上升)
- 429/超时比例上升 = 该降速或熔断了
★ 一句话总结:
★"永远不要写无界并发;用 Semaphore 控制资源占用、
用令牌桶匹配对方的速率规则、用有界队列获得背压和恒定内存;
每个请求都要有超时,重试必须退避加抖动并走同一个限流器。"★
决策流程很清晰:任务数 100 以内直接 gather(但仍要有超时)、几千用 gather + Semaphore、很大或流式用有界队列 + worker 池、对方按 QPS 限流就叠加令牌桶、多副本部署要用 Redis 做分布式限流。并发数怎么定:先用利特尔法则估算(并发 ≈ 目标 QPS × 平均延迟),再看下游和本地的硬限制,最后压测扫参数——吞吐不再上升但 P99 开始恶化的那个点就是最优并发。检查清单里最容易漏的是:每个请求都有超时(否则慢请求占死信号量)、重试走同一个限流器、worker 的 task_done() 放在 finally 里、结束时 cancel worker 别泄漏任务。监控上要盯四个指标:队列长度持续增长(消费能力不足)、限流等待时间上升(打到速率上限)、在途请求数长期贴着上限(下游变慢了)、429 和超时比例上升(该降速或熔断)。
记忆钩子:**「★asyncio.gather(一万个任务) 是无界并发★,后果是阶跃式的:①本地 fd 耗尽(ulimit -n 默认 1024)和临时端口耗尽 ②一万个协程对象同时驻留内存(一个 aiohttp 请求缓冲约几十 KB → 500MB)③对下游 DDoS(被限流/封 IP/打挂)④事件循环堆积回调、延迟飙升;而且重试还会让并发翻倍形成★重试风暴★。★三个要控制的量是独立的:并发数(同时多少在飞)、速率(每秒多少次)、内存占用★。★Semaphore 只管并发数★——两个最常见的错:①★把 async with sem 包在『创建任务』的地方★(只限制了创建速度,任务一创建就跑)②★以为它能限 QPS★(利特尔法则:★并发 = 速率 × 平均延迟★,Semaphore(100) 配 10ms 延迟就是 10000 QPS)。★限速率要用令牌桶★(按时间差惰性补令牌、用 time.monotonic、允许突发)或漏桶(完全平滑),『每分钟 N 次』这类规则用滑动窗口最贴合;★多副本部署时单进程限流会 N 倍超发★,要用 Redis 做全局限流。★超大批量用『有界队列 + 固定 worker 池』★——因为 gather+Semaphore 的★所有协程对象仍然同时在内存里★,而队列方案内存恒定(只有 N 个 worker + 队列里的 M 项),更重要的是★有界队列在满时让 put() 阻塞、生产者自动放慢=背压★;实现细节:★task_done() 必须放 finally★(否则 join 永远挂起)、worker 里吞异常、结束时 cancel worker。★每个请求都必须有超时★(否则慢请求占死信号量、有效并发降为 0),超时值分层递减;重试要★指数退避 + 抖动★(没抖动会让所有失败请求同时重试形成同步尖峰)、★走同一个限流器★、4xx 不重试、429 按 Retry-After 等、非幂等写操作要带幂等键。并发数用★压测扫参数找『吞吐不再上升但 P99 开始恶化』的拐点★。」*
七、常见误区与追问
- 误区:
await asyncio.gather(*tasks)会自动帮我控制并发。gather没有任何并发限制——它把所有传入的协程全部包装成 Task 并立即交给事件循环,一万个 URL 就是一万个协程同时发起连接。后果是阶跃式的:本地文件描述符耗尽(ulimit -n默认 1024,报OSError: Too many open files)、临时端口耗尽、内存暴涨(一万个协程对象加缓冲区轻松几百 MB)、下游被打挂或封 IP、事件循环堆积海量回调导致整体延迟飙升。而且「本地测 100 个没问题」完全不能说明问题——100 和 10000 之间是资源的阶跃式崩溃而不是线性劣化。任何面向外部的批量操作都必须有明确的并发上限。 - 误区:把
async with sem:写在创建任务的循环里就限制了并发。 这是最常见的实现错误:async with sem: tasks.append(asyncio.create_task(work()))—— 信号量在这里只保护了「创建 Task 这个动作」,而create_task一执行任务就立刻被调度运行了,信号量在下一行就被释放。结果是限制了「创建任务的速度」(毫无意义,因为创建本身几乎不耗时),实际并发完全没有上限。正确写法是把信号量放进任务协程的内部,包住真正的 IO:async def limited(x): async with sem: return await fetch(x)。判断方法很简单:信号量应该在「阻塞的 await」外面,而不是在「创建协程/任务」外面。 - 误区:
Semaphore(100)就等于限制了 100 QPS。Semaphore限制的是「同时进行的数量」,不是「单位时间的次数」。根据利特尔法则,并发数 = 速率 × 平均延迟——同样是Semaphore(100),如果每个请求只要 10ms,实际速率是 100/0.01 = 10000 QPS;如果每个请求要 1 秒,才是 100 QPS。所以当下游的限流规则是「每秒最多 100 次」时,光靠信号量根本控制不住,请求延迟一变快你就会被限流甚至封禁。要限速率必须用令牌桶/漏桶/滑动窗口,并且通常要和信号量叠加使用(信号量保护资源占用、令牌桶匹配对方规则)。反过来也一样:只限 QPS 而不限并发,遇到下游变慢时在途请求会无限堆积。 - 误区:用了
Semaphore之后,内存就不会有问题了。 信号量只挡住了「同时执行 IO」的数量,所有协程对象仍然在gather的那一刻被全部创建出来——一百万个任务就是一百万个协程对象、闭包和参数同时驻留内存,即使其中 999980 个在信号量上排队等待。对超大批量或流式输入,正确方案是「有界队列 + 固定数量的 worker」:worker 数量固定、队列大小固定,内存占用与总任务数无关;而且有界队列在满时会让生产者的put()阻塞,天然形成背压把压力传递回上游(读文件、拉数据库的速度自动放慢)。这是gather无论如何都做不到的——gather需要一个「有限的、已知的」任务列表,而队列方案可以处理无限流。 - 误区:请求失败了就重试,重试逻辑和限流是两回事。 重试会让实际并发翻倍甚至数倍:200 个并发配上每个失败重试 3 次,峰值可能达到 800 个请求;而重试往往发生在下游已经变慢或出故障的时候——你的重试正好在对方最脆弱时加倍施压,形成重试风暴,把本来能自愈的抖动变成彻底的雪崩。三条纪律:① 重试必须走同一个
Semaphore和限流器(不能绕过);② 必须指数退避 + 抖动——没有抖动时所有失败请求会在同一时刻重试,形成同步化的尖峰(AWS 推荐 full jitter:sleep(random.uniform(0, min(cap, base * 2**i))));③ 熔断——连续失败就快速失败、不再打下游,给对方恢复的时间。另外要区分可重试的错误:4xx(除 429)重试无意义,非幂等的写操作必须带幂等键。 - 追问:并发数应该设多少?有没有可参考的方法? 分四步。① 用利特尔法则估算:
并发数 ≈ 目标吞吐 × 平均延迟——目标 200 QPS、平均延迟 100ms,理论并发就是 20;这给你一个起点数量级。② 看硬限制:下游 API 文档写明的并发上限、数据库连接池大小、对方的 QPS 配额,以及本地的ulimit -n、内存、连接池limit。③ 压测扫参数:并发从 5 → 10 → 20 → 50 → 100 逐档测,记录吞吐量和 P99 延迟——吞吐不再上升但 P99 开始明显恶化的那个点就是最优并发(再往上加只会增加排队时间,不会提升吞吐)。④ 留余量并动态观察:生产环境把并发设在拐点的 70%~80%,同时监控「在途请求数」——如果它长期贴着上限,说明下游变慢了,应该告警而不是继续加压。要避免的做法是「拍脑袋设一个很大的数」或「设成 CPU 核数」——IO 密集型的并发数和 CPU 核数没有关系。 - 追问:令牌桶和漏桶有什么区别,怎么选? 令牌桶是「以固定速率
rate往桶里放令牌,最多攒capacity个;请求消耗令牌,没有就等」——它的特点是允许突发:如果一段时间没有请求,令牌会攒起来,之后可以一次性发出最多capacity个请求。漏桶是「请求以固定速率流出」——完全平滑,不允许任何突发(实现上就是记录next_time,每次sleep到那个时刻再+= 间隔)。选择取决于对方的限流实现和业务特点:如果对方也是令牌桶(大多数 API 网关是),你用令牌桶能充分利用突发额度、吞吐更高;如果对方是严格的匀速限流、或者你想保护一个脆弱的下游(比如数据库),漏桶更安全。还有第三种是滑动窗口(记录最近 N 次请求的时间戳),它最贴合「每 60 秒最多 100 次」这类文档描述的语义,代价是 O(N) 内存。实现时的三个共同要点:用time.monotonic()(不受系统改时间影响)、按时间差惰性补充而不是开定时器、用asyncio.Lock保护状态的读-改-写。 - 追问:为什么说「有界队列」提供了背压,这在实践中意味着什么? 背压(backpressure)就是「下游忙不过来时,让上游慢下来」。设想一个流水线:从数据库读 1000 万行 → 调用外部 API 处理 → 写结果。如果中间用无界队列,生产者(读数据库)会以远快于消费者(调 API)的速度往队列里塞,队列无限增长直到 OOM——而且这个 OOM 往往发生在跑了几十分钟之后,非常浪费。换成有界队列(
asyncio.Queue(maxsize=100)),队列满时生产者的await q.put(item)会挂起,直到有 worker 取走一项——于是读数据库的速度自动被 API 的处理速度限制住,内存占用恒定在「队列大小 + worker 数」的量级。这个机制会沿着流水线一级级向上传递:API 慢 → worker 慢 → 队列满 → 读取慢 → 最终反映到数据源。实践中的关键参数是maxsize:太小会让 worker 空转等待(吞吐下降),太大会占内存并延迟背压的生效,经验值是 并发数 × 2~5。
八、加强记忆
asyncio.gather(*一万个任务) 是无界并发,后果是阶跃式的:本地 fd 耗尽(ulimit -n 默认 1024)和临时端口耗尽、一万个协程对象同时驻留内存(一个 aiohttp 请求缓冲约几十 KB,一万并发就是 500MB)、对下游造成 DDoS(被限流、封 IP 或直接打挂)、事件循环堆积回调导致延迟飙升;而重试还会让并发翻倍形成重试风暴。要建立的心智模型是:你要控制的是三个独立的量——并发数、速率、内存占用。Semaphore 只管并发数,两个最常见的错误是:把 async with sem: 包在「创建任务」的地方(只限制了创建速度,任务一创建就跑)和以为它能限 QPS(利特尔法则:并发 = 速率 × 平均延迟,Semaphore(100) 配 10ms 延迟就是 10000 QPS)。限速率要用令牌桶(按时间差惰性补充令牌、用 time.monotonic()、允许突发)或漏桶(完全平滑),「每分钟 N 次」这类规则用滑动窗口最贴合;多副本部署时单进程限流会 N 倍超发,必须用 Redis 做全局限流。超大批量要用「有界队列 + 固定 worker 池」——因为 gather + Semaphore 的所有协程对象仍然同时在内存里,而队列方案内存恒定(只有 N 个 worker 加队列里的 M 项),更重要的是有界队列满时 put() 会阻塞、生产者自动放慢,这就是背压;实现细节:task_done() 必须放在 finally 里(否则 q.join() 永远挂起)、worker 里要吞掉异常、结束时要 cancel worker 防泄漏,maxsize 取并发数的 2~5 倍。每个请求都必须有超时(否则慢请求会占死信号量、有效并发降为 0),且超时值要分层递减;重试要指数退避 + 抖动(没有抖动会让所有失败请求同时重试形成同步尖峰)、走同一个限流器、4xx 不重试、429 按 Retry-After 等待、非幂等写操作要带幂等键。最后,并发数用压测扫参数找「吞吐不再上升但 P99 开始恶化」的拐点,并监控在途请求数和队列长度。