Asyncio并发控制实战:背压、批处理与熔断的稳定性设计
发布时间:2026/9/26 4:57:36 作者:尧图编辑部 阅读量:1,286

去年我重构一个数据同步服务时最开始把 Asyncio 的并发拉满了每来一条消息就create_task最后用asyncio.gather一把梭觉得这样才充分利用了异步性能。结果生产环境上线不到十分钟数据库连接池爆了Redis 内存涨了百分之三十下游接口开始甩 429。我盯着监控面板反复确认了好几次才接受一个事实异步并发不是越快越好背压和批处理才是真正要考虑的稳定性设计。这个教训说来也简单Asyncio 的并发能力强不代表下游有能力承接这么多请求。如果上游把几千个协程同时打到一个只能扛几十 QPS 的服务上结果必然是连接池打满、请求堆积、内存飙升。网上讲 Asyncio 的教程很多都在讲怎么并发更多真正缺的往往是怎么控制并发、怎么应对下游变慢。这篇文章我想把背压backpressure、批处理、熔断这三个关键词串起来结合我实际踩过的坑和可抄的代码讲清楚它们分别解决什么问题以及参数到底怎么调。适合这样几种朋友一是曾经用asyncio.gather一把梭被线上告警教做人的二是想把异步程序里的并发控制做得更规范、不想靠拍脑袋设参数的同学三是想在项目里引入背压和批处理策略但不知道从哪入手的初学者。代码以 Python 3.10 的 asyncio 为主最后一节也会专门讲一讲参数推导的方法。1. 先讲一个把并发拉满然后翻车的真实场景之前那个数据同步服务业务场景不算复杂从消息队列里读用户 ID查 MySQL 里的明细数据再把结果写进 Redis同时同步一份给第三方接口。第一版实现特别自信所有环节都是最高并发策略。数据读取这块我在 one-shot 任务里直接用了asyncio.gather每条消息一个协程查询数据库、写 Redis、调第三方接口全部并行。压测环境的数字确实好看吞吐量比老版本涨了好几倍于是开心地上线了。结果生产流量一进来问题立刻冒头MySQL 连接池被打满大量查询开始排队P99 延迟从 50ms 一路跳到 5 秒多Redis 内存曲线持续上涨检查发现有大量写失败重试任务滞留还有成堆的asyncio.Task对象没有被及时回收每个任务都占着一块协程栈和 Future 状态第三方接口开始返回 429短时间请求频率超过了对方的配额限制一部分同步数据永久丢失。我最开始的第一反应是连接池不够大、Redis 不够快于是把连接池从 50 调到 200给服务加了两台机器。结果问题只是被推迟了并没有真正消失。连接池调大了以后查询确实不怎么排队了但打了更多请求到 Redis 和第三方接口下游的抖动反而更严重。这时候我才意识到问题不在资源配给而在我根本没有给上游设节流阀。真正翻车的地方有两个第一请求从业务里发出之后我没有控并发的手段几千个协程同时在下游排队等 I/O第二重试逻辑是失败就整批重试重启之后积压的消息和重试的消息叠加直接把缓存和下游都拖崩了。这一套操作下来我从并发拉满变成了故障拉满。后面我把代码改成有界队列加固定 worker 的模型配合信号量限制并发再把写 Redis 和调第三方接口的批量收集做起来同样一批流量下P99 延迟降到了 300ms 以内内存曲线变成一条直线。这个反差让我印象极深也是这篇文章想写出来的核心经验。2. 背压的本质不是并发数不够而是下游消化不了2.1 事件循环不会帮你限流很多人对 Asyncio 有一个误解觉得协程是异步的所以并发越多越好反正不阻塞线程。实际上asyncio的核心组件——事件循环——只是一个调度器。它在一个线程里轮流调度所有就绪的协程谁在等 I/O 就把谁挂起谁的 I/O 完成了就唤醒谁。它只负责转起来并不理解下游的能力边界。打个比方事件循环像银行大堂经理协程像在等待区坐着的客户。大堂经理可以让一万个客户同时坐在大厅里但窗口一共就三个柜员一分钟只能处理十笔业务。客户坐多久、大厅会不会被挤爆经理是不管的他只负责排队叫号。下游的真实处理能力才是整个链路的天花板。所以当你用asyncio.gather(fetch_all())一次性拉起几千个协程本质上是让上游的请求瞬间全部涌到下游。数据库有连接池上限Redis 有单实例吞吐上限第三方 API 更是有明确的配额。你发过去的请求越多被拒绝的也越多重试也越多最终形成恶性循环。2.2 任务堆积的代价你以为在等 I/O其实在等内存还有一个很容易被忽略的问题asyncio.Task对象不是免费的。每个没有结束的 Task 都会持有自己的栈帧、局部变量、Future 对象。并发一万个任务就是一万份维持在内存里的运行现场。如果下游响应很慢任务迟迟不结束GC 也回收不了它们。内存涨上去只是第一步更麻烦的是这些任务还在反复超时、反复重试把日志和异常处理器也一起拖垮。我见过一个极端例子某服务下游接口 P99 是 3 秒结果上游一次性投递了两万条消息每条消息都创建了一个协程去请求接口。两万个协程同时在等 I/O内存一路涨到 2GB 之后开始频繁 Full GC最后整个服务被 OOM 杀掉。你说这是并发能力的问题吗明显不是。这是没有给生产端设置节流阀让压力一瞬间砸到了下游。背压要做的事就是把这个一瞬间砸过去改成下游消化多少上游就放多少。上游不再一股脑地生产而是感知到下游队列快满了自己停下来等一下。这就是背压的核心含义让生产者按照消费者的节奏工作。2.3 背压不是性能优化而是稳定性设计想清楚这一点之后我对背压的定位就变了。它不只是调优性能的技巧更是一种稳定性设计在任何环节都可能出现突发流量系统要能扛得住而不是靠下游加班加点。设计顺序通常是先确定下游的吞吐上限再在上游设置并发闸门和有界队列闸门之外的任务要么排队、要么快速失败、要么进入批处理。没有这一步后面所有的高性能异步代码都是空中楼阁。3. 用有界队列和 Semaphore 把失控并发管起来3.1 Semaphore最简单也最有效的并发闸门Asyncio 里限制并发数最直接的工具是asyncio.Semaphore。用法很简单import asyncio sem asyncio.Semaphore(20) async def fetch_one(item): async with sem: # 这里最多只有 20 个协程同时进入 return await call_api(item) async def main(items): tasks [asyncio.create_task(fetch_one(item)) for item in items] results await asyncio.gather(*tasks)这里async with sem做了两件事如果当前活跃的协程数不到 20就直接放行如果已经到 20 了就让新协程在acquire这里等着直到某个协程释放信号量。用这个方式哪怕你create_task创建了五千个任务真正同时在跑的也只有 20 个其余只是排队等待。这个方案比完全无限制已经好很多但它有一个问题任务还是全被创建出来了只是在信号量门口排队。换句话说协程数量本身没有减少只是并发执行数被限制了。如果任务积压太多内存问题依然存在数量级大了以后事件循环本身的调度开销也会上来。3.2 有界队列让生产者也感受到压力要解决生产端失控的问题必须要用到有界队列asyncio.Queue(maxsizeN)。关键行为在put()当队列满的时候put()会阻塞生产者必须等消费者消费掉一部分再继续塞。这样压力就不是只压在下游而是反向传导到了上游。一个经典的生产者-消费者结构是这样的import asyncio from random import randint MAX_QUEUE 100 WORKERS 20 async def producer(queue): for item in range(10000): await queue.put(item) # 队列满时这里会阻塞 await asyncio.sleep(0.01) async def worker(qid, queue): while True: item await queue.get() try: # 模拟下游耗时 await asyncio.sleep(randint(1, 5) * 0.01) finally: queue.task_done() async def main(): queue asyncio.Queue(maxsizeMAX_QUEUE) workers [asyncio.create_task(worker(i, queue)) for i in range(WORKERS)] await producer(queue) await queue.join() for w in workers: w.cancel() asyncio.run(main())MAX_QUEUE100意味着队列中最多攒着 100 个未处理的任务。如果消费者处理不过来生产者会被阻塞在put()这比无界队列安全得多任何时刻积压的数据量上限是队列容量 正在处理的并发数内存占用完全可控。我自己的实践里生产者和消费者之间这个有界队列几乎是所有消息处理类 Asyncio 程序的标准配置。不需要一开始就把并发数调到很大先把队列打满观察消费速度再逐步调整这样整个系统的水位是可见的。3.3 为什么不推荐全量 create_task gather有人会说我不用队列用asyncio.gather配合Semaphore不就行了吗短期看可以长期看有几个隐患create_task一次创建几千个任务事件循环调度开销会上升协程上下文的来回切换不是零成本排查问题时极难定位五千个协程的堆栈一起 dump 出来很难分清哪些是真正在执行哪些是在等信号量一旦某个任务异常导致gather取消所有未完成任务一起取消很可能连带把已经发出、还没返回的请求状态搞乱Semaphore只能限制并发数不能限制任务总量内存水位还是不可控。所以我在正式项目里更推荐有界队列 固定数量 worker的模型worker 数量就是并发度任务在队列里只是普通数据不是调度对象。这个模型从结构上就逼着你想清楚生产速度和消费能力后面调参也有抓手。4. 批处理把 1000 次小请求合并成 10 次大请求4.1 为什么批处理能提吞吐省的是 RTT 和往返开销背压管住了并发但吞吐量还有挖掘空间。很多异步程序性能上不去不是因为下游处理慢而是因为请求次数太多了。每次调用都是一次网络往返RTT即使单条请求只花 1ms来回 1000 次也要 1 秒如果合并成 10 个批次、每批 100 条网络往返就只剩 10 次总耗时可能降到 100ms 上下。这个道理在数据库写入、Redis、消息队列、第三方 API 上都成立。典型场景MySQL/PostgreSQL 批量INSERT而不是一条一条执行Redis 用pipeline或MULTI把几十个命令一次发给服务端日志采集攒一批再 flush而不是每行日志一次 I/O第三方 API 支持数组入参时把单个 ID 查询合并成批量查询。批处理省的不只是 RTT还有序列化/反序列化的固定开销、连接建立和释放的开销甚至下游锁竞争的开销。所以在并发被压住的前提下把小请求攒成大请求往往是吞吐提升最大的一步。4.2 一个可以抄的批量收集器实现Asyncio 里做批处理核心是三个参数最大条数、最大等待时间、消费函数。攒够最大条数就立即发出去如果一直凑不满超过最大等待时间也会发出去保证低流量时数据不会被无限期积压。我写过一个比较通用的批收集器分享出来import asyncio import time class BatchCollector: def __init__(self, flush_func, max_size100, max_wait0.5): self.flush_func flush_func self.max_size max_size self.max_wait max_wait self._queue asyncio.Queue() self._worker None async def start(self): self._worker asyncio.create_task(self._run()) return self async def _run(self): buffer [] last_flush time.monotonic() while True: try: # 最多等 max_wait 秒如果没数据就继续等 item await asyncio.wait_for(self._queue.get(), timeoutself.max_wait) buffer.append(item) should_flush len(buffer) self.max_size except asyncio.TimeoutError: should_flush bool(buffer) if should_flush and buffer: await self.flush_func(buffer) buffer [] async def submit(self, item): await self._queue.put(item) async def close(self): if self._worker: self._worker.cancel()用法async def flush_redis_pipeline(keys): # 一次 pipeline 写入 async with redis_client.pipeline() as pipe: for k in keys: pipe.set(k, k) await pipe.execute() collector await BatchCollector(flush_redis_pipeline, max_size50, max_wait0.2).start() async def on_message(key): await collector.submit(key)这个实现里flush_func就是真正的下游写入函数。你可以根据下游类型自由替换比如批量写数据库、批量发 MQ、批量调第三方接口。核心思想是submit返回特别快真正的 I/O 攒到一批做。4.3 注意批次大小和半满载的平衡批处理不是越大越好。批次太大单次请求耗时会变长下游如果是一个超时时间很短的 HTTP 接口可能直接超时内存里攒的数据也会多一些。批次太小又体现不出合并的优势。我一般这样定先看下游单次能接受的上限比如 Redis pipeline 建议不超过几百个命令数据库批量 insert 一次几百行到一千行再结合单条数据处理耗时把max_size设为单批执行耗时在 100-300ms 之间那个档位。半满载问题靠max_wait兜底max_wait通常设在几十毫秒到一秒之间够攒一批数据又不会让低峰期的数据等太久。这里还有一个容易踩的坑flush_func如果失败了要不要把整批数据放回队列我的建议是不要让收集器背这个锅而是把失败明细通过异常抛给上层由调用方决定是重试、落本地文件还是丢弃。收集器只负责攒批和调用不负责可靠传输否则重试逻辑会把队列搅乱。5. 熔断当下游已经不行了别再发请求了5.1 背压是预防层熔断是保护层背压和批处理解决的是正常波动下上下游节奏不匹配的问题。但还有一种更恶劣的情况下游已经挂了或者网络分区了。这时候你再怎么限流、再怎么批处理每一次请求都是在浪费时间同时还加重下游故障恢复的负担。这就是熔断要解决的问题。熔断器的语义来自电路当错误率达到阈值直接跳闸打开电路后续请求不再发往下游快速失败返回过一段时间进入半开状态放少量试探请求试探成功就关闭熔断、恢复正常失败则继续打开。背压和熔断的区别可以理解为背压是交通拥堵时的缓行熔断是前方道路坍塌后的封路。前者是调节节奏后者是止损。把这两个概念分开系统的故障处理逻辑才清晰。5.2 用 Asyncio 写一个轻量熔断器网上有不少现成的熔断库但如果只是想在一个协程密集的场景里防住下游故障自己写一个几十行的就够了。我的简化版长这样import asyncio import time class CircuitBreaker: def __init__(self, fail_threshold5, reset_timeout10.0): self.fail_threshold fail_threshold self.reset_timeout reset_timeout self.fail_count 0 self.state closed # closed / open / half_open self.opened_at 0.0 def _trip(self): self.state open self.opened_at time.monotonic() async def call(self, func, *args, **kwargs): if self.state open: if time.monotonic() - self.opened_at self.reset_timeout: self.state half_open else: raise RuntimeError(circuit open) try: result await func(*args, **kwargs) if self.state half_open: self.state closed self.fail_count 0 return result except Exception: self.fail_count 1 if self.fail_count self.fail_threshold: self._trip() raise用法breaker CircuitBreaker(fail_threshold5, reset_timeout10.0) async def safe_call(item): return await breaker.call(call_api, item)用连续失败次数做阈值优点是简单直观缺点是如果某一秒突然来 100 个请求失败 5 个和失败 50 个都要等 10 秒才恢复反应不够灵敏。生产环境更讲究一点的话可以用滑动窗口统计最近 N 秒的错误率超过百分比再跳闸。但核心逻辑是一样的。5.3 熔断之后做什么优雅降级比裸报错强熔断只是第一步跳闸之后怎么处理请求同样重要。我常用的降级策略有三种缓存兜底查不到实时数据时返回缓存中的历史值适合监控类、报表类场景跳过非核心任务记录一条日志或者把任务塞回重试队列主流程继续走静默失败对于可丢弃的数据比如埋点日志熔断期间直接丢弃保护系统整体可用性。降级不是随便做要在业务层面讲清楚哪些数据可以丢、哪些必须保。我在数据同步服务里就是核心账目绝不静默丢弃非核心统计可以丢代码里两条路径用不同的开关控制。这样熔断打开的时候系统至少不会崩溃核心链路还能走。6. 参数怎么定从流量模型到实测调优6.1 先算一下你的下游瓶颈别拍脑袋背压、批处理、熔断都写了代码最头疼的问题是参数填多少。我见过很多系统把 Semaphore 并发数设成 1000理由是机器内存够大这基本等于没设。合理的做法是从下游的吞吐上限往上推。假设下游接口的 P99 延迟是 200ms也就是单并发下每秒最多处理 5 个请求。如果你需要每秒处理 200 个请求理论并发数至少是 200 × 0.2 40。再考虑连接池大小比如 30、下游服务实例数比如 3 台我会把并发数定在 min(40, 30 × 0.8) 附近大概 24 到 30。估算公式可以整理成所需并发数 ≈ 目标QPS × 单请求P99耗时 实际设置值 ≈ min(所需并发数, 下游连接池可用连接数 × 0.8)留下 20% 的余量是给突发流量和 GC 停顿留的缓冲。不要正好踩满否则稍微有点抖动就会连锁超时。批处理大小也可以算一算单批请求耗时 网络RTT 下游处理耗时 × 批量大小 / 下游并发能力目标是把单批耗时控制在你的超时阈值的 1/3 到 1/2 之间。比如下游 HTTP 超时是 1 秒那一批请求最好在 300-500ms 内完成。用这个条件反推max_size是多少。6.2 实测调优的步骤从压测到线上观察参数公式只是起点最终要靠实测。我的调优流程一般是这样先用保守参数上线并发数往小里设队列长度往大里设保证功能不受影响写一个压测脚本把生产参数档位的流量打上去观察消费者的实际吞吐和内存曲线逐步调大并发数每调一档跑 5-10 分钟记录 P99 延迟、队列积压曲线、下游错误率找到一个再加大并发吞吐不再明显提升延迟开始上涨的临界点把并发数设在这个点的 80%-90%线上观察时重点看队列剩余容量有没有持续打满。如果持续打满说明消费速度跟不上要扩容而不是继续加并发。这个思路说起来简单执行时很容易被压测工具看到吞吐上去了误导。记住吞吐上去了不代表稳定延迟曲线和错误率才是关键指标。我每次调完参数都会把批大小、并发数、P99 延迟放在同一张图表里看寻找那个收益拐点。6.3 几个容易翻车的细节Semaphore 的acquire没有超时机制如果一直等不到信号量协程会一直挂在那里。生产上建议外层再包一个asyncio.wait_for给排队任务一个超时避免任务无限等待。Redis pipeline 批次过大Redis 单次 pipeline 塞几千个命令服务端可能直接报错或者大幅增加内存开销。批大小要实测不是越大越好。超时和重试要放在批处理外层一条数据写失败了是整个批次重试还是只重试失败条如果是整批重试容易造成重复写入最好是收集失败的下标单独重试那几条。asyncio.TaskGroup和旧式 gather 的取消语义不同Python 3.11 的TaskGroup在某个任务异常时会等待所有任务结束再抛异常和gather的行为不一样。用之前先确认清楚你的取消边界。最后补充一点个人习惯这套背压与批处理策略我在不同项目里反复用过之后最后的体会就一句话先画一条清清楚楚的生产者-消费者边界再谈并发参数。只要边界清晰背压、批处理、熔断其实是水到渠成的事边界模糊的并发优化最后都会变成靠运气上线。现在做新的异步项目时我一般会先在白板上画出数据流上游从哪里来中间队列多大几个 worker 消费下游写入用的是什么方式失败之后走哪条降级路径。画清楚以后再填代码和参数。最近我还会在批处理收集器里加监控埋点队列积压数、批量大小、单批耗时、熔断打开次数全部打进日志和指标系统。这样每次调完参数不是靠感觉而是靠数据说话下一个接手代码的人也能通过监控曲线快速理解这套策略的设计意图少踩几个坑。