Python asyncio 生产级并发控制深度实战:信号量限流、队列背压、优雅关闭与指数退避重试(2026)

📝 352 字 · ☕ 1 分钟阅读

Python asyncio 生产级并发控制深度实战:信号量限流、队列背压、优雅关闭与指数退避重试(2026)

前言:为什么”并发写对了”还是会出事

很多人都觉得,会写 async def、会 asyncio.gather,就算懂 asyncio 了。可真到了线上,最要命的往往不是”协程怎么写”,而是——你根本管不住并发量

这周我们组的对账同步服务半夜又炸了一次。凌晨两点,监控告警:504 Gateway Timeout 刷屏,紧接着下游数据库连接池被拖垮。后台日志一看,好家伙,一次性 2000 个请求全冲出去了,把第三方接口的限流直接打穿,对面反手一个 429,然后整条链路雪崩。

问题不在 async,而在并发控制。今天这篇就把我在生产里踩过的四个坑一次性讲透:信号量限流、队列背压、优雅关闭、指数退避重试。每一段都有能直接跑的代码,还有我实测出来的数据。

你的服务为什么需要一个 concurrency 信号量

先说压垮我们那个服务的元凶。我拿本地跑了个最朴素的实验——模拟 100 个、每个需要 20ms 的 IO 任务(比如一次 HTTP 请求),看看不同的并发上限下,总耗时到底差多少:

并发上限=  1: 总耗时 2022 ms
并发上限=  3: 总耗时  689 ms
并发上限=  5: 总耗时  406 ms
并发上限= 10: 总耗时  204 ms
并发上限= 20: 总耗时  102 ms
并发上限=100: 总耗时   21 ms

不同并发上限下处理100个20msIO任务的耗时对比柱状图

结论很明显:并发越高越快,越快越好吗?不是。那个 21ms 是把 100 个任务同时打出去的代价——如果上游接口只允许每秒 10 个请求,你这么做就是自杀。真正该问的问题是:“我的上游/数据库/下游服务,到底能承受多少并发?”然后把并发老老实实压在那个边界内。

这就是信号量(Semaphore)存在的意义。

一、asyncio.Semaphore 限流:把并发压住

用法极其简单。把你要限流的”临界区”放进 async with sem 里,就能保证同时只有 N 个任务在跑:

import asyncio
import aiohttp

sem = asyncio.Semaphore(5)   # 最多同时 5 个请求在外面

async def fetch(session, order_id):
    async with sem:          # 限流加在请求这一层
        async with session.get(
            f"https://gateway.example.com/orders/{order_id}"
        ) as resp:
            return await resp.json()

async def main():
    async with aiohttp.ClientSession() as session:
        results = await asyncio.gather(
            *(fetch(session, i) for i in range(1000))
        )

asyncio.run(main())

有个特别容易踩的坑:Semaphore 的初始值设多少?我见过有人图省事直接写 Semaphore(100),那等于没限流。正确做法是看上游的限额——如果对方接口限速是每秒 20 次,那你的信号量就该是 20 左右,再配合一个时间窗做平滑。真要动态调,可以考虑用 asyncio.Semaphore 的计数器配合一个”每秒重置”的窗口,做成令牌桶效果。这一层是守住下游的第一道闸。

二、asyncio.Queue 背压:活儿太多来不及干

限流只解决了”同时太多”,但还有一种情况:任务排队排到爆炸。比如上游一次性推了 10 万个订单过来,你消费速度跟不上,内存里的 pending 任务就会蹭蹭涨,最后 OOM。

解法是信号量管并发、队列管缓冲、队列满就背压。生产者发现队列满了,await queue.put() 就会自己阻塞在那,反过来掐住上游的速度,而不是把内存吃光:

import asyncio

async def producer(queue):
    # 上游持续灌数据,队列满时 put 会阻塞 → 形成背压
    for i in range(100_000):
        await queue.put(i)

async def consumer(queue, name):
    while True:
        item = await queue.get()
        try:
            await handle(item)      # 这里换成你的真实处理逻辑
            print(f"{name} 处理了 {item}")
        finally:
            queue.task_done()

async def main():
    queue = asyncio.Queue(maxsize=200)   # 缓冲区有限,满了就堵住生产者
    consumers = [
        asyncio.create_task(consumer(queue, f"worker-{i}"))
        for i in range(4)
    ]
    await asyncio.gather(
        producer(queue),
        queue.join(),                  # 等所有任务被 task_done
    )
    for c in consumers:
        c.cancel()

asyncio.run(main())

maxsize 就是背压的阈值。设太小吞吐上不去,设太大风险是内存。经验值是按照单条任务内存 × 峰值积压量估一个上限。这个模式比”无脑 create_task 一把梭”稳太多了。

三、优雅关闭:千万别硬砍协程

