☰
异步函数与异步生成器:从事件循环到流式数据处理的实战指南
2026/9/28 7:32:52 网站建设 项目流程

异步函数与异步生成器,这对概念我琢磨了挺长时间才真正吃透。说实话,网上的教程一抓一大把,但绝大多数都停留在“讲语法”的层面,看完你会用async def和await,但一碰到真实场景,比如高并发请求、流式数据处理,还是不知道怎么组织代码。这篇文章我不打算给你念文档,咱就从一个实际写异步程序的人的角度,聊聊这两个东西到底是什么、怎么用、踩过哪些坑。

这篇内容适合这几类人看:刚学完 Python 基础、想搞懂 asyncio 到底怎么回事的初学者;已经会用 requests 但被性能瓶颈卡住、想切 aiohttp 的爬虫开发者;以及写后端服务时被阻塞调用拖垮、想用异步优化但不知道怎么设计代码的朋友。核心就三件事:异步函数怎么正确声明和调用、异步生成器怎么帮你处理流式数据、以及这两者在真实项目里怎么配合。

1. 先搞清楚什么是异步函数:不是“快”,而是“不等”

1.1 为什么同步代码会卡住

要理解异步函数,得先理解同步代码的痛点。你写一个普通的函数,它会从头到尾顺序执行,碰到耗时的操作(网络请求、文件读写、数据库查询)就傻乎乎地等结果,这个等待的时间里,CPU 什么事也没干。更糟糕的是,如果你在 Web 服务里这么写,一个请求卡在 I/O 上,整个进程的其它请求都得陪着等。

我最早写爬虫的时候就是这种体验:用 requests 循环发请求,100 个 URL 要一个个等响应,少说几十秒。那会儿我以为是网络慢,后来才明白,问题根本不在网速,是我的代码在“傻等”。

1.2 异步函数的本质:协作式调度

异步函数的核心机制是协作式调度,四个字拆开解释:

  • 协作:每个任务自觉地在 I/O 等待点“让出”控制权,而不是被操作系统强行打断(那是抢占式调度)。
  • 调度:一个叫事件循环(Event Loop)的东西负责管理所有任务,谁可以继续执行、谁还得等,都是它说了算。

用生活化的类比来说,同步编程像你在食堂窗口排队打饭,你得等前面的人打完,才轮到你;异步编程像你点完单拿了个号,先去别的窗口吃小吃,等餐好了再回来取。你并没有“同时”做很多事,但你的时间没有浪费在干等上。

注意:异步不是多线程,它是单线程内的并发。异步函数看起来是“同时”跑的,底层其实是事件循环在快速切换,同一时刻只有一段代码在执行。

1.3 async def 和 await 的“搭配规矩”

定义异步函数很简单,把def换成async def就够了。但真正的小坑在于await的使用规则,我见过无数新手在这里翻车。

async def fetch_data(): resp = await http_get("https://example.com/api") data = resp.json() return data

三条红线必须记住:

  • await只能在async def函数内部使用,在普通函数里用await直接语法错误。
  • await的对象必须是“可等待对象”,包括协程(coroutine)、任务(Task)、Future 这三种。
  • 调用一个异步函数,不会执行函数体,只是创建一个协程对象。你必须有await或把它交给事件循环,函数体才会真正跑起来。

这最后一条是新手最容易懵的地方。很多人这么写:

async def main(): fetch_data() # 错!这里只是创建了协程对象,函数根本没执行

正确的写法是await fetch_data()或者asyncio.create_task(fetch_data())。这个区别我后面还会反复提到,因为它是理解整个异步编程模型的钥匙。

2. 从函数到任务:让多个异步函数真正跑起来

2.1 await 是“等”,不是“并行”

很多初学者误以为写了async def就自动并行了,其实大错特错。看这段代码:

import asyncio async def task_a(): print("A start") await asyncio.sleep(2) print("A end") async def task_b(): print("B start") await asyncio.sleep(1) print("B end") async def main(): await task_a() await task_b() asyncio.run(main())

