← 返回题目列表

异步并发怎么限流?一次 gather 一万个任务会发生什么?

中等 第 16 / 27 题 更新于 2026/08/01
限流Semaphore令牌桶并发控制

简化版

asyncio.gather(*[fetch(u) for u in urls])同时创建并启动所有任务——一万个 URL 就是一万个协程同时发起连接。后果是四重的:① 瞬间打满本地的文件描述符(每个连接一个 fd,ulimit -n 通常 1024);② 一万个协程对象和缓冲区同时驻留内存③ 对下游发起 DDoS(对方限流、封 IP,或者被你打挂);④ 事件循环里堆积海量回调,延迟飙升。而且 gather先把所有任务创建出来再等,所以「内存峰值」在第一时间就到了。控制并发的标准工具是 asyncio.Semaphoresem = 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 和连接池的限制会叠加aiohttpTCPConnector(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 开始恶化」的拐点,并监控在途请求数和队列长度。