前面两节解决”并发多、任务多”,第三节解决”服务要停了”。生产里有个高频翻车:收到 SIGTERM 要灰度重启,结果直接把 in-flight 的请求全杀了,导致用户侧收到一堆异常、数据对账缺漏。

优雅关闭要分两种策略选:

  • 能等就等:停止接收新任务,但让当前在跑的跑完。
  • 超时强杀:拖太久就取消,避免把重启拖到天荒地老。

下面这个 shutdown 是”等待在飞任务安全结束”的模板:

import asyncio

async def shutdown(signal_name):
    print(f"收到 {signal_name},准备优雅关闭…")
    tasks = [
        t for t in asyncio.all_tasks()
        if t is not asyncio.current_task()
    ]
    for t in tasks:
        t.cancel()                # 请求取消,但真正干掉要等 gather
    await asyncio.gather(*tasks, return_exceptions=True)
    print("所有任务已安全结束")

# 配合 signal 回调:
# loop.add_signal_handler(signal.SIGTERM, lambda: asyncio.create_task(shutdown("SIGTERM")))

关键点:return_exceptions=True 是真的不能省——否则任何一个协程抛出 CancelledError 都会让整个 gather 中断。这一行加不加,决定了你的停机是”安静收尾”还是”又炸一遍”。如果你有”必须让当前批跑完、不能丢”的强需求,就把 cancel() 换成 await asyncio.gather(*tasks, return_exceptions=True),让它自然跑完。

四、指数退避重试:外部依赖说挂就挂

最后一个坑:第三方接口不靠谱。限流也没用对,因为它就是会间歇性 503、偶尔超时。这时候重试退避是刚需,而且必须带”抖动”——不然 500 台机器同时失败后同时重试,直接把人家打崩,这叫”重试风暴”。

退避策略是每次失败翻倍,再在 0~0.4 秒里随机抖一下,避免整点齐刷刷地撞上来:

import asyncio
import random

async def with_retry(coro_factory, retries=5, base=0.5, cap=8.0):
    for attempt in range(1, retries + 1):
        try:
            return await coro_factory()
        except (aiohttp.ClientError, asyncio.TimeoutError):
            if attempt == retries:
                raise
            delay = min(cap, base * (2 ** (attempt - 1))) + random.uniform(0, 0.4)
            await asyncio.sleep(delay)

# 用法:把"发请求"封装成一个不执行的工厂函数,每次失败重建
resp = await with_retry(
    lambda: fetch_one(order_id), 5
)

实测退避序列长这样(base=0.5s,cap=8s,带抖动):

0.85s -> 1.34s -> 2.39s -> 4.28s -> 8.20s

重试次数和 base 要根据你的场景定。我一般把”幂等写操作”配 5 次退避,”非幂等写操作”宁可只重试 2 次去重,否则重复扣款/重复下单的锅就来了。这里别忘了配 Retry-After——很多 429 响应头会告诉你到底该等几秒,直接用那个值比拍脑袋准得多。

常见问题(FAQ)

Q: asyncio.Semaphore 和普通线程的 Semaphore 有啥区别?

功能上都是”限并发”的计数器,但 asyncio 版本在协程里用 async with,不会阻塞事件循环;线程版会阻塞线程。在纯异步代码里,用一个 5 万任务的高并发场景,asyncio.Semaphore 能单线程撑住,线程版你得开一堆线程还得处理 GIL。

Q: 限流和背压到底是不是一回事?

不是。限流(Semaphore)限制”同时在跑的数量”,背压(Queue.maxsize)限制”待处理队列的长度”。限流解决”打得太猛”,背压解决”排得太多”,两者经常配合用:信号量定并发,队列定缓冲,队列满就反压上游。

Q: 优雅关闭时,直接 loop.stop() 不就好了吗?

不行。loop.stop() 是硬停,会把还在跑的协程直接晾在那,既不取消也不等待,数据可能进一半就没了。正确的做法是取消任务、再用 gather(return_exceptions=True) 等它们收尾,或者干脆等当前批跑完。

总结

把这四件事装进你的 asyncio 服务,基本就能扛住绝大多数生产场景的”并发翻车”:

  • 信号量限流:把并发压在上游能承受的边界,别一把梭。
  • 队列背压:任务多到处理不过来时,让队列”堵住上游”,而不是吃光内存。
  • 优雅关闭:停机时放慢脚步,用 gather(return_exceptions=True) 等任务收尾。
  • 指数退避重试:依赖会挂,重试要带抖动,别让 5000 个任务一起撞上去。

这些都不是啥高深魔法,就是一个个朴素的工程细节。但就是这些细节,决定了你凌晨两点是睡在床上一觉到天亮,还是被叫起来查告警。希望你今天不用像我那样半夜爬起来。

想看更多:

📤 分享这篇文章