Python
用 `asyncio.to_thread` 包装不支持异步的工具
这篇笔记记录一下我在写 AI 工作流项目时,突然想通的一个问题。 当时项目里有很多异步函数,到处都能看到 async 和 await。但是有些工具调用并不是直接 await tool(),而是写成: 我一开始有点疑惑:为什么不直接 await? 后来才发现,不是所有工具都提供异步方法。有些库只有普通的同步函数。直接调用

这篇笔记记录一下我在写 AI 工作流项目时,突然想通的一个问题。
当时项目里有很多异步函数,到处都能看到 async 和 await。但是有些工具调用并不是直接 await tool(),而是写成:
result = await asyncio.to_thread(tool)我一开始有点疑惑:为什么不直接 await?
后来才发现,不是所有工具都提供异步方法。有些库只有普通的同步函数。直接调用它们会卡住事件循环,这时就可以用 asyncio.to_thread() 包一层,把阻塞操作交给另一个线程执行。
asyncio.to_thread() 从 Python 3.9 开始提供。
先说结论
如果一个工具本身提供异步方法,直接 await:
result = await client.fetch_data()如果一个工具只有同步方法,而且执行时可能需要等待文件、网络或者其他 I/O,可以用 asyncio.to_thread():
result = await asyncio.to_thread(client.fetch_data)要注意,to_thread 不是把同步函数变成了真正的异步函数,也不是把它放进事件循环里运行。
它做的事情更像是:
- 把同步函数交给另一个线程执行。
- 当前协程先暂停等待结果。
- 事件循环继续处理其他任务。
- 线程执行完成后,当前协程再继续往下走。
为什么不能直接 await
await 后面需要放一个可等待对象,比如协程。
一个普通同步函数不是协程:
import time
def slow_tool():
time.sleep(3)
return "执行完成"下面这样写是不对的:
result = await slow_tool()因为 Python 会先正常调用 slow_tool()。三秒后,它返回一个普通字符串。然后程序相当于尝试:
await "执行完成"字符串不能被 await。
更麻烦的是,在这三秒里,事件循环已经被阻塞了。
什么叫阻塞事件循环
可以把事件循环理解成一个一直在处理任务的调度员。
有些任务执行到等待网络响应的位置时,会主动让出执行权:
await asyncio.sleep(3)事件循环不会在这里傻等三秒,而是先去处理其他任务。
但是普通同步函数不知道怎么让出执行权:
time.sleep(3)如果在异步函数里直接调用它,整个事件循环都会停在这里。
看一个例子:
import asyncio
import time
def slow_tool():
time.sleep(3)
return "工具执行完成"
async def print_status():
for index in range(3):
await asyncio.sleep(1)
print(f"状态更新:{index + 1}")
async def main():
status_task = asyncio.create_task(print_status())
result = slow_tool()
print(result)
await status_task
asyncio.run(main())运行时,前三秒什么都不会打印。slow_tool() 执行完成后,状态更新才会继续出现。
这说明同步函数把事件循环堵住了。
用 asyncio.to_thread 包一层
改成:
import asyncio
import time
def slow_tool():
time.sleep(3)
return "工具执行完成"
async def print_status():
for index in range(3):
await asyncio.sleep(1)
print(f"状态更新:{index + 1}")
async def main():
status_task = asyncio.create_task(print_status())
result = await asyncio.to_thread(slow_tool)
print(result)
await status_task
asyncio.run(main())这次 slow_tool() 会在线程里等待。事件循环没有被堵住,所以每隔一秒仍然可以打印一次状态更新。
状态更新:1
状态更新:2
状态更新:3
工具执行完成到第三秒时,最后一次状态更新和工具执行完成的顺序不一定固定。重要的是,等待工具期间,其他协程仍然可以继续执行。
在 AI 工作流里有什么用
AI 工作流通常不只是调用一次模型,还会串联很多工具:
- 请求模型 API。
- 查询数据库。
- 读取本地文件。
- 调用第三方 SDK。
- 执行向量检索。
- 调用浏览器自动化工具。
- 运行已有的同步脚本。
有些库提供了异步方法:
result = await async_client.search(query)这种情况直接用异步方法就好。
但是有些工具只有同步方法:
result = sync_client.search(query)如果它需要等待网络响应,直接放进异步工作流里调用,就会卡住其他任务。可以改成:
result = await asyncio.to_thread(sync_client.search, query)to_thread 后面可以继续传参数:
result = await asyncio.to_thread(
sync_client.search,
query,
limit=10,
)这对接入旧代码或者第三方 SDK 很有用。工具本身不支持异步,也不一定要重写一遍。
同时执行多个同步工具
如果几个工具彼此没有依赖关系,可以一起执行:
import asyncio
async def collect_context(query: str):
web_result, file_result, database_result = await asyncio.gather(
asyncio.to_thread(search_web, query),
asyncio.to_thread(search_files, query),
asyncio.to_thread(query_database, query),
)
return {
"web": web_result,
"files": file_result,
"database": database_result,
}这几个同步工具会在线程中执行。只要它们主要是在等待 I/O,就能减少互相等待的时间。
我之前对多线程的疑惑
我以前知道 Python 有 GIL,也就是全局解释器锁。
我的理解一直是:同一时间只有一个线程能执行 Python 字节码,那多线程和单线程好像也差不多,多线程是不是没什么用?
后来我发现,这个理解只说对了一部分。
如果任务是大量计算,比如用 Python 循环做复杂运算:
def calculate():
total = 0
for number in range(100_000_000):
total += number * number
return total这种 CPU 密集型任务,就算开多个线程,通常也不会因为线程变多而明显加速。因为多个线程仍然会受到 GIL 限制。
但是很多程序慢,不是因为 CPU 一直在计算,而是在等待:
- 等网络响应。
- 等磁盘读写。
- 等数据库返回结果。
- 等第三方 SDK 完成请求。
- 等操作系统处理某个 I/O 操作。
一个线程等待 I/O 时,其他线程仍然可以工作。
所以多线程不是多余的。它很适合处理同步阻塞的 I/O 工具,尤其适合把这些工具接进异步项目。
CPU 密集型任务怎么办
asyncio.to_thread() 默认更适合 I/O 密集型任务。
如果是普通 Python 代码做大量计算,优先考虑进程池:
import asyncio
from concurrent.futures import ProcessPoolExecutor
def calculate():
return sum(number * number for number in range(100_000_000))
async def main():
loop = asyncio.get_running_loop()
with ProcessPoolExecutor() as pool:
result = await loop.run_in_executor(pool, calculate)
print(result)
asyncio.run(main())不同进程有各自的 Python 解释器,可以绕开同一个 GIL 的限制。
也有一些特殊情况。比如某些用 C、C++ 或 Rust 编写的扩展模块,在执行耗时操作时会主动释放 GIL。这类函数即使是 CPU 密集型,也可能从线程中受益。
to_thread 不是无限开线程
asyncio.to_thread() 使用线程池执行任务,不是每调用一次就无条件新建一个线程。
但是线程池容量还是有限的。如果一下子塞进去太多阻塞任务,后面的任务仍然需要排队。
例如工作流里一次要处理几百个文件,可以用 Semaphore 限制并发数量:
import asyncio
semaphore = asyncio.Semaphore(10)
async def read_one_file(path: str):
async with semaphore:
return await asyncio.to_thread(read_file, path)
async def read_all_files(paths: list[str]):
return await asyncio.gather(
*(read_one_file(path) for path in paths)
)这样同一时间最多处理十个文件,不会一下把太多任务塞进线程池。
还有几个容易踩的坑
1. 优先使用库本身提供的异步方法
如果一个库同时提供同步和异步客户端,优先使用异步版本:
result = await async_client.search(query)to_thread 更适合兼容没有异步接口的工具,不需要什么都包一层。
2. await to_thread() 不代表当前任务不等待
下面这段代码仍然会等待 slow_tool() 执行完成:
result = await asyncio.to_thread(slow_tool)只是等待期间,事件循环可以去处理别的协程。
如果希望当前代码先继续做其他事情,可以先创建任务:
tool_task = asyncio.create_task(
asyncio.to_thread(slow_tool)
)
print("先做其他事情")
result = await tool_task3. 取消协程不一定能立刻停掉线程
如果外层异步任务被取消,正在运行的同步函数不一定会马上停止。
线程里的函数没有自动收到一套通用的强制终止机制。所以对耗时很长的工具,最好使用工具本身支持的超时参数,或者把任务拆小。
4. 线程里的代码仍然要注意线程安全
如果多个线程同时修改同一个全局变量、缓存或者文件,还是可能出现问题。
to_thread 解决的是事件循环被阻塞的问题,不会自动解决线程安全问题。
和 run_in_executor 的关系
在 asyncio.to_thread() 出现以前,也可以写:
loop = asyncio.get_running_loop()
result = await loop.run_in_executor(None, slow_tool)asyncio.to_thread() 可以理解成更方便的高层写法:
result = await asyncio.to_thread(slow_tool)如果只是想把一个同步阻塞函数放到线程里,优先用 to_thread,代码更清楚。
如果需要自己管理线程池、进程池或者执行器,再考虑 run_in_executor()。
我最后记住的判断方式
碰到一个工具函数时,我可以先问自己:
- 它本身有没有异步方法?
- 它主要是在等待 I/O,还是在做大量计算?
- 它会不会阻塞事件循环?
可以简单记成:
| 场景 | 处理方式 |
|---|---|
| 工具本身提供异步方法 | 直接 await |
| 只有同步方法,主要等待 I/O | await asyncio.to_thread(...) |
| 普通 Python 代码做大量计算 | 考虑进程池 |
| 少量很快的同步操作 | 直接调用通常也没问题 |
我之前觉得 Python 有 GIL,多线程好像没什么必要。现在想通以后,感觉它的位置很清楚:
多线程不一定是为了让 Python 同时做更多计算,也可以是为了让一个线程等待的时候,其他任务还能继续执行。
在异步工作流里,asyncio.to_thread() 就是一个很好用的连接方式。它可以把不支持异步的同步工具接进来,又不把整个事件循环堵住。