这个程序总耗时是 3 秒,不是 2 秒。因为main里先await task_a(),意思是“等 task_a 执行完再说”,task_a 结束前 task_b 根本没机会跑。注意,await task_a()时控制权确实回到了事件循环,但main自己也在等待,事件循环里除了 task_a 没有别的任务可切换,所以只能干等。

sleep(2)结束之前 task_b 完全没有被调度。这就是传说中的“有 await 但不并发”。

2.2 create_task 和 gather:并发从哪来

真正让多个异步函数并发跑的办法,是让它们成为“任务”(Task),或者用gather、asyncio.wait这类高层 API 把它们分组。

async def main(): # 方式一:手动创建任务 a_task = asyncio.create_task(task_a()) b_task = asyncio.create_task(task_b()) await a_task await b_task # 方式二:gather 一行搞定 await asyncio.gather(task_a(), task_b())

create_task的作用是:把协程包装成 Task,并立刻把它注册到事件循环的调度队列里。这样main进行到await a_task时,虽然它自己不往下走,但事件循环里已经有 task_a 和 task_b 两个任务了,task_a 每次await asyncio.sleep()让出控制权时,事件循环就会切到 task_b 去跑。两个任务的耗时就被重叠起来了。

用gather更省事,它内部帮你处理了任务创建和结果收集,返回值是对应协程结果的列表。

results = await asyncio.gather( fetch_data("url1"), fetch_data("url2"), fetch_data("url3"), ) # results = [data1, data2, data3]

关键认知:并发是任务级别的,不是函数定义级别的。同样是两个异步函数,写await f1()再await f2(),就是串行;用create_task包一层,就是并发。差别就藏在你怎么调用它们。

2.3 超时与取消:异步里的“刹车型号”

并发跑起来了,新问题跟着就来了:如果某个任务卡住了怎么办?同步代码里你最多就是等,但异步场景下,一个任务卡死可能拖垮整个事件循环。

asyncio.wait_for给 await 加了时限:

try: result = await asyncio.wait_for(fetch_data("url"), timeout=3) except asyncio.TimeoutError: # 处理超时 result = fallback_data()

wait_for内部会在超时后取消对应的任务,所以它不仅是“等不到就放弃”,还会主动给任务发取消信号。如果你的异步函数里写了finally块,就有机会做资源清理。

取消机制同样可以直接用:

task = asyncio.create_task(fetch_data("url")) await asyncio.sleep(1) task.cancel() try: await task except asyncio.CancelledError: pass

这里要特别提醒:asyncio.CancelledError在 Python 3.8 之后不再继承自Exception,而是继承自BaseException。所以如果你的代码里有except Exception,是接不住取消信号的。这算是一个隐蔽的坑。

3. 异步生成器:把“懒加载”和“异步”结合起来

3.1 为什么要用异步生成器

先说普通生成器。yield让函数变成迭代器,每次next()才执行到下一个yield,这种“惰性求值”在数据处理时能省内存。

异步生成器就是把yield和await结合起来。函数里既能await异步操作,又能通过yield逐步产出结果,消费者用一个async for就能拿到值。

它解决的核心问题是什么?我举一个特别典型的场景:边下载边处理。

比如你要从一个接口拉取大量数据,一次性全部拉回来内存扛不住,或者响应时间太长体验很差。这时候用异步生成器,每拿到一批数据就yield给上层处理,上层可以边消费边准备下一步请求。用同步生成器做这种事,得先同步地等每一次网络响应,无法并发;用普通异步函数做,又没法分批次产出。异步生成器正好把这两者缝在一起。

3.2 async for 的本质

先看一个最简单的异步生成器:

async def countdown(n): while n > 0: await asyncio.sleep(0.5) yield n n -= 1 async def main(): async for num in countdown(3): print(num)

