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(total=10) async with session.get(url, timeout=timeout) 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, f"ERROR: {type(exc).__name__}", 0 async def main(): urls = [ f"https://httpbin.org/delay/{i % 3}" # 模拟不同响应速度的接口 for i in range(50) ] # 关键点 1:Session 只创建一次,所有请求复用同一个连接池 # 关键点 2:limit 控制底层 TCP 连接总数,和信号量是两个维度的限制 connector = aiohttp.TCPConnector(limit=10, ttl_dns_cache=300) semaphore = asyncio.Semaphore(5) async with aiohttp.ClientSession(connector=connector) 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 = [f"https://example.com/page/{i}" for i in range(10)] queue = asyncio.Queue(maxsize=3) results = [] producer_task = asyncio.create_task(producer(queue, urls)) consumer_tasks = [ asyncio.create_task(consumer(f"worker-{i+1}", 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", encoding="utf-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_for,3.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(), timeout=3) 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.timeout,3.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_retries=4): # 信号量放最外面,限制的是“总请求”的并发,包括重试请求 async with semaphore: for attempt in range(max_retries): try: timeout = ClientTimeout(total=5) async with session.get(url, timeout=timeout) as response: # 某些服务器在过载时会返回 429/503,这些也值得重试 if response.status in (429, 500, 502, 503, 504): raise aiohttp.ClientError(f"HTTP {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"第 {attempt+1} 次失败: {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(limit=5)又用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()。症状是:程序刚开始还能输出,跑了一会儿就完全卡住了,Ctrl+C 都难响应。原因仍然是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(limit=20) async with aiohttp.ClientSession(connector=connector) as session: await run_all_tasks(session) # 所有请求都用这个 session还要注意 Session 不是线程安全的,但在同一个事件循环里让不同协程共享 Session 是安全的,因为协程之间是协作式调度,不会同时执行两个请求的读写操作。当然这也是个双刃剑:如果一个请求卡住了,事件循环里其他协程也无法切换走,这又回到超时控制的重要性上——超时、限流、复用连接池,这三件事要一起配齐,异步项目才算有了最基础的“安全护栏”。
我个人在实际操作中的体会是,这四大场景不是孤立的,一个像样的异步项目往往是它们的组合体:用 aiohttp 发请求 + asyncio.Queue 做调度 + to_thread 处理本地文件 + 超时重试限流全程兜底。把这个组合当成一套固定模板反复用,等你对事件循环的调度方式形成了直觉,回头再写同步代码反而会觉得不顺手。如果你正在从同步转向异步,建议先拿这四个场景里的“高并发网络请求”练手,把它跑通、改熟,再逐步加入队列和重试逻辑。