
1. 并发拉满却更慢了一次线上事故让我重新理解 Asyncio先讲个真实经历。之前维护的一个数据同步服务用 asyncio 写了一个定时任务需要从消息队列拉取几万条数据再逐条写入下游存储。第一版代码写得很自信直接把信号量上限设成了 5000心想 asyncio 是单线程协程反正不占系统线程资源并发越高吞吐肯定越大。结果上线第一分钟下游存储的延迟从平均 5ms 涨到了 800ms本地 CPU 倒是没怎么动但任务整体耗时反而比原来 500 并发时慢了三倍。那天下午我盯着监控面板第一次意识到一个反直觉的事实Asyncio 的并发不是越快越好盲目的高并发只是把压力从本地转移到了下游然后下游再把压力反弹回来。这个问题后来花了两周才彻底解决核心就是标题里那两个词背压Backpressure和批处理Batch Processing。这篇就把当时的完整思路和落地方案整理出来包括背压的本质原理、批处理触发器的设计、以及给 Redis 客户端做优雅背压与熔断的实际做法。适合正在用 asyncio 写爬虫、消息消费、数据同步任务的读者也适合那些曾经天真地以为并发数拉满就等于性能拉满的人。1.1 真正的瓶颈不在本地在下游很多人对 asyncio 的最大误解是把并发数直接等同于处理能力。实际上 asyncio 是单线程的它靠事件循环在协程执行 I/O 等待时主动让出控制权从而实现并发。这个机制决定了两个基本事实本地 CPU 计算密集型任务asyncio 没有优势协程切换的开销反而可能让它更慢I/O 密集型任务的真实吞吐上限从来不由本地并发数决定而由下游服务的承接能力决定。拿水管打比方你开 5000 个并发协程相当于在一根只有 100 容量的小水管上硬接 5000 个水龙头。水龙头全拧开水管并不会变大只会全线憋压。asyncio 里表现得更隐蔽——所有协程都在等待 I/O 返回事件循环需要频繁遍历就绪队列、创建 Future、处理回调这些调度开销会随着并发数增长而肉眼可见地上升。我后来用asyncio.get_event_loop().slow_callback_duration做过一次粗略统计并发从 500 升到 5000 时事件循环每次轮询处理回调的耗时增加了近一个数量级。1.2 改进前的代码错在哪当时第一版代码大概是这个风格import asyncio import aiohttp async def fetch_one(session, url): async with session.get(url) as resp: return await resp.json() async def main(): async with aiohttp.ClientSession() as session: tasks [fetch_one(session, fhttp://api.example.com/item/{i}) for i in range(50000)] results await asyncio.gather(*tasks)问题一眼就能看出来一次性把 50000 个协程全部提交没有任何流量控制。asyncio.gather会创建 50000 个 Task每个 Task 都往事件循环里塞然后一起涌向下游。下游一旦慢下来所有协程卡在等待 I/O内存里堆满了未完成的 Future整个进程变成一锅粥。这个错误非常典型我在代码评审里见过无数次把并发控制完全交给下游自己不设任何防线。正确的思路很朴素——本地必须有一个机制让生产者能够感知消费者的承受能力这就是背压的起点。2. 背压的本质不是限速而是让生产者感知消费者状态背压Backpressure这个词最早来自流体力学用在系统设计里指的是当下游处理不过来时把这种处理不过来的信号反向传递给上游让上游主动降低生产速度。很多人一提背压就想到限流其实两者有微妙差别限流通常保护的是本服务自身不被打垮而背压保护的是整条链路——它关心下游是否健康。在 asyncio 生态里做背压没有现成的框架级方案需要自己组合三种工具asyncio.Semaphore信号量、asyncio.Queue有界队列、以及条件变量。它们分别对应三种不同的背压策略。2.1 信号量最简单的闸门Semaphore的用法非常简单它维护一个计数器每次acquire()减一每次release()加一当计数降到 0 时后续acquire()调用会阻塞在协程层面不会占系统线程。import asyncio async def worker(sem, item): async with sem: # 真正执行 I/O 操作 await asyncio.sleep(0.1) return item async def main(): sem asyncio.Semaphore(200) # 最多允许 200 个并发 I/O tasks [asyncio.create_task(worker(sem, i)) for i in range(50000)] await asyncio.gather(*tasks)注意这里的变化虽然 tasks 还是 50000 个但真正同时打到下游的协程最多只有 200 个。剩下的协程停在async with sem这一行等闸门放行。用信号量做背压优点是简单、直观、几乎无脑缺点也很明显——它是一个全局闸门不区分任务优先级不关心下游健康状态的变化更没法知道下游到底还能扛多少。如果你明确知道下游的并发上限比如数据库连接池大小用信号量是最快的解法但在复杂的生产环境里它往往只是兜底的那一层。2.2 有界队列让生产者亲自排队感受压力有界队列的思路是维护一个固定大小的asyncio.Queue生产者往里放任务消费者从里面取任务。当队列满了put()会阻塞生产者自然停下来等待——这就是最直观的背压信号生产者直接被阻塞亲身体会到下游满了。import asyncio import random async def producer(queue: asyncio.Queue, total: int): for i in range(total): await queue.put(i) # 队列满时这里会阻塞生产速度被自动调节 if i % 100 0: print(fproduced {i}, queue size{queue.qsize()}) async def consumer(queue: asyncio.Queue): while True: item await queue.get() try: # 模拟一个耗时随机的 I/O 操作 await asyncio.sleep(random.uniform(0.01, 0.1)) finally: queue.task_done() async def main(): queue asyncio.Queue(maxsize500) # 核心背压参数 consumers [asyncio.create_task(consumer(queue)) for _ in range(20)] await producer(queue, 10000) await queue.join() for c in consumers: c.cancel()这个模式下maxsize就是整条链路的缓冲水位。水位设得越小背压越灵敏吞吐越受限水位设得越大下游抖动时缓冲越充足但代价是任务延迟变高而且队列里的任务可能因为长时间等待而失效比如 token 过期。我个人的经验是maxsize一般取下游允许的并发量乘以 2 到 5 倍。比如 Redis 连接池限制 100 个连接队列水位设在 300~500 比较合适。太低会让生产者频繁阻塞吞吐波动大太高会让背压失效退化成普通缓冲区。2.3 生产-消费者模型里最容易踩的两个坑用有界队列做背压有两个坑几乎人人都会踩到。第一个坑是忘记处理队列消费完毕后的退出。常见写法是消费者用while True死循环任务处理完不知道该什么时候退出最后只能用cancel()硬砍。更好的做法是使用哨兵对象async def consumer(queue): while True: item await queue.get() if item is None: # 哨兵表示没有更多任务 queue.task_done() break # 处理任务 queue.task_done()生产者发完所有数据后往队列里放 N 个NoneN 是消费者数量每个消费者拿到哨兵就退出。第二个坑是把队列当成了无限缓冲。有人觉得队列越大越好顺手写成asyncio.Queue()默认 maxsize 为 0意思是无限。一旦生产速度持续大于消费速度内存会被队列里的待处理对象撑爆。这个坑我见过不止一次生产环境 OOM 之后查半天才发现是无限队列的问题。没有背压的队列本质上是把内存当成了缓冲迟早出事。3. 批处理策略从来一个打一个到攒一批统一处理背压解决的是下游扛不住的问题批处理解决的是单次操作成本太高的问题。两者经常配合使用背压限制下游的瞬间压力批处理降低单位数据的处理成本从而显著提升整体吞吐。3.1 为什么批处理能提升吞吐所有 I/O 操作都有固定成本以 Redis 为例一次网络往返大约 0.1ms~1ms。如果你逐条发送 1000 条命令就需要 1000 次网络往返如果用 Pipeline 批量发送一次往返就能带走几十上百条命令。省掉的是往返时间提升的是单位时间内的命令吞吐量。同理批量写入数据库、批量上报日志、批量调用第三方接口本质都在摊销固定的往返开销。asyncio 虽然让你能同时发起大量异步请求但每次请求依然是一次独立的 I/O 往返。并发再高也改变不了1000 次独立往返的事实而批处理可以把 1000 次往返压缩成 20 次。3.2 最实用的聚合器模式时间窗口 数量窗口批处理最简单的实现是攒够 N 个再一起发但纯数量触发的缺点是如果流量不足任务会一直攒不够延迟越来越大。所以生产环境里更常用的是数量窗口 时间窗口双触发满足任一条件就立即发送。import asyncio import time class BatchCollector: def __init__(self, max_batch_size100, max_wait_time0.1, sinkNone): self.max_batch_size max_batch_size self.max_wait_time max_wait_time self.sink sink # 异步批量发送函数 self.buffer [] self.lock asyncio.Lock() self._flush_task None async def add(self, item): async with self.lock: self.buffer.append(item) if len(self.buffer) self.max_batch_size: batch, self.buffer self.buffer, [] await self._flush(batch) elif self._flush_task is None: self._flush_task asyncio.create_task(self._schedule_flush()) async def _schedule_flush(self): await asyncio.sleep(self.max_wait_time) async with self.lock: if self.buffer: batch, self.buffer self.buffer, [] else: batch [] self._flush_task None if batch: await self._flush(batch) async def _flush(self, batch): # 实际执行批量发送 await self.sink(batch)这个设计有几个细节值得展开add()里用asyncio.Lock保护 buffer避免多协程并发追加时数据错乱数量达到阈值时立即 flush不等待定时器定时器任务只创建一个flush 完成后置回None避免重复创建定时任务定时窗口从第一条数据进来才开始计时不是全局固定周期这样空闲时不会有空转的定时器。3.3 批处理 背压如何组合才算完整如果把批处理直接接到无限队列后面又会出现攒批过程吞掉背压的问题队列里的数据只进不出消费端一直在攒批生产者依然不管不顾地生产。所以完整的链路应该是先背压、后批处理。我最后落地的结构是这样生产者协程 - 有界队列(背压闸门) - 消费者协程 - 批处理聚合器 - 批量 I/O 操作生产者的生产速度受有界队列水位约束消费者从队列取数据后不做逐条 I/O而是丢进聚合器攒批聚合器攒满一批或超时就一次性发往下游。这样既保证下游不会瞬间塞入海量请求也保证了单次操作的摊销成本最小。背压管住量批处理管住效率两者缺一不可。4. Redis 客户端实战优雅背压与熔断机制的落地方案前面讲的是通用策略这一节专门说 Redis 场景。之所以单独拎出来聊是因为 Redis 的高性能和简单协议容易让人放松警惕——很多人觉得Redis 那么快根本不需要什么背压和熔断。直到某次流量高峰Redis 所在宿主机 CPU 被打满Redis 开始响应超时服务端大量重试然后形成一个可怕的循环重试加剧 Redis 压力Redis 压力导致更多超时更多超时触发更多重试。这种连锁故障光靠背压是拦不住的必须上熔断。4.1 给 Redis 客户端加背压的两种方式第一种方式是复用前面的有界队列给 Redis 调用层加一个统一入口import asyncio from redis.asyncio import Redis class RedisGate: def __init__(self, redis: Redis, max_queue_size500, max_concurrent50): self.redis redis self.queue asyncio.Queue(maxsizemax_queue_size) self.sem asyncio.Semaphore(max_concurrent) self._workers [] for _ in range(5): self._workers.append(asyncio.create_task(self._worker())) async def execute(self, command, *args): await self.queue.put((command, args)) # 背压在这里生效 result await self._dispatch(command, args) return result但这个设计有个问题如果要拿到返回值就得让每个调用者等待一个 Future而队列里存放的又是命令Future的元组复杂度直线上升。更实用的做法是给execute包装一层信号量并发数被限制住再用 Redis Pipeline 做批处理。第二种方式是直接吃透redis.asyncio的连接池参数Redis 客户端本身会用连接池维护到服务器的连接。把连接池上限调小等于在客户端侧做了一个隐形的并发闸门但它的粒度是连接而不是命令所以还是要配合 Pipeline 使用。4.2 熔断器三态切换的完整实现熔断器Circuit Breaker是应对下游故障的最后一道防线它有三种状态关闭Closed正常调用记录失败次数失败率达到阈值后切换到打开打开Open直接拒绝请求快速失败不真正访问下游等待冷却时间后进入半开半开Half-Open放少量试探请求如果成功判定下游恢复切回关闭如果失败重新回到打开。import asyncio import time class RedisCircuitBreaker: def __init__(self, failure_threshold10, recovery_timeout10, half_open_max3): self.failure_threshold failure_threshold self.recovery_timeout recovery_timeout self.half_open_max half_open_max self.state CLOSED self.failure_count 0 self.half_open_count 0 self.state_until 0.0 self._lock asyncio.Lock() async def call(self, func, *args, **kwargs): async with self._lock: now time.monotonic() if self.state OPEN: if now self.state_until: raise RuntimeError(circuit breaker open, rejected) self.state HALF_OPEN self.half_open_count 0 if self.state HALF_OPEN: if self.half_open_count self.half_open_max: raise RuntimeError(circuit breaker half-open, rejected) self.half_open_count 1 try: result await func(*args, **kwargs) except Exception: async with self._lock: self.failure_count 1 if self.state HALF_OPEN: self.state OPEN self.state_until time.monotonic() self.recovery_timeout self.failure_count 0 elif self.failure_count self.failure_threshold: self.state OPEN self.state_until time.monotonic() self.recovery_timeout self.failure_count 0 raise else: async with self._lock: self.failure_count 0 if self.state HALF_OPEN: self.state CLOSED self.half_open_count 0 return result这个实现刻意用了一把锁把状态切换串行化防止并发请求同时进入半开导致试探流量翻车。实际使用时把 Redis 查询函数包一层breaker RedisCircuitBreaker() async def safe_get(key): return await breaker.call(redis.get, key)熔断生效后请求会在本地立刻抛出异常不再打到 Redis。这时候配合降级策略——比如返回缓存旧值、返回默认值、或直接把请求丢进一个重试队列——就能把故障影响控制在单次调用范围内。4.3 背压、熔断、批处理三者如何联动熔断和背压看上去都是保护下游但它们分工完全不同背压下游缓慢但没死主动削减并发让下游有时间恢复熔断下游已经明显故障快速失败切断所有流量防止故障扩散批处理下游健康时降低单位请求成本提高吞吐上限。合理的联动顺序是正常状态靠批处理提效波动状态靠背压削峰故障状态靠熔断止损。我在 Redis 客户端上的最终设计就是这三层套在一起外层是熔断器中层是信号量控制的并发闸门内层是 Pipeline 批处理。熔断器判定 Redis 能用了才放行到信号量信号量控制 Redis 上的瞬间命令数批处理把多条命令合成一次 Pipelined 发送。三层各管一档事互相不干扰。5. 实测对比并发数、批次大小与整体吞吐的真实关系说了这么多理论没有实测数据没有说服力。我在调整完成后用本地模拟下游服务做了一组基准测试下游每次请求固定耗时 50ms模拟真实网络延迟单次可批处理上限 100 条。测试过程不复杂但结果很有参考价值列在下面。5.1 并发数对吞吐的影响并发协程数平均请求延迟成功率单任务总耗时10000 条现象5052ms100%52s延迟稳定资源利用率低20055ms100%20.5s延迟略增吞吐稳步上升50062ms100%12.4s接近瓶颈延迟抖动开始出现2000180ms99.2%21s延迟暴涨偶发超时重试5000640ms96.5%45s下游过载成功率明显下降这个表格能说明很多问题并发从 200 提到 500吞吐还在涨但提到 2000 时延迟翻了约 3 倍总耗时反而比 500 节流时更差。下游的承接上限就在 500~800 这个区间超过它多出来的并发全变成了排队时间和失败重试。这里有一个更隐蔽的成本重试会进一步放大下游压力。失败一次客户端重试一次等于把请求数又翻了一倍。在高并发场景下重试风暴是压垮下游的最后一根稻草。所以实测之后我给这个服务定死的并发上限是 400留了 20%~30% 的冗余给下游的短期波动。5.2 批次大小与时间窗口的权衡批处理的两个参数——批次大小和时间窗口——是互相制约的批次太小摊销效果不明显批次太大单次 I/O 的峰值延迟变高下游内存占用也上升时间窗口太长数据在缓冲区内滞留端到端延迟变差时间窗口太短可能攒不满一批就发出去了批处理退化成逐条发送。我的调参思路是先定延迟上限再反推时间窗口。假设业务允许的最大延迟是 200ms那时间窗口就不能超过 100ms留一半给批量发送本身。然后看在这 100ms 内平均能攒到多少条数据再把批次大小设成这个数值的 3~5 倍保证大部分批次是数量触发而不是时间触发。5.3 最终效果调整完成后同一套数据同步任务总耗时从最初的 45s 降到 8.5s下游平均延迟从峰值 800ms 回落到 40ms 以内错误率从 3.5% 降到 0。这个提升不是靠更高并发得来的恰恰相反是靠降低并发放宽了下游的呼吸空间再靠批处理把单位成本打下来。数据自己会说话盲目追求并发数是误区找到那个平衡点才是正解。6. 那些文档里不会写的细节协调锁、取消安全与监控埋点最后聊几个实践细节都是我在踩坑过程中总结出来的常规文档和教程里很少讲。6.1 协调锁的正确打开方式批处理聚合器里的asyncio.Lock以及熔断器里的状态锁作用都是保护共享状态。但别把小锁用成大锁——如果你在持有锁的期间去做真正的 I/O所有协程都会卡在锁上等于自己把自己并发降成了 1。正确的用法是锁只包住状态修改的临界区真正的 I/O 操作移到锁外。这是我反复强调的一点很多人写着写着就把await self.sink(batch)写进了锁里面导致批量性能断崖式下跌。6.2 任务取消要留后路asyncio 的cancel()是协作式取消意味着协程得在某个await点才能感知到取消信号。如果你给消费者协程用了while True循环取消时可能会在任意一个await处抛CancelledError这时候队列里可能还有数据没处理完也可能已经取出来了但还没task_done()。我的建议是两步走一是统一用哨兵机制优雅退出而不是依赖外部cancel()二是如果必须cancel()用asyncio.shield()保护住关键清理逻辑确保队列计数不会错乱。否则下一次执行任务时会报task_done() called once but ...这种让人一头雾水的错。6.3 监控埋点不能省背压系统最怕的是看起来正常实际上所有请求都在队列里排队。我给队列、批次、熔断三个关键点加了计数器埋点队列当前积压量和put()阻塞累计时长批次实际大小分布是否经常低于设定值熔断器每次状态切换的日志时间线。尤其是批次大小分布这个指标作用超出预期。它直接暴露了一个大问题大部分批次根本攒不满是因为生产速率波动太大而不是批次参数设得不对。顺着这个指标我把生产端做了一层平滑批次的实际填充率从 40% 提升到了 85% 以上。6.4 最后分享一个小技巧如果你是第一次给 asyncio 服务加背压别急着写代码先在现有服务里加一个监控指标记录下游处理单个请求的 P99 延迟与请求量的关系曲线。通常你会看到一条明显拐点——拐点右侧就是下游的真正承受上限。背压参数设在哪直接参考这个拐点。我后来给好几个服务调优都是先画这条曲线再定参数比拍脑袋设并发数靠谱一万倍。异步编程的乐趣在于用极少的线程调度海量任务但海量不等于无限。找到系统链条里最弱的那一环用背压保护它用批处理释放它这比把并发数调大几个数量级有用得多。这套方法论不仅适用于 asyncio任何异步框架、任何分布式链路背后都是同一个道理。