async for语法上很像for,但它内部做的事情完全不同。普通for循环迭代时调用的是__next__(),同步地取下一个值;async for每次迭代调用的是__anext__(),这是一个异步方法,返回的是一个可等待对象,事件循环会等它完成再继续。

换句话说,async for是对“异步迭代器协议”的语法糖。手动实现它大概是这样的:

gen = countdown(3) try: while True: num = await anext(gen) print(num) except StopAsyncIteration: pass

anext()是 Python 3.10 加入的内置函数,等价于普通生成器里的next()。在 3.10 之前,你得gen.__anext__()然后await它,丑得很。

3.3 anext() 和异步生成器的手动迭代

为什么需要手动迭代?因为有些场景async for不够用。比如你要跳过前几个元素,或者要根据条件提前终止迭代,async for的固定模式做不到。

# 跳过第一个元素 agen = async_source() await anext(agen) # 丢弃第一个 async for item in agen: process(item)

这种需求在实际数据流处理里挺常见的,比如数据源的第一行往往是列名,不是真正的数据。

也要小心的一个细节:异步生成器是冷启动的。创建它时不执行任何代码,必须等第一次await anext()(或async for第一次迭代)才开始执行到第一个yield。这个特性和异步函数一样——“只声明不动手”。

4. 实战:一个带异步生成器的流式下载场景

4.1 场景设计

光讲概念不落地,等于白看。我们设计一个真实场景:有一个分页接口,每页返回 100 条记录,总共有几千页。我们要做的是:

  • 用异步函数并发请求多个页面
  • 用异步生成器把结果逐条产出
  • 上层拿到每一条记录就立即处理(比如入库)

为什么用异步生成器而不用列表收集?因为几千页的数据全收进内存,几十万条记录可能直接把内存打爆。流式处理的需求就是“来一条处理一条”,内存占用恒定。

4.2 完整代码示例

import asyncio import aiohttp async def fetch_page(session, page: int): url = f"https://api.example.com/items?page={page}" async with session.get(url) as resp: resp.raise_for_status() data = await resp.json() return data.get("items", []) async def item_stream(session, max_page: int): """异步生成器:按页拉取数据,逐条产出 item""" for page in range(1, max_page + 1): items = await fetch_page(session, page) if not items: break for item in items: yield item async def main(): async with aiohttp.ClientSession() as session: # 并发拉取,但处理时仍是逐条流式 sem = asyncio.Semaphore(5) # 控制并发数,别把服务器打挂 async def limited_fetch(page): async with sem: return await fetch_page(session, page) # 用 gather 并发拿多页数据,然后逐条处理 pages = await asyncio.gather( *(limited_fetch(p) for p in range(1, 11)) ) count = 0 for page_items in pages: for item in page_items: # 这里做逐条处理,比如写入数据库 count += 1 print(f"processed {count} items") # 如果数据量极大,建议用异步生成器方式 async for item in item_stream(session, 10): # 每条 item 到达就立即处理 pass asyncio.run(main())

4.3 关键细节解读

这个例子看似简单,里面藏了至少四个重要的设计决策:

第一,ClientSession用async with管理。aiohttp 的 session 底层维护了连接池,必须正确关闭,否则会有连接泄漏风险。async with会保证会话结束后连接被释放,这个不能省。

第二,用Semaphore控制并发度。并发不是越大越好。我实测过,并发 50 个请求和并发 5 个请求,总耗时差异不大(瓶颈在目标服务器),但并发过高会导致大量连接被重置,甚至被封 IP。给并发加上限,是对自己和目标服务器双方的尊重。

第三,gather和异步生成器的取舍。gather适合“并发等全部结果再统一处理”的场景,优点是代码简单、好理解;异步生成器适合“边拉边处理”的场景,优点是内存占用低、管道流式传递。两者不是替代关系,是不同场景下的不同工具。

