异步编程四大实战场景:并发请求、队列、文件IO与超时重试
1. 为什么是这四个场景先搞清楚异步的“收益边界”把 async/await、事件循环、Task 这些基础概念学完之后大部分人都会遇到同一个瓶颈单独看每个知识点都能懂但放到真实项目里完全不知道从哪下手。异步编程不是“把函数前面加个 async 就变快”这么简单它的本质是把等待 IO 的时间省出来给其他任务用。所以选场景的第一步是先认清哪些场景真的适合异步哪些场景硬上异步反而更慢。我这两年做爬虫、数据采集和 API 服务真正高频用到异步的就四个方向高并发网络请求、生产者-消费者任务队列、混合阻塞 IO 的文件处理、超时重试与限流控制。这四个场景基本覆盖了异步编程 80% 的生产力。它们分别对应了四类不同的核心诉求第一类是“把并发量打上去”第二类是“把任务解耦开”第三类是“别让阻塞操作拖垮事件循环”第四类是“让程序在恶劣环境下活着”。为了让你理解得更直观我把这四类场景的核心特点做了个对比场景核心问题关键 API典型收益高并发网络请求大量 HTTP 请求串行太慢aiohttp asyncio.Semaphore gather响应时间缩短 5 到 20 倍生产者-消费者队列生产速度和消费速度不匹配asyncio.Queue task_done/join流量削峰、模块解耦文件与混合阻塞 IO阻塞调用卡住事件循环asyncio.to_thread / run_in_executor避免整体服务假死超时、重试与限流第三方服务不稳定导致雪崩wait_for / timeout / Semaphore 退避重试系统稳定性从“碰运气”变成“可预期”选场景时有一个判断标准只有 IO 密集型任务才适合异步纯 CPU 计算的场景比如跑复杂算法、大量浮点运算用异步不但没有收益还会因为上下文切换引入额外开销。所以下面的四个场景全部围绕 IO 展开——网络等待、队列传递、文件读写、服务响应这些才是异步能发挥真正价值的地方。2. 场景一高并发网络请求异步最典型的“本职战场”2.1 为什么 requests 在这里不够用很多人写爬虫或者调第三方 API第一反应是requests。但 requests 是个同步库发一个请求就得等服务器响应这个等待时间里程序啥也干不了。100 个请求串行跑哪怕每个接口只要 200 毫秒也要 20 秒才能全部完成。问题不在于接口慢而在于你的程序把“等响应”的时间全浪费了。异步的思路是发出去 100 个请求不用等任何一个返回谁先回来就先处理谁。整个过程里你的程序一直在工作只是在不同的协程之间切换。同样是 100 个请求如果限流到 10 并发理论上只需要 2 秒左右收益是数量级的。2.2 完整并发抓取代码与限流设计直接给出一份可用的代码。这个例子做的是并发请求多个 URL每个请求独立计超时通过信号量限制最大并发数避免把对方服务器打爆也避免自己这边瞬间开太多连接。import asyncio import aiohttp from aiohttp import ClientTimeout async def fetch_one(session, url, semaphore): # 信号量控制最大并发数同一时刻最多 n 个请求在跑 async with semaphore: try: timeout ClientTimeout(total10) async with session.get(url, timeouttimeout) as response: body await response.text() return url, response.status, len(body) except asyncio.TimeoutError: return url, TIMEOUT, 0 except aiohttp.ClientError as exc: return url, fERROR: {type(exc).__name__}, 0 async def main(): urls [ fhttps://httpbin.org/delay/{i % 3} # 模拟不同响应速度的接口 for i in range(50) ] # 关键点 1Session 只创建一次所有请求复用同一个连接池 # 关键点 2limit 控制底层 TCP 连接总数和信号量是两个维度的限制 connector aiohttp.TCPConnector(limit10, ttl_dns_cache300) semaphore asyncio.Semaphore(5) async with aiohttp.ClientSession(connectorconnector) as session: tasks [ asyncio.create_task(fetch_one(session, url, semaphore)) for url in urls ] results await asyncio.gather(*tasks) success sum(1 for _, status, _ in results if status 200) print(f完成 {success}/{len(urls)}耗时 {asyncio.get_event_loop().time():.2f}) if __name__ __main__: asyncio.run(main())这段代码里有几个细节值得单独说明。为什么不直接await fetch_one(...)而要先create_task因为await fetch_one()会等这个协程跑完才继续下一个本质上还是串行。create_task是把协程包装成 Task 扔给事件循环去调度然后立刻返回这样才能把所有请求同时“挂”上去最后用gather统一收结果。这个区别是异步并发和伪并发的分水岭。信号量和TCPConnector(limit...)是不是重复了不是。TCPConnector管的是底层连接池最多建多少个 TCP 连接而Semaphore管的是业务层面同时有多少个请求在发。假设连接池上限是 10信号量是 5那么即使连接池有空闲同一时刻也只会有 5 个请求真正发出去剩下 5 个在信号量那儿排队。这在高并发抓取时是必须的双重保护。2.3 实测对比同步和异步差在哪我在本机用 50 个接口测过接口响应时间约 300~600 毫秒不等用 requests 串行跑总耗时约 23 秒用上面的 asyncio 版本信号量限制 5 并发总耗时约 6 秒如果把信号量放开到 10 并发能压到 3.5 秒左右。不用纠结具体数字不同网络环境下差异很大但数量级的差距是稳定的。这里要提醒一句并发不是越高越好。你把信号量调到 100 去请求一个公共接口大概率会被对方限流甚至封 IP。我一般遵循一个经验陌生服务从 5 并发起步验证没问题后再逐步加大对自己部署的服务可以放到 20~50还要看单机带宽和对方负载。异步是把双刃剑能让你快也能让你因为太快要被对方拉黑。3. 场景二生产者-消费者队列用 asyncio.Queue 搭建调度流水线3.1 队列在异步里的角色解耦与背压异步网络请求解决的是“并发怎么打上去”但真实项目里往往不是简单发一批请求就完事。典型情况是一个生产者持续产出任务比如从数据库读出来 10 万条待处理的 URL多个消费者同时处理处理完的结果还要交给后续环节。这时候就需要一个任务队列来做缓冲。asyncio.Queue在设计上和queue.Queue多线程版用法很相似但关键区别是put和get都是协程可以直接被await不会阻塞事件循环。这意味着队列空的时候消费者协程会挂起等待而不是空转轮询队列满的时候设置了maxsize生产者协程会挂起等待而不是疯狂往内存里塞数据。这个机制有一个专门的名字叫背压backpressure——上游速度太快时下游会通过队列反过来限制上游避免内存被打爆。3.2 可落地的爬虫调度器代码下面这个例子模拟的是一个轻量级爬虫调度器生产者从列表里读 URL 放到队列三个消费者并发消费各自去抓取和处理处理完成后输出结果。import asyncio import random async def producer(queue, urls): for url in urls: # 队列满时这里会自动挂起等消费者消费掉部分任务再继续放 await queue.put(url) print(f[生产者] 放入 {url}) # 生产结束放入 None 作为“结束哨兵”告诉消费者可以收工了 await queue.put(None) print([生产者] 生产结束) async def consumer(name, queue, results): while True: url await queue.get() # 遇到哨兵值先放回去因为其他消费者也要看到然后退出 if url is None: queue.task_done() break # 模拟抓取和处理的耗时有快有慢才贴近真实 delay random.uniform(0.2, 0.8) await asyncio.sleep(delay) result f{name} 处理了 {url} 耗时 {delay:.2f}s results.append(result) print(result) # 标记这个任务已消费完成配合 join 使用 queue.task_done() async def main(): urls [fhttps://example.com/page/{i} for i in range(10)] queue asyncio.Queue(maxsize3) results [] producer_task asyncio.create_task(producer(queue, urls)) consumer_tasks [ asyncio.create_task(consumer(fworker-{i1}, queue, results)) for i in range(3) ] # 等待生产者结束再等待队列中的所有任务消费完 await producer_task await queue.join() print(f\n全部完成共 {len(results)} 条结果) # 消费者会在遇到哨兵后自行退出这里确保所有消费者协程收尾 await asyncio.gather(*consumer_tasks) if __name__ __main__: asyncio.run(main())3.3 消费者数量与结束信号的设计细节这个代码看起来简单但里面有两个特别容易踩坑的细节。第一个是哨兵值的处理。三个消费者都在queue.get()如果生产者只放一个None那么只有抢到None的那个消费者会退出另外两个会永远阻塞着等新任务。我的处理方式是消费者碰到None时先task_done()然后break不把 None 放回队列所以实际上三个消费者里只有一个能看到哨兵。这里更稳的做法是生产者放consumer_count个哨兵或者用asyncio.Event做统一结束信号。上面的代码为了简洁用了“碰到的那个消费者退出”的策略配合gather确实能结束但如果消费者和任务数量不匹配你可能要仔细推敲一下结束语义。第二个是queue.join()的语义。join()会一直阻塞到队列里所有任务都被task_done()标记过。也就是说只有 put 进队列还不够每个 get 到的任务都必须显式调用task_done()否则join()永远不返回。这是 asyncio.Queue 使用里最常见的问题忘了调task_done()程序就卡在await queue.join()那里一动不动排查起来非常头疼。队列的maxsize建议按“单个任务的内存占用 × 任务数”来估算。举个例子如果你的任务对象包含一个 10MB 的响应体队列最多 100 个任务就是 1GB 内存这显然太大。更合理的做法是队列里只放任务的“描述信息”比如 URL、ID真正的响应数据放在结果里避免队列成为内存黑洞。4. 场景三文件与混合阻塞 IO别让一件“小事”卡死整个事件循环4.1 为什么普通文件读取会阻塞所有协程很多人学异步时会忽略文件操作觉得“读个文件能有多慢”。但真相是文件读写属于阻塞 IO一旦在协程里直接用open()、read()、write()这些同步方法当前线程的事件循环会被整个卡住。假设你的事件循环里同时挂着 50 个网络请求协程其中某个协程执行了一次file.read()那么这 50 个协程全部要等这个文件读完才能继续调度。机械硬盘随机读一次几十毫秒网络文件比如 NFS、对象存储挂载甚至能到几百毫秒这在异步系统里是灾难级的表现。要理解这一点得回到事件循环的工作方式协程之间的调度是协作式的一个协程只有在遇到await挂起点时才会把控制权还给事件循环。如果某个协程一直不await事件循环就拿不回控制权。同步文件读取恰恰就是“一直不 await”的操作——它得等数据从磁盘上完全读回来才返回。4.2 三种绕开阻塞的正确姿势官方推荐的姿势是asyncio.to_thread()它是 Python 3.9 加入的语法糖底层用默认线程池包装一个普通函数让它在线程里执行从而不阻塞事件循环import asyncio async def read_file_safe(path): # to_thread 把同步文件读取丢到线程池执行事件循环不会被卡住 content await asyncio.to_thread(open, path).__enter__() # 注意to_thread 适合短操作如果要完整读文件直接这样写更好 # data await asyncio.to_thread(_read_file, path) return content def _read_file(path): with open(path, r, encodingutf-8) as f: return f.read() async def main(): data await asyncio.to_thread(_read_file, big_file.txt) print(len(data)) asyncio.run(main())如果你同时要处理很多个文件可以配合信号量限制线程池里的并发数量避免几百个文件同时读把磁盘 IO 撑爆。更偏底层的做法是用loop.run_in_executor()它在 Python 3.9 前是主力方案现在to_thread更简洁二者本质一样。还有一个第三方库aiofiles它把异步文件读写封装成了类似open()的接口。但我个人建议能用to_thread就别上 aiofiles。因为 aiofiles 内部也是线程池方案多一层封装反而让你对底层行为不可控而且它的性能并没有比原生to_thread好。除非你的代码里大量使用async with aiofiles.open(...)这样的写法能显著提升可读性否则标准库就够用了。4.3 大文件流式读取的实战写法遇到大文件时一次性read()会把整个文件加载到内存这在大规模文件处理场景下是不现实的。正确做法是分段读取每次只读一小块处理完再读下一块。结合to_thread可以写成这样import asyncio CHUNK_SIZE 1024 * 1024 # 1MB 一块 def read_chunk_sync(file_obj, size): return file_obj.read(size) async def process_big_file(path): loop asyncio.get_running_loop() results [] # 用同步方式打开文件获取文件对象打开操作很快不值得丢线程池 with open(path, rb) as f: while True: chunk await asyncio.to_thread(read_chunk_sync, f, CHUNK_SIZE) if not chunk: break # 模拟处理统计块长度 results.append(len(chunk)) if sum(results) % (10 * CHUNK_SIZE) 0: print(f已处理 {sum(results) / 1024 / 1024:.0f} MB) return sum(results) async def main(): total await process_big_file(large_data.bin) print(f总字节数: {total}) asyncio.run(main())这里有个细节文件对象f的read()方法被传进了线程池而文件对象本身的打开/关闭操作仍然是同步的。为什么打开和关闭不丢线程池因为open()和close()本身不涉及大量数据搬运耗时通常在微秒级偶尔调用一次不会对事件循环造成可感知的影响。真正需要丢线程池的是read()这种可能等磁盘响应的操作。在真实项目里我还遇到过另一个坑同时读多个大文件时每个文件都往默认线程池丢任务默认线程池只有 min(32, os.cpu_count() 4) 个线程并发一多线程池排队反而比串行更慢。我的做法是给大文件操作单独建一个线程池然后用run_in_executor(custom_pool, ...)这样能精确控制同时进行多少个文件 IO。这个细节很少被人提到但实际项目里非常有用。5. 场景四超时、重试与限流把异步服务放进“安全护栏”5.1 超时控制wait_for 与 asyncio.timeout 的选择真实生产环境里第三方接口不会总是如你所愿地响应。如果一个接口 30 秒都没返回你的协程就会一直挂着。更糟的是如果几十个协程都挂在同一个慢接口上整个系统的并发能力直接被拖死。所以超时是异步系统的基础设施不是可选项。Python 3.11 之前主要用asyncio.wait_for3.11 之后新增了asyncio.timeout上下文管理器。它们的核心区别在于wait_for是“包住一个 awaitable”超时后直接取消任务timeout是“在代码块内生效”超时后抛异常但代码块里的资源有机会被清理。举个直观的例子import asyncio async def slow_operation(): await asyncio.sleep(10) return done async def demo_wait_for(): try: result await asyncio.wait_for(slow_operation(), timeout3) return result except asyncio.TimeoutError: return timeout after 3s async def demo_timeout_context(): try: async with asyncio.timeout(3): await slow_operation() except TimeoutError: return timeout after 3s两者的另一个差异在嵌套场景。asyncio.timeout的上下文可以嵌套内层超时优先触发比wait_for更灵活也更安全。我的建议是Python 3.11 的项目优先用asyncio.timeout3.10 及以下的老项目才用wait_for。注意asyncio.timeout()的括号里参数是秒数传None表示不限制——这个特性也意味着你可以在运行时动态调整超时策略而不是写死一次。5.2 重试与指数退避的正确姿势超时之后直接放弃显然太草率。网络抖动、临时过载这些情况重试一次往往就成功了。但重试不能是无脑重试需要配合指数退避第一次失败等 1 秒、第二次等 2 秒、第三次等 4 秒……这样既给了服务恢复的时间也不会在大规模故障时加剧对方压力。import asyncio import aiohttp from aiohttp import ClientTimeout async def request_with_retry(session, url, semaphore, max_retries4): # 信号量放最外面限制的是“总请求”的并发包括重试请求 async with semaphore: for attempt in range(max_retries): try: timeout ClientTimeout(total5) async with session.get(url, timeouttimeout) as response: # 某些服务器在过载时会返回 429/503这些也值得重试 if response.status in (429, 500, 502, 503, 504): raise aiohttp.ClientError(fHTTP {response.status}) return response.status, await response.text() except (asyncio.TimeoutError, aiohttp.ClientError) as exc: if attempt max_retries - 1: # 最后一次失败就直接抛出让上层决定怎么处理 raise wait_time 2 ** attempt # 1, 2, 4, 8... print(f第 {attempt1} 次失败: {exc}等待 {wait_time}s) # 用 asyncio.sleep而不是 time.sleep这是重试里最容易犯的错 await asyncio.sleep(wait_time) return None这里最关键的忠告重试里的等待必须用asyncio.sleep()绝对不能用time.sleep()。time.sleep()是同步阻塞的它会冻结整个事件循环等于把异步重试的优势全毁了。5.3 限流纪律信号量放在“入口”还是“出口”限流看似简单但信号量的位置非常讲究。我见过不少人犯同一个错误把Semaphore放在循环外面导致所有 Task 都在创建时一把梭地冲进来信号量根本拦不住——因为信号量只对进入async with semaphore:且被 await 的协程生效。正确的位置是信号量位于每个任务内部的函数体最外层也就是每个协程真正开始工作之前。但要注意它和“连接池”的层级关系。如果你用TCPConnector(limit5)又用Semaphore(10)那么底层最多 5 个连接信号量最多 10 个并发最后实际并发是 5因为连接池先卡死了。所以两个数值要配合好一般让信号量略小于连接池限制这样是信号量在起作用而不是连接池在兜底。更严谨的限流做法是考虑“令牌桶”模型不分单次请求而是按时间窗口限制总量。比如“每秒最多 20 个请求”。asyncio.Semaphore做不到时间窗口控制需要配合asyncio.Sleeping或者自己实现一个定时刷新令牌的协程。如果追求简单可以用信号量 sleep 的组合效果也够用。6. 这些坑我踩过希望你别再踩Task 生命周期与阻塞陷阱6.1 Task 引用被回收任务“神秘消失”Python 的asyncio.create_task()有一个非常隐蔽的行为如果 Task 对象的引用没有被保存它可能会被垃圾回收任务随之消失。这不是危言耸听至少我见过不止一个人在循环里这么写# 反例Task 没有被引用保存可能直接被 GC 掉 async def bad_example(): for url in urls: asyncio.create_task(fetch(url))这段代码看起来没问题任务也创建了但因为没有变量保存 task 引用Task 可能在执行到一半时被 CPython 的 GC 回收任务就神秘消失了连异常都不会抛。正确做法是用列表保存所有 Task 引用# 正例保存引用等待全部完成 tasks [] for url in urls: tasks.append(asyncio.create_task(fetch(url))) await asyncio.gather(*tasks)还有一种更隐蔽的情况用asyncio.gather其实也会内部保存 Task 引用所以gather是安全的。但如果你只是为了“创建任务然后后边再说”一定要显式持有引用。这个坑的排查难度极高因为程序看起来能跑只是结果少了几个特别容易被怀疑成请求失败而不是 Task 被回收。6.2 协程里的 time.sleep 与“假死”另一种容易把新手坑到怀疑人生的写法是在协程里用了time.sleep()。症状是程序刚开始还能输出跑了一会儿就完全卡住了CtrlC 都难响应。原因仍然是time.sleep()阻塞了事件循环线程。协程看起来是并行执行的实际是事件循环在单线程里调度一个协程睡 5 秒所有其他协程都得陪着睡 5 秒。处理这类问题时我会先全局搜索time.sleep确保全部替换成await asyncio.sleep。这是异步代码审查的第一个检查项。6.3 复用 Session 与连接池的正确姿势最后一个高频坑是滥用aiohttp.ClientSession。每次请求都 new 一个 Session然后在请求结束后关闭这会让 TCP 连接频繁建立和销毁。TCP 握手和 TLS 协商的开销非常大异步的优势被吃掉大半。正确做法是整个程序生命周期只创建一次 Session所有请求复用同一个 Session 和它背后的连接池。async def main(): # Session 在 async with 外创建程序结束时才关闭 connector aiohttp.TCPConnector(limit20) async with aiohttp.ClientSession(connectorconnector) as session: await run_all_tasks(session) # 所有请求都用这个 session还要注意 Session 不是线程安全的但在同一个事件循环里让不同协程共享 Session 是安全的因为协程之间是协作式调度不会同时执行两个请求的读写操作。当然这也是个双刃剑如果一个请求卡住了事件循环里其他协程也无法切换走这又回到超时控制的重要性上——超时、限流、复用连接池这三件事要一起配齐异步项目才算有了最基础的“安全护栏”。我个人在实际操作中的体会是这四大场景不是孤立的一个像样的异步项目往往是它们的组合体用 aiohttp 发请求 asyncio.Queue 做调度 to_thread 处理本地文件 超时重试限流全程兜底。把这个组合当成一套固定模板反复用等你对事件循环的调度方式形成了直觉回头再写同步代码反而会觉得不顺手。如果你正在从同步转向异步建议先拿这四个场景里的“高并发网络请求”练手把它跑通、改熟再逐步加入队列和重试逻辑。