python-asyncio与事件循环(Event Loop)
·
早期 Python(3.7 前)没有 asyncio.run(),需要手动写 “创建→启动→关闭” 流程:
asyncio.gather(*aws) 中的 aws(awaitables)参数只能是:
- 协程(
async def定义的函数) Task对象(通过asyncio.create_task()创建)Future对象(包括通过线程池包装的Future)-
# 获取事件循环 loop = asyncio.get_event_loop() # 运行协程直到完成 loop.run_until_complete(my_coroutine()) # 关闭事件循环 loop.close()
asyncio 高层 API 与底层事件循环方法对应表
| 功能分类 | 高层 API | 底层事件循环(Event Loop)方法 | 核心作用说明 |
|---|---|---|---|
| 事件循环启动 | asyncio.run(main()) | loop.run_until_complete(main()) + loop.close() | 自动创建事件循环、运行协程、结束后关闭循环,是 Python 3.7+ 推荐的启动方式,底层帮你完成了 “创建 - 运行 - 关闭” 的完整流程。 |
| 协程运行 | await coro(在协程内部) | loop.create_task(coro) / loop.ensure_future(coro) | 高层用 await 直接挂起 / 恢复协程,底层通过 create_task 将协程包装成 Task 并加入事件循环调度。 |
| 任务创建 | asyncio.create_task(coro) | loop.create_task(coro) | 完全对应,高层 API 只是将底层方法直接暴露,用于创建可并发调度的 Task 对象。 |
| 休眠(延时) | asyncio.sleep(secs) | loop.call_later(secs, callback) | 高层 sleep 是封装后的协程,底层通过 call_later 注册一个延时回调,到时间后唤醒协程。 |
| 同步原语 | asyncio.Lock() | loop.create_lock() | 高层 Lock 类内部调用底层 create_lock() 方法,实现协程间的互斥同步。 |
asyncio.Event() | loop.create_event() | 同理,高层 Event 依赖底层 create_event(),实现协程间的信号通知。 | |
| 回调注册 | asyncio.get_running_loop().call_soon(cb) | loop.call_soon(cb) | 高层需先通过 get_running_loop() 获取当前循环,再调用底层 call_soon,将回调加入 “立即执行” 队列。 |
asyncio.get_running_loop().call_later(secs, cb) | loop.call_later(secs, cb) | 同上,高层需先获取循环,再调用底层延时回调方法。 | |
| 子进程管理 | asyncio.create_subprocess_exec(...) | loop.subprocess_exec(...) | 高层 create_subprocess_exec 是对底层 subprocess_exec 的封装,简化子进程创建与 IO 重定向。 |
asyncio.create_subprocess_shell(cmd) | loop.subprocess_shell(cmd) | 同理,对应底层通过 Shell 执行命令的方法。 | |
| Future 操作 | asyncio.Future() | loop.create_future() | 高层 Future 类的实例化,本质是调用底层 create_future() 生成循环绑定的 Future 对象。 |
| 信号处理 | asyncio.get_running_loop().add_signal_handler(sig, cb) | loop.add_signal_handler(sig, cb) | 高层需先获取循环,再调用底层方法注册信号处理器,实现对系统信号(如 SIGINT)的响应。 |
task与future
| 类型 | 本质 | 适用场景 |
|---|---|---|
Future | 通用的 “结果占位符”(基类) | 对接非 asyncio 原生任务(进程 / 线程执行的同步函数) |
Task | 包装协程的 Future(子类) | 执行 async def 定义的协程(异步 IO 任务) |
Future:
async def main():
loop = asyncio.get_running_loop()
# 1. 创建进程池执行器
executor = ProcessPoolExecutor(max_workers=1)
# 2. 把进程池任务包装成 Future
future = loop.run_in_executor(executor, cpu_task)
# 3. 等待 Future 完成(本质是等子进程返回结果)
result = await future
print(f"子进程 PID: {result}")
executor.shutdown()
Task:
async def fetch_data(url):
await asyncio.sleep(1) # 模拟网络请求
return f"数据来自 {url}"
async def main():
# 创建Task(包装协程,自动加入事件循环调度)
task1 = asyncio.create_task(fetch_data("https://api.example.com/1"))
task2 = asyncio.create_task(fetch_data("https://api.example.com/2"))
# 并发执行,等待所有结果
result1 = await task1
result2 = await task2
print(result1) # 输出:数据来自 https://api.example.com/1
print(result2) # 输出:数据来自 https://api.example.com/2
asyncio.run(main())
run_in_threadpool
把同步阻塞的函数放到独立的线程中执行,返回一个可被 await 的异步对象(Future),让主线程的异步事件循环不被阻塞。
async def execute(*, db: AsyncSession, pk: int) -> None:
workers = await run_in_threadpool(celery_app.control.ping, timeout=0.5)
if not workers:
raise errors.ServerError(msg='Celery Worker 暂不可用,请稍后重试')
多个 Task 可以被事件循环并发调度(在 await 时切换)
事件循环会维护一个 “可运行的 Task 队列”,每次 await 后,当前 Task 会被移出队列(让出cpu控制器),等它等待的操作完成后,再重新加入队列,等待下一次调度。
asyncio实现异步处理长时间的耗时任务
更多推荐

所有评论(0)