早期 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实现异步处理长时间的耗时任务

更多推荐