第四,fetch_page里的await resp.json()不能省略。aiohttp 的响应对象拿到手时,body 可能还没下载完。await resp.json()会等 body 全部完成再解析。这一点和 requests 不一样,requests 是阻塞的,响应返回时 body 已经完整了。

5. 常见问题与排查技巧实录

5.1 问题:async 函数里直接写 time.sleep

太经典了,我自己早期也中过招。在异步函数里用time.sleep(1),整个事件循环会被阻塞 1 秒。事件循环是单线程的,你阻塞了它,所有其它任务全部暂停。

# 错:阻塞整个事件循环 async def bad_task(): time.sleep(1) # 事件循环卡死 1 秒 # 对:让出控制权 async def good_task(): await asyncio.sleep(1) # 挂起当前任务,循环可以跑别的

排查方法:程序总耗时和预期不符,比自己算的并发时间慢好几倍,先怀疑是不是有同步阻塞调用混在异步代码里。

5.2 问题:忘了 await,返回的是 coroutine

result = fetch_data("url") # result 是一个 coroutine 对象,不是数据!

这可能是异步编程最常见的新手错误。忘了await,变量拿到的是协程对象,而不是函数返回值。更微妙的是,如果这个协程对象没被任何地方 await 或调度,Python 会抛RuntimeWarning: coroutine was never awaited。如果你用的是 PyCharm 这类 IDE,它会用波浪线提示你,别忽略它。

排查方法:看到coroutine was never awaited警告,逐个检查所有调用异步函数的地方,确认是否都有await。

5.3 问题:async for 里不小心用了同步大循环

异步生成器本身没毛病,但消费者如果写了一个耗时的同步循环,同样会阻塞事件循环。

async for item in item_stream(...): # 这里如果做重 CPU 计算,比如解析超大 JSON、压缩图片 # 事件循环一样被卡住 heavy_compute(item)

遇到这种情况,要么把重计算丢给线程池(asyncio.to_thread),要么就根本不该用 asyncio,直接上多进程。这点我在后面还会展开。

排查方法:事件循环卡顿时,用py-spy dump --pid <PID>看当前线程堆栈,如果栈停在某个同步循环里,基本就实锤了。

5.4 问题:异步生成器不自动执行,必须被消费

异步生成器不同async for或anext驱动,函数体永远不执行。很多人写了个异步生成器,在main里调用了一下,发现什么都没有,怀疑代码有问题。实际上只是没消费。

gen = asyncio.gather(...) # 不跑 agen = item_stream(session, 10) # 不跑 async for x in agen: # 跑起来了 pass

记住:创建异步生成器对象是零成本操作,真正的执行发生在第一次迭代时。这个特性在调试时很头疼,因为代码不执行到那个位置,你根本看不出来有没有 bug。

5.5 排查工具与技巧

调试异步程序比调试同步程序难度高一个量级,因为调用栈不连续,任务切换会把你的思维打断。我的经验是:

第一个基础技巧是在关键位置打印日志,并且带上任务标识。用asyncio.current_task()拿到当前任务对象,打印出任务的名字或 id,这样看日志时能分清哪一段属于哪个任务。

import asyncio async def task_with_log(name): task = asyncio.current_task() print(f"[{name}] task_id={id(task)} start") await asyncio.sleep(1) print(f"[{name}] task_id={id(task)} end")

第二,用asyncio.run(main(), debug=True)开启调试模式。它会提醒你哪些协程没 await、哪些调用耗时超过 100ms(默认阈值)、是否发生事件循环阻塞。对于初学者排查“代码没反应”这类问题非常有帮助。

第三,遇到诡异问题别瞎猜,用asyncio.all_tasks()打印出所有待处理任务,看看是不是有任务卡住了或者泄漏了。

pending = [t for t in asyncio.all_tasks() if not t.done()] for t in pending: print("pending task:", t)

第四个技巧:如果一个任务怎么都找不到退出条件,试试给wait_for包一层超时,让它强制失败。这既是排查手段,也是保险措施。

