1. Task
调用协程函数只会创建协程对象。直接 await 协程时,当前协程会等它完成;如果希望多个相互独立的协程在同一段时间内推进,就需要把它们交给事件循环调度。
Task 是协程与事件循环之间的执行对象。它主要负责两件事:
- 推进协程运行,直到协程暂停、完成、失败或被取消;
- 保存协程最终的结果、异常和取消状态。
可以用下面的表格区分协程对象和 Task:
| 对象 | 如何得到 | 是否已经被调度 | 主要用途 |
|---|---|---|---|
| 协程对象 | 调用 async def 函数 | 否 | 描述一次异步调用 |
| Task | asyncio.create_task(coroutine) | 是 | 让事件循环独立推进协程 |
01import asyncio0203async def fetch():04print("fetch 开始")05await asyncio.sleep(1)06return "请求结果"0708async def main():09coroutine = fetch()10print(type(coroutine).__name__) # coroutine1112task = asyncio.create_task(coroutine)13print(type(task).__name__) # Task1415result = await task16print(result)1718asyncio.run(main())
Task 不是线程,也不会让 Python 代码自动并行。一个事件循环线程同一时刻仍然只运行一个 Task;当前 Task 执行到尚未完成的 await 后,事件循环才有机会推进其他 Task。
2. 创建
asyncio.create_task() 接收协程对象,把它包装成 Task,并安排到当前正在运行的事件循环:
01import asyncio0203async def fetch(name, delay):04print(f"{name}:开始")05await asyncio.sleep(delay)06print(f"{name}:完成")07return name0809async def main():10first = asyncio.create_task(fetch("任务一", 1), name="first-fetch")11second = asyncio.create_task(fetch("任务二", 2), name="second-fetch")1213print("两个任务已经创建")14first_result = await first15second_result = await second1617print(first_result, second_result)1819asyncio.run(main())
两个 Task 在 main() 等待 first 之前就已经被安排执行。即使代码先 await first、再 await second,second 也会在这段时间内继续由事件循环推进,所以总耗时接近 2 秒,而不是 3 秒。
create_task() 必须在已经运行的事件循环中调用。下面在普通同步代码中直接创建 Task 会失败:
1import asyncio23async def work():4return "完成"56# task = asyncio.create_task(work())7# RuntimeError: no running event loop
普通程序应该在 asyncio.run(main()) 管理的协程中创建 Task;FastAPI 等框架已经运行事件循环,可以在异步接口或异步依赖中创建。
name 参数可以给 Task 设置便于调试的名称:
01import asyncio0203async def work():04await asyncio.sleep(0.1)0506async def main():07task = asyncio.create_task(work(), name="load-user")08print(task.get_name()) # load-user09await task1011asyncio.run(main())
通常可以把 create_task() 理解为「安排任务尽快运行」,而不是保证函数体在这一行立即执行。Task 何时第一次获得执行机会由事件循环调度;代码不应该依赖创建 Task 后、下一行代码前是否已经执行了任务体。
同一个协程对象完成后不能再次 await,但 Task 会保存最终结果,因此完成后的 Task 可以被多次等待:
01import asyncio0203async def work():04await asyncio.sleep(0.1)05return "完成"0607async def main():08task = asyncio.create_task(work())0910print(await task) # 完成11print(await task) # 完成,不会重新执行 work()1213asyncio.run(main())
第二次 await task 直接取得已保存的结果,不会重新执行协程。如果需要把工作再执行一次,应该调用协程函数并创建新的 Task。
3. 等待
如果任务数量固定,并且需要等待它们全部完成,可以使用 asyncio.gather():
01import asyncio0203async def fetch(name, delay):04await asyncio.sleep(delay)05print(f"{name} 完成")06return name0708async def main():09results = await asyncio.gather(10fetch("A", 1),11fetch("B", 0.2),12fetch("C", 0.5),13)14print(results) # ['A', 'B', 'C']1516asyncio.run(main())
传给 gather() 的协程会自动被调度。任务可能按 B、C、A 的顺序完成,但返回结果仍然按照参数 A、B、C 的顺序排列。
默认情况下,某个任务首次抛出异常时,正在等待 gather() 的调用方会立即收到这个异常。需要注意:这不表示其余任务会自动取消,它们通常会继续运行:
01import asyncio0203async def fail():04await asyncio.sleep(0.1)05raise ValueError("任务失败")0607async def slow():08await asyncio.sleep(0.3)09print("慢任务仍然完成了")1011async def main():12slow_task = asyncio.create_task(slow())1314try:15await asyncio.gather(fail(), slow_task)16except ValueError as error:17print(error)1819await slow_task2021asyncio.run(main())
如果设置 return_exceptions=True,异常不会立即抛出,而会和正常结果一起返回:
01import asyncio0203async def success():04return "成功"0506async def fail():07raise ValueError("失败")0809async def main():10results = await asyncio.gather(11success(),12fail(),13return_exceptions=True,14)1516for result in results:17if isinstance(result, Exception):18print(f"异常:{result}")19else:20print(f"结果:{result}")2122asyncio.run(main())
这种方式适合每个任务可以独立成功或失败的批处理。业务逻辑必须主动识别异常对象,否则可能把失败误当成正常结果。
如果需要谁先完成就先处理谁,可以使用 asyncio.as_completed();如果只想等待部分完成或需要更细的等待控制,可以使用 asyncio.wait()。这两个 API 会在后面的「并发」文章中详细讲解。
4. TaskGroup
TaskGroup 从 Python 3.11 开始提供,用于管理一组生命周期相互关联的 Task:
01import asyncio0203async def fetch(name, delay):04await asyncio.sleep(delay)05return name0607async def main():08async with asyncio.TaskGroup() as group:09first = group.create_task(fetch("A", 0.3), name="fetch-a")10second = group.create_task(fetch("B", 0.1), name="fetch-b")1112print(first.result())13print(second.result())1415asyncio.run(main())
进入 async with 后可以创建子任务。退出代码块之前,TaskGroup 会等待所有子任务完成,因此代码块之后读取 result() 是安全的。
这种模式叫结构化并发:创建子任务的代码范围,同时也是负责等待、取消和收集这些任务的范围。任务不会在调用关系中无边界地散落。
TaskGroup 与 gather() 最重要的差别出现在失败时:
| 场景 | gather() 默认行为 | TaskGroup 行为 |
|---|---|---|
| 一个子任务失败 | 向等待者抛出异常,其他任务通常继续 | 取消其余未完成任务并等待它们结束 |
| 多个子任务失败 | 通常先传播首个异常 | 使用 ExceptionGroup 汇总可报告异常 |
| 生命周期 | 调用方自行管理 | 离开代码块前统一收尾 |
01import asyncio0203async def fail():04await asyncio.sleep(0.1)05raise ValueError("请求失败")0607async def slow():08try:09await asyncio.sleep(10)10finally:11print("慢任务执行清理")1213async def main():14try:15async with asyncio.TaskGroup() as group:16group.create_task(fail())17group.create_task(slow())18except* ValueError as errors:19for error in errors.exceptions:20print(error)2122asyncio.run(main())
fail() 失败后,TaskGroup 会取消 slow(),等待它执行 finally 清理,再把异常作为 ExceptionGroup 的一部分抛出。Python 使用 except* 处理异常组中的特定异常类型。
当这些任务共同完成一个业务目标时,通常优先使用 TaskGroup;当每个操作可以独立失败,或者需要按照固定顺序收集所有结果时,gather() 可能更合适。
5. 状态
Task 对外最重要的状态可以简化为:尚未完成、正常完成、异常完成和已取消。
01import asyncio0203async def work():04await asyncio.sleep(0.1)05return "完成"0607async def main():08task = asyncio.create_task(work(), name="example-task")0910print(task.done()) # False11print(task.cancelled()) # False1213result = await task1415print(result) # 完成16print(task.done()) # True17print(task.result()) # 完成18print(task.exception()) # None1920asyncio.run(main())
常见状态方法包括:
| 方法 | 作用 |
|---|---|
done() | Task 是否已经正常结束、异常结束或被取消 |
cancelled() | Task 是否最终被取消 |
result() | 读取结果;任务失败时重新抛出异常 |
exception() | 读取任务异常;正常完成时返回 None |
get_name() | 读取任务名称 |
不要在 Task 尚未完成时调用 result() 或 exception(),否则会抛出 InvalidStateError。业务代码通常直接 await task,状态方法更多用于监控、调试和底层管理。
可以用 asyncio.current_task() 获取当前 Task,用 asyncio.all_tasks() 查看当前事件循环中尚未完成的 Task:
01import asyncio0203async def main():04current = asyncio.current_task()05print(current.get_name())0607for task in asyncio.all_tasks():08print(task.get_name(), task.done())0910asyncio.run(main())
后台任务
不要创建一个 Task 后立刻丢掉引用。事件循环只对 Task 保留弱引用,没有其他引用的 Task 可能在完成前被垃圾回收。可以把后台任务保存在集合中:
01import asyncio0203background_tasks = set()0405async def send_log(message):06await asyncio.sleep(0.1)07print(message)0809async def main():10task = asyncio.create_task(send_log("日志已发送"))11background_tasks.add(task)12task.add_done_callback(background_tasks.discard)1314await task1516asyncio.run(main())
回调会在任务完成后把引用从集合中删除,避免集合不断增长。不过,真正的「创建后不等待」还需要处理任务异常,否则失败时可能出现 Task exception was never retrieved。
对于必须可靠完成、需要重试或不能随进程退出而丢失的后台工作,不要只依赖进程内 Task。FastAPI 进程重启、worker 被终止或部署更新时,Task 都可能中断;这类工作应该使用消息队列和独立 worker。
6. 取消
task.cancel() 发出取消请求,并不表示 Task 在调用这一行立即停止。通常在协程下一次运行到可取消的等待点时,Python 会在协程内部抛出 asyncio.CancelledError:
01import asyncio0203async def long_task():04try:05while True:06print("处理中")07await asyncio.sleep(1)08finally:09print("释放连接")1011async def main():12task = asyncio.create_task(long_task())13await asyncio.sleep(0.1)1415task.cancel()16print(task.done()) # 此刻不一定已经结束1718try:19await task20except asyncio.CancelledError:21print("任务已取消")2223print(task.cancelled()) # True2425asyncio.run(main())
协程应该使用 try...finally 清理网络连接、文件和临时状态。如果确实需要捕获 CancelledError,清理完成后通常应该重新抛出:
1import asyncio23async def work():4try:5await asyncio.sleep(10)6except asyncio.CancelledError:7print("记录取消信息")8raise
CancelledError 直接继承自 BaseException,普通的 except Exception 通常不会捕获它。不要无意中吞掉取消异常,TaskGroup 和 asyncio.timeout() 等结构化并发工具内部都依赖取消机制。
取消是协作式请求,并不保证 Task 最终一定进入取消状态。如果协程捕获 CancelledError 后没有重新抛出,它可能继续运行并正常返回,此时 task.cancelled() 会是 False。除非业务明确要求忽略取消,否则不应该这样处理。
超时
超时通常也是通过取消等待中的 Task 实现的:
01import asyncio0203async def fetch():04await asyncio.sleep(5)05return "完成"0607async def main():08try:09async with asyncio.timeout(1):10print(await fetch())11except TimeoutError:12print("请求超时")1314asyncio.run(main())
asyncio.timeout() 从 Python 3.11 开始提供。超时发生时,代码块外收到 TimeoutError;被等待的协程仍然应该正确响应取消并清理资源。
屏蔽调用方取消
asyncio.shield() 可以防止调用方被取消时,把取消继续传播到指定 Task:
01import asyncio0203async def save_result():04await asyncio.sleep(1)05print("结果已经保存")0607async def main():08task = asyncio.create_task(save_result())09await asyncio.shield(task)1011asyncio.run(main())
shield() 不会让调用方免于收到 CancelledError,也不能阻止 Task 被其他代码直接取消。它只隔断一条取消传播路径,而且仍然必须保留 Task 引用。取消是系统释放资源和停止无用工作的关键机制,因此只有在确实必须让子任务继续时才使用 shield()。