agent异步调用学习--asyncio 模块
·
为什么需要 asyncio?
在传统的同步编程中,当一个任务需要等待 I/O 操作(如网络请求)完成时,程序会阻塞,直到操作完成。这会导致程序的效率低下,尤其是在需要处理大量 I/O 操作时。
asyncio 通过引入异步编程模型,允许程序在等待 I/O 操作时继续执行其他任务,从而提高了程序的并发性和效率。
想象一下你正在经营一家餐厅:
- 同步模式(普通函数): 你只有一个厨师。客人 A 点了一份牛排,厨师开始煎牛排(这需要等待 5 分钟)。在煎牛排的这 5 分钟里,厨师完全被占用,不能做任何其他事,即使客人 B 只想点一杯水,也必须干等着。
- 异步模式(asyncio): 你有多个厨师(实际上还是一个,但非常聪明)。厨师开始煎客人 A 的牛排后,发现需要等待,他立刻把这份牛排标记为等待中,然后转头去给客人 B 倒水。倒完水回来,看看牛排是不是快好了,如果还没好,又可以去处理客人 C 的订单。这样,在等待 I/O(如煎牛排、网络请求、读写文件)的时间里,厨师(CPU)一直在高效地工作。
asyncio 就是 Python 用来实现这种聪明工作模式的标准库,它允许你编写 单线程并发 的代码,特别适用于网络爬虫、Web 服务器、微服务等 I/O 密集型场景。而且相对应的在当下agent需要频繁调用相关IO操作则必须学会掌握它来提高系统运行的效率。
asyncio 的核心概念
- 协程(Coroutine)
协程是 asyncio 的核心概念之一。它是一个特殊的函数,可以在执行过程中暂停,并在稍后恢复执行。协程通过 async def 关键字定义,并通过 await 关键字暂停执行,等待异步操作完成。
如何让协程真正执行?
- 方式一(推荐):使用 asyncio.run(),它是 asyncio 提供的顶层入口,负责创建事件循环、运行协程并清理资源。
asyncio.run(say_hello()) # 直接运行协程
- 方式二:在另一个异步函数中 await 它,然后通过 asyncio.run() 启动这个顶层异步函数。
await say_hello() # 在异步上下文中等待协程完成
if __name__ == '__main__':
asyncio.run(main())```
3. 方式三:使用 asyncio.create_task() 将协程包装为任务并调度执行(但通常仍需要 await 任务或保持事件循环运行)。
```python
import asyncio
async def say_hello() :
print("hello") await asyncio.sleep(1) print("world")
# 由于在 Jupyter Notebook 中运行,不需要在主程序中运行协程,notebook已经是一个事件循环不需要使用方式一来运行,py不允许嵌套事件循环会报错误
await say_hello()
hello world
- 事件循环(Event Loop)
事件循环是 asyncio 的核心组件,负责调度和执行协程。它不断地检查是否有任务需要执行,并在任务完成后调用相应的回调函数。
实例
await say_hello()
asyncio.run(main())
- 任务(Task)
任务是对协程的封装,表示一个正在执行或将要执行的协程。你可以通过 asyncio.create_task() 函数创建任务,并将其添加到事件循环中。(Task可以类比为真正执行的人,他既具备future也接管了协程的执行权)
实例
task = asyncio.create_task(say_hello()) await task```
4. Future(Future可以类比为一个承诺,例如外卖订单完成后,会通知你)而协程这是一份提前写好的购物清单,
Future 是一个表示异步操作结果的对象。它通常用于底层 API,表示一个尚未完成的操作。你可以通过 await 关键字等待 Future 完成。
实例
```pythonasync def main():
future = asyncio.Future() await future```
### 概念深度解析
asyncio.Future(底层基础件)
- 是什么:它是 asyncio 中最底层的可等待对象(Awaitable)。它是一个容器,专门用来存放未来某个时刻才会产生的值。
- 核心机制:它内部有一个状态机(Pending -> Finished / Cancelled)。它不关心这个值是怎么来的(可能是网络请求返回的,可能是硬盘读取的,也可能是另一个线程传过来的)。
- 关键特征:Future 不会主动执行任何代码。它只是被动地等待别人调用 set_result() 或 set_exception() 来给它填值。
- 谁在用:通常我们不直接在业务代码中创建 Future,而是由底层库(如 asyncio.open_connection 或 loop.run_in_executor)返回。
asyncio.Task(协程执行器)
- 是什么:Task 是 Future 的子类。它把 Future 和协程函数结合在了一起。
- 核心机制:当你用 asyncio.create_task(some_coro()) 时,Task 会主动将你的协程交给事件循环(Event Loop),并立即开始调度执行。
- 关键特征:Task 拥有一个 coro 属性,它通过 send() 方法驱动协程一步步运行,直到遇到 await 暂停,或者运行结束。当协程返回结果时,Task 会自动调用底层的 set_result(),将自己这个 Future 标记为完成。
核心区别(面试/开发必知)
维度 Future Task
身份关系 父类(基类) 子类(Future 的子类)
执行主动性 被动(惰性)。只存储结果,不执行代码。 主动。立即把协程丢给事件循环去跑,是“驱动者”。
创建者 通常由底层 loop 或库创建(如 loop.run_in_executor)。 由开发者显式创建(asyncio.create_task() 或 loop.create_task())。
典型应用 连接底层 IO(网络、文件),回调转异步。 并发执行多个后台协程(如同时爬 10 个网页)。
是否包含协程 不包含,只是一个空盒子。 包含(内部持有 _coro)。
### 同步版本的示例
```python
import asyncio
import requests
import time
def fetch(url):
# 模拟耗时的网络请求
print(f"start 获取 {url}") time.sleep(2) print(f"end 获取 {url}")
def main_sync():
urls=["https://www.baidu.com","https://www.taobao.com","https://www.jd.com"] results=[] start =time.time() for url in urls: results.append(fetch(url)) end =time.time() print(f"耗时 {end-start} 秒")
print(results)main_sync()
start 获取 https://www.baidu.com end 获取 https://www.baidu.com start 获取 https://www.taobao.com end 获取 https://www.taobao.com start 获取 https://www.jd.com end 获取 https://www.jd.com 耗时 6.004304885864258 秒
[None, None, None]
异步版本的示例
import asyncio
import time
import aiohttp
async def fetch_async(session,url):
print(f"异步start 获取 {url}") async with session.get(url) as response: await asyncio.sleep(2) text=await response.text() print(f"异步end 获取 {url}") return f"{url} 获取到的内容是:{text}"
async def main_async():
urls=["https://www.baidu.com","https://www.taobao.com","https://www.jd.com"] results=[] async with aiohttp.ClientSession() as session: # 创建一个异步会话对象
tasks=[] for url in urls: task=asyncio.create_task(fetch_async(session,url)) tasks.append(task)
print("所有任务已经创建")
results=await asyncio.gather(*tasks) # *task 将列表进行结包,然后返回多个独立的参数传给gather
return resultsif __name__ == '__main__':
start =time.time()
final_results=await main_async() end=time.time() print(f"耗时 {end-start} 秒")
所有任务已经创建
异步start 获取 https://www.baidu.com 异步start 获取 https://www.taobao.com 异步start 获取 https://www.jd.com 异步end 获取 https://www.taobao.com 异步end 获取 https://www.jd.com 异步end 获取 https://www.baidu.com 耗时 2.1354458332061768 秒
| 函数 | 主要作用 | 常用参数说明 |
|---|---|---|
asyncio.run(coro, *, debug=False) | 运行一个顶层协程,管理事件循环的生命周期。是异步程序的主入口。 | coro: 要运行的协程对象。debug: 设为 True 可启用事件循环的调试模式(如慢回调检测)。 |
asyncio.create_task(coro, *, name=None) | 将协程包装成一个 Task 对象,并排入事件循环等待调度。这是实现并发的主要方式。 | coro: 要包装的协程对象。name:(Python 3.8+)为任务指定一个名称,便于调试和日志追踪。 |
asyncio.gather(*aws, return_exceptions=False) | 并发运行多个异步任务(可接受协程、任务等),等待所有完成,返回结果列表(顺序与传入顺序一致)。 | *aws: 可变参数,传入多个异步对象。return_exceptions: 默认为 False,任何异常会立即传播;设为 True 时,异常会作为结果元素返回而不中断整体。 |
asyncio.sleep(delay, result=None) | 异步地休眠指定秒数。非阻塞(与阻塞线程的 time.sleep 有本质区别),常用于模拟 IO 等待或让出控制权。 | delay: 休眠的秒数(可以是浮点数)。result: 休眠结束后返回给 await 调用者的值。 |
asyncio.wait(aws, *, timeout=None, return_when=ALL_COMPLETED) | 并发运行任务,并等待满足指定条件。返回一个元组 (done, pending),分别是已完成和未完成的任务集合。 | aws: 可迭代的异步对象集合。timeout: 超时时间(秒),超时后未完成的任务会留在 pending 中。return_when: 触发返回的条件,可选 FIRST_COMPLETED、FIRST_EXCEPTION、ALL_COMPLETED(默认)。 |
asyncio.to_thread(func, /, *args, **kwargs) | (Python 3.9+)将普通的、可能阻塞的同步函数放到一个独立的线程中运行,返回一个可 await 的协程。用于处理 CPU 密集型或遗留的阻塞式 IO 操作。 | func: 要在线程中运行的同步函数。*args, **kwargs: 传递给该函数的位置参数和关键字参数。 |
asyncio的基本用法
- 运行协程–要运行协程,你可以使用asyncio.run()函数,他会创建一个事件循环,并且运行指定的协程
# 示例
import asyncio
async def main():
print("start") await asyncio.sleep(2) print("end")asyncio.run(main())
- 并发执行多个任务–你可以使用asyncio.gather()函数来并发执行多个任务,等待所有任务完成,返回结果列表。
import asyncio
async def task1():
print("start task1") await asyncio.sleep(2) print("end task1") return "task1 result"async def task2():
print("start task2") await asyncio.sleep(2) print("end task2") return "task2 result"async def task3():
print("start task3") await asyncio.sleep(2) print("end task3") return "task3 result"async def main():
await asyncio.gather(task1(),task2(),task3())asyncio.run(main())
start task1 start task2 start task3 end task1 end task2 end task3
- 超时控制–你可以使用asyncio.wait_for()函数来设置任务的超时时间,超过时间未完成的任务会被取消。如果协程在指定时间内未完成,将引发 asyncio.TimeoutError 异常。
import asyncio
async def long_task():
await asyncio.sleep(10) print("Task finished")
async def main():
try: await asyncio.wait_for(long_task(), timeout=5) except asyncio.TimeoutError: print("Task timed out")
asyncio.run(main())
常用类、方法和函数
1. 核心函数
| 方法/函数 | 说明 | 示例 |
|---|---|---|
asyncio.run(coro) | 运行异步主函数(Python 3.7+) | asyncio.run(main()) |
asyncio.create_task(coro) | 创建任务并加入事件循环 | task = asyncio.create_task(fetch_data()) |
asyncio.gather(*coros) | 并发运行多个协程 | await asyncio.gather(task1, task2) |
asyncio.sleep(delay) | 异步等待(非阻塞) | await asyncio.sleep(1) |
asyncio.wait(coros) | 控制任务完成方式 | done, pending = await asyncio.wait([task1, task2]) |
2. 事件循环(Event Loop)
| 方法 | 说明 | 示例 |
|---|---|---|
loop.run_until_complete(future) | 运行直到任务完成 | loop.run_until_complete(main()) |
loop.run_forever() | 永久运行事件循环 | loop.run_forever() |
loop.stop() | 停止事件循环 | loop.stop() |
loop.close() | 关闭事件循环 | loop.close() |
loop.call_soon(callback) | 安排回调函数立即执行 | loop.call_soon(print, "Hello") |
loop.call_later(delay, callback) | 延迟执行回调 | loop.call_later(5, callback) |
3. 协程(Coroutine)与任务(Task)
| 方法/装饰器 | 说明 | 示例 |
|---|---|---|
@asyncio.coroutine | 协程装饰器(旧版,Python 3.4-3.7) | @asyncio.coroutinedef old_coro(): |
async def | 定义协程(Python 3.5+) | async def fetch(): |
task.cancel() | 取消任务 | task.cancel() |
task.done() | 检查任务是否完成 | if task.done(): |
task.result() | 获取任务结果(需任务完成) | data = task.result() |
4. 同步原语(类似 threading)
| 类 | 说明 | 示例 |
|---|---|---|
asyncio.Lock() | 异步互斥锁 | lock = asyncio.Lock()async with lock: |
asyncio.Event() | 事件通知 | event = asyncio.Event()await event.wait() |
asyncio.Queue() | 异步队列 | queue = asyncio.Queue()await queue.put(item) |
asyncio.Semaphore() | 信号量 | sem = asyncio.Semaphore(5)async with sem: |
5. 网络与子进程
| 方法/类 | 说明 | 示例 |
|---|---|---|
asyncio.open_connection() | 建立 TCP 连接 | reader, writer = await asyncio.open_connection('host', 80) |
asyncio.start_server() | 创建 TCP 服务器 | server = await asyncio.start_server(handle, '0.0.0.0', 8888) |
asyncio.create_subprocess_exec() | 创建子进程 | proc = await asyncio.create_subprocess_exec('ls') |
6. 实用工具
| 方法 | 说明 | 示例 |
|---|---|---|
asyncio.current_task() | 获取当前任务 | task = asyncio.current_task() |
asyncio.all_tasks() | 获取所有任务 | tasks = asyncio.all_tasks() |
asyncio.shield(coro) | 保护任务不被取消 | await asyncio.shield(critical_task) |
asyncio.wait_for(coro, timeout) | 带超时的等待 | try: await asyncio.wait_for(task, 5)except asyncio.TimeoutError: |
更多推荐


所有评论(0)