try: await asyncio.wait_for(suspicious_task(), timeout=3) except asyncio.TimeoutError: print("任务超时,强制退出,问题在任务内部")

6. 一段关于性能边界的经验之谈

6.1 异步不是银弹

我得说点泼冷水的话。异步编程最大的好处是让 I/O 密集型任务的吞吐量飙升,但它不是万能的。

如果任务是 CPU 密集型的,比如加密解密、图像处理、复杂数值计算,asyncio 不仅帮不上忙,反而因为 Python GIL 的限制,多任务之间几乎没法真正并行。你写一堆async def跑密集计算,事件循环照样得一个一个算,白白增加调度开销。

一个简单的判定方法:你的任务里,是“等数据”的时间多,还是“算数据”的时间多?如果 90% 的时间在等网络响应、等数据库返回、等磁盘 I/O,异步是正确选择;如果 90% 的时间在算东西,异步没什么意义,应该用concurrent.futures.ProcessPoolExecutor或直接上多进程。

6.2 什么时候该用多线程/多进程

用 asyncio 的asyncio.run_coroutine_threadsafe和asyncio.to_thread可以在异步程序里桥接多线程:

# 在异步代码里跑一个阻塞函数,丢到线程池去执行 result = await asyncio.to_thread(requests.get, "https://example.com")

to_thread是实现“同步库适配异步框架”的经典技巧。比如你还在用 requests 但想用 async 框架(FastAPI),你不需要把所有代码改成 aiohttp,直接把 requests 调用丢给线程池就行。它会在线程池里跑,阻塞的只是池里的线程,不会卡事件循环。

多进程的话,ProcessPoolExecutor配合asyncio用是有的,但复杂度很高,数据要 pickle 序列化、进程间通信麻烦、调试困难,非必要不推荐。

6.3 我的取舍经验

最后说说我自己做技术选型时的经验。单机搞并发,我会先问三个问题:

  • 是 I/O 密集还是 CPU 密集?前者选异步,后者选多进程。
  • 现有代码库大量使用了同步库吗?如果是,优先考虑to_thread而不是全量重写。
  • 团队熟悉异步吗?异步代码写了一段时间你会发现,它维护起来比同步代码难得多,特别是并发 bug、竞态条件和难看的调用栈。如果团队没有异步经验,可能用线程池 + 同步代码反而更靠谱。

异步使吞吐量上来了,但代码可读性和心智负担也上来了。这个 trade-off,只有上过生产环境的人才有体会。网上那些“一招教你高并发爬虫”的教程从来不提这个,但我见过太多人用 asyncio 把代码写得一团糟,最后性能没上去多少,反而 bug 翻倍。

再补一个我踩过的真实坑:在异步代码里如果用了requests这类同步库,遇到网络超时会非常恶心。requests 默认没有超时,卡住的话整个任务一直挂着,事件循环里其他任务好好的,但这个任务永远结束不了。你在wait_for里设置了超时,它能取消 await,但取消不掉底层线程池里的实际请求——那个线程还在等响应。这是异步 I/O 和同步 I/O 混合使用时最大的暗坑。我的建议是:一旦决定走全异步路线,所有网络调用必须用 aiohttp 或 httpx 的异步模式,不要混用同步库,否则出了问题排查成本极高。

异步函数和异步生成器的核心模型其实不复杂,复杂的是在真实场景里怎么设计任务边界、怎么控制并发、怎么处理取消和超时。这篇文章没有覆盖所有细节,比如asyncio.Queue的生产者消费者模式、async with配合异步上下文管理器实现资源自动回收,这些内容任何一个单独拿出来都能再写一篇长文。

把握住一个主线就够用了:异步让 I/O 等待不再阻塞其他任务,异步生成器让数据流可以按需加工、按需消费,两者叠加,无论处理多大规模的数据,内存和吞吐量都能在可控范围。剩下的细节,在实战里慢慢体会就行。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询