1. 阻塞
事件循环通过一个线程依次推进 Task。协程遇到尚未完成的异步 IO 时会主动暂停,事件循环便可以运行其他任务;普通同步函数并不知道这套协作规则,它会一直占用当前线程,直到函数返回。
01import asyncio02import time0304def read_data():05time.sleep(2)06return "数据"0708async def heartbeat():09for _ in range(3):10print("心跳")11await asyncio.sleep(0.5)1213async def main():14heartbeat_task = asyncio.create_task(heartbeat())1516result = read_data()17print(result)1819await heartbeat_task2021asyncio.run(main())
read_data() 中的 time.sleep(2) 会阻塞事件循环线程,因此心跳不能每隔 0.5 秒正常执行。即使把调用它的函数写成 async def,同步函数也不会自动获得暂停能力。
常见阻塞操作包括:
- 只提供同步接口的 HTTP 或数据库客户端;
- 普通磁盘文件读写;
- 调用旧版 SDK 或操作系统命令;
- 压缩、图片处理和大型 JSON 解析;
- 调用内部包含长时间计算的普通函数。
解决问题时应该先寻找原生异步接口。原生异步库可以直接把 socket 等 IO 注册给事件循环,通常比占用一个线程等待更节省资源。只有暂时无法替换同步调用时,才考虑把它移到线程。
Task 与线程并不是同一种对象:
| 对比 | asyncio Task | 操作系统线程 |
|---|---|---|
| 调度者 | 事件循环 | 操作系统 |
| 切换方式 | 通常在 await 处协作切换 | 可以被操作系统抢占 |
| 内存开销 | 相对较小 | 相对较大 |
| 是否共享进程内存 | 是 | 是 |
| 常见用途 | 大量异步 IO | 隔离同步阻塞 IO |
线程不是事件循环的替代品,而是异步程序兼容同步阻塞代码的一座桥。
2. to_thread
asyncio.to_thread() 接收一个普通同步函数,把它交给独立线程执行,并返回一个可以等待的协程:
01import asyncio02import threading03import time0405def read_data(name):06print(f"{name} 运行于:{threading.current_thread().name}")07time.sleep(1)08return f"{name} 完成"0910async def main():11print(f"事件循环运行于:{threading.current_thread().name}")1213result = await asyncio.to_thread(read_data, "文件读取")14print(result)1516asyncio.run(main())
可以把执行过程理解为:
- 当前 Task 调用
to_thread(); - 同步函数被提交给事件循环使用的默认线程池;
- 当前 Task 等待线程池返回的 Future,并暂停执行;
- 事件循环继续处理其他 Task;
- 工作线程完成后,事件循环恢复等待结果的 Task。
to_thread() 不会为每次调用永久创建一个新线程。默认情况下,它借助事件循环的线程池执行,同一个线程可以被后续工作重复使用。
函数的位置参数和关键字参数都可以直接传入:
01import asyncio0203def format_message(message, *, prefix="[INFO]"):04return f"{prefix} {message}"0506async def main():07result = await asyncio.to_thread(08format_message,09"处理完成",10prefix="[任务]",11)12print(result)1314asyncio.run(main())
开始时机
to_thread() 返回的是协程对象。只调用它并不会保证同步函数已经开始,必须等待这个协程,或者把它创建成 Task:
01import asyncio02import time0304def blocking_work():05print("同步函数开始")06time.sleep(0.5)07return "完成"0809async def main():10coroutine = asyncio.to_thread(blocking_work)1112print("只创建了协程")13await asyncio.sleep(0.1)1415result = await coroutine16print(result)1718asyncio.run(main())
如果希望多个同步函数在同一段时间内运行,可以把这些协程交给 gather() 或 TaskGroup:
01import asyncio02import time0304def blocking_work(name):05time.sleep(1)06return f"{name} 完成"0708async def main():09results = await asyncio.gather(10asyncio.to_thread(blocking_work, "任务 A"),11asyncio.to_thread(blocking_work, "任务 B"),12)1314print(results)1516asyncio.run(main())
这里的等待时间可以重叠,但最终能同时运行多少函数仍受线程池容量限制。
3. 上下文
服务端经常使用 contextvars.ContextVar 保存请求编号、链路追踪编号或当前用户等上下文。to_thread() 会把调用方当前的 contextvars.Context 复制到工作线程:
01import asyncio02from contextvars import ContextVar0304request_id = ContextVar("request_id", default="unknown")0506def write_log():07print(f"request_id={request_id.get()}")0809async def main():10token = request_id.set("req-1001")1112try:13await asyncio.to_thread(write_log)14finally:15request_id.reset(token)1617asyncio.run(main())
工作线程能读取提交任务时的上下文副本。在线程中重新设置 ContextVar,不会把修改反向同步回原 Task;它不是跨线程共享可变状态的工具。
普通全局变量和对象则仍然属于同一个进程,线程之间可以看到并修改它们。是否传播 ContextVar 与是否共享进程内存,是两个不同问题。
这也是 to_thread() 相比手动提交线程池更方便的一点。底层 run_in_executor() 不会自动完成同样的上下文复制;需要时应手动使用 contextvars.copy_context(),或者优先使用 to_thread()。
4. Executor
asyncio.to_thread() 是高层快捷方式。如果需要指定线程池、使用进程池,或者更精细地控制提交方式,可以调用 loop.run_in_executor():
01import asyncio02import time0304def blocking_work(name):05time.sleep(0.5)06return f"{name} 完成"0708async def main():09loop = asyncio.get_running_loop()1011future = loop.run_in_executor(12None,13blocking_work,14"默认线程池",15)1617print(await future)1819asyncio.run(main())
第一个参数是 Executor。传入 None 表示使用事件循环的默认 ThreadPoolExecutor;如果默认线程池尚未创建,事件循环会按需创建。返回值是 asyncio.Future,可以直接 await。
和 to_thread() 不同,run_in_executor() 是普通方法,调用时就会把工作提交给 Executor,不需要等到第一次 await 才提交。
两者可以这样选择:
| 对比 | asyncio.to_thread() | loop.run_in_executor() |
|---|---|---|
| 抽象层级 | 高层 API | 低层事件循环 API |
| 返回对象 | 协程 | asyncio Future |
| 提交时机 | 协程开始执行时 | 调用方法时 |
| 关键字参数 | 直接支持 | 需要 functools.partial() |
| ContextVar | 自动复制当前上下文 | 需要调用者自行处理 |
| Executor | 使用默认线程池 | 可以显式指定线程池、进程池等 |
需要独立控制容量时,可以创建自己的线程池:
01import asyncio02import time03from concurrent.futures import ThreadPoolExecutor0405def blocking_work(name):06time.sleep(0.5)07return name0809async def main():10loop = asyncio.get_running_loop()1112with ThreadPoolExecutor(13max_workers=2,14thread_name_prefix="document",15) as pool:16futures = [17loop.run_in_executor(pool, blocking_work, name)18for name in ["A", "B", "C"]19]2021print(await asyncio.gather(*futures))2223asyncio.run(main())
with 代码块退出时会关闭线程池,并等待已经提交的工作结束。线程池是需要管理生命周期的资源,不应该在每次函数调用中重复创建。Web 服务通常在应用生命周期开始时创建,在关闭阶段统一释放。
也可以通过 loop.set_default_executor() 替换该事件循环的默认 Executor,但从 Python 3.11 开始,它必须是 ThreadPoolExecutor 实例或其子类。应用依赖框架时,不要在不了解框架线程池用途的情况下全局替换。
5. 容量
ThreadPoolExecutor 的 max_workers 表示最多同时运行多少个工作线程。超过容量的任务会进入 Executor 内部队列等待,不会继续创建无限线程。
在 Python 3.13 及当前版本中,未指定 max_workers 时,ThreadPoolExecutor 默认使用 min(32, (os.process_cpu_count() or 1) + 4)。这是标准库的默认值,不是适合所有业务的性能结论,应用不应该依赖它解决容量规划。
线程数过少会让阻塞调用排队,线程数过多则会增加内存、上下文切换和下游资源竞争。应该结合下面的信息决定:
- 同步操作平均阻塞多久;
- HTTP、数据库和文件句柄容量;
- 单个进程可能同时收到多少请求;
- 服务启动了多少 worker;
- 下游接口的并发与频率限制。
线程池控制的是正在运行的线程数量,但提交队列仍可能积压大量工作。面对没有上限的输入流,可以在提交前使用 Semaphore 或 Queue 建立反压:
01import asyncio02import time0304def blocking_work(item):05time.sleep(0.2)06return item0708async def run_limited(item, limit):09async with limit:10return await asyncio.to_thread(blocking_work, item)1112async def main():13limit = asyncio.Semaphore(4)1415results = await asyncio.gather(16*(run_limited(item, limit) for item in range(10))17)1819print(results)2021asyncio.run(main())
这段代码把当前业务的线程并发限制为 4,不会改变默认线程池供其他功能使用的总容量。Semaphore 是当前事件循环和当前进程内的限制,多 worker 部署时每个进程都会有一份。
还要避免在线程池工作中同步等待同一个线程池的另一个 Future。如果池中所有线程都在互相等待,剩余工作没有线程可以执行,就会形成线程池死锁。工作之间存在依赖时,应由事件循环在池外组织依赖关系。
6. 结果
同步函数的返回值会成为 await asyncio.to_thread(...) 的结果,抛出的异常也会在等待处重新抛出:
01import asyncio0203def parse_data():04raise ValueError("数据格式错误")0506async def main():07try:08await asyncio.to_thread(parse_data)09except ValueError as error:10print(error)1112asyncio.run(main())
这使线程调用在协程看来仍然像普通异步函数,可以继续使用 try...except、gather() 和 TaskGroup 管理结果。
取消线程
取消等待 to_thread() 的 Task,不等于强制停止已经运行的线程函数。Python 没有一个安全通用的 API 可以从外部终止任意线程:
01import asyncio02import time0304def slow_work():05print("线程开始")06time.sleep(1)07print("线程仍然完成")0809async def main():10task = asyncio.create_task(asyncio.to_thread(slow_work))1112await asyncio.sleep(0.1)13task.cancel()1415try:16await task17except asyncio.CancelledError:18print("等待线程的 Task 已取消")1920await asyncio.sleep(1)2122asyncio.run(main())
如果工作还在 Executor 队列里、尚未开始,取消有机会阻止它执行;一旦底层调用已经进入运行状态,就不能通过 Future 取消。asyncio.timeout() 也只能停止当前协程继续等待,不能杀死已经运行的线程。
需要支持停止时,同步函数本身必须协作检查一个线程安全信号:
01import asyncio02import threading0304def work(stop_event):05while not stop_event.wait(0.1):06print("处理一批数据")0708print("线程收到停止信号")0910async def main():11stop_event = threading.Event()12task = asyncio.create_task(13asyncio.to_thread(work, stop_event)14)1516await asyncio.sleep(0.35)17stop_event.set()18await task1920asyncio.run(main())
线程中使用的是 threading.Event,不是 asyncio.Event。如果第三方同步函数没有超时参数或停止机制,提交后就只能等待它自行返回,因此调用同步网络库时尤其要配置库自身的连接和读取超时。
7. GIL
GIL 是 CPython 中的全局解释器锁。在默认 CPython 构建中,同一进程通常只有一个线程可以在某个时刻执行 Python 字节码。因此,把纯 Python CPU 计算交给多个线程,通常不能获得多核并行加速,反而可能增加切换开销。
1def calculate():2total = 034for number in range(20_000_000):5total += number * number67return total
线程仍然适合阻塞 IO,因为线程等待网络或磁盘时不会持续执行 Python 字节码。某些使用 C、C++ 或 Rust 实现的扩展库也会在耗时计算期间主动释放 GIL,这时线程可能获得并行收益,但要以对应库的文档为准。
从 Python 3.13 开始,CPython 提供可以禁用 GIL 的 free-threaded 构建,使线程能够并行执行 Python 代码,但它不是默认构建。不能仅凭 Python 版本号就假设生产环境已经无 GIL,还要确认解释器构建、依赖兼容性和实际性能。
默认 CPython 中的纯 Python CPU 计算,通常使用 ProcessPoolExecutor:
01import asyncio02from concurrent.futures import ProcessPoolExecutor0304def calculate(limit):05return sum(number * number for number in range(limit))0607async def main():08loop = asyncio.get_running_loop()0910with ProcessPoolExecutor(max_workers=2) as pool:11results = await asyncio.gather(12loop.run_in_executor(pool, calculate, 10_000_000),13loop.run_in_executor(pool, calculate, 10_000_000),14)1516print(results)1718if __name__ == "__main__":19asyncio.run(main())
进程池中的函数、参数和返回值通常需要能够被 pickle 序列化,启动进程和传输数据也有额外成本。if __name__ == "__main__" 入口保护对于安全启动子进程非常重要。
Python 3.14 还提供了 InterpreterPoolExecutor。每个工作线程拥有隔离的解释器和独立 GIL,因此可以获得多核并行,但可变对象不能直接共享,调用和结果同样需要跨解释器传递。它是新的高级选项,不应在没有测试依赖兼容性和数据传输成本时直接替代进程池。
| 工作类型 | 优先选择 |
|---|---|
| 原生异步 HTTP、数据库或模型 SDK | 直接 await 异步接口 |
| 暂时无法替换的同步 IO | asyncio.to_thread() |
| 需要独立线程池配置的同步工作 | run_in_executor() 加 ThreadPoolExecutor |
| 默认 CPython 中的纯 Python CPU 计算 | ProcessPoolExecutor 或独立 worker |
| 长时间、需要重试和持久化的工作 | 消息队列与独立任务进程 |
8. 共享
线程共享同一个进程的列表、字典和对象,所以会产生数据竞态。GIL 不能把多条 Python 语句组成的业务操作自动变成原子事务。
01import asyncio02import time0304counter = 00506def increase():07global counter0809current = counter10time.sleep(0)11counter = current + 11213async def main():14await asyncio.gather(15*(asyncio.to_thread(increase) for _ in range(100))16)1718print(counter) # 可能小于 1001920asyncio.run(main())
多个线程可能读到相同的旧值,再分别写回同一个新值,导致部分更新丢失。线程共享状态应使用 threading.Lock:
01import asyncio02import threading0304counter = 005lock = threading.Lock()0607def increase():08global counter0910with lock:11counter += 11213async def main():14await asyncio.gather(15*(asyncio.to_thread(increase) for _ in range(100))16)1718print(counter) # 1001920asyncio.run(main())
两类锁不能互换:
| 锁 | 协调对象 | 等待方式 |
|---|---|---|
asyncio.Lock | 同一事件循环中的 Task | await,不会阻塞事件循环线程 |
threading.Lock | 同一进程中的线程 | 同步阻塞当前线程 |
不要从工作线程直接操作 asyncio.Lock,它不是线程安全对象;也不要在事件循环线程中同步等待竞争激烈的 threading.Lock,这会阻塞整个事件循环。让线程内同步逻辑使用 threading 原语,让协程侧的共享状态使用 asyncio 原语。
多进程 worker 不共享普通全局变量,进程内 threading.Lock 也无法保护其他进程。数据库状态仍然需要事务、约束、条件更新或其他跨进程协调机制。
9. 回调
线程中的函数不能随意调用事件循环对象。需要从其他线程安排普通回调时,使用 loop.call_soon_threadsafe():
01import asyncio02import threading03import time0405def complete(future, result):06if not future.done():07future.set_result(result)0809def worker(loop, future):10time.sleep(0.5)11loop.call_soon_threadsafe(12complete,13future,14"线程结果",15)1617async def main():18loop = asyncio.get_running_loop()19future = loop.create_future()2021thread = threading.Thread(22target=worker,23args=(loop, future),24)25thread.start()2627print(await future)28thread.join()2930asyncio.run(main())
工作线程没有直接调用 future.set_result(),而是让事件循环在线程安全入口中执行完成操作。取消 Task、设置 asyncio Future 结果或修改事件循环管理的状态时,也应遵循这个边界。
如果工作线程需要提交一整个协程,可以使用 asyncio.run_coroutine_threadsafe():
01import asyncio0203async def save_data(value):04await asyncio.sleep(0.2)05return f"已保存:{value}"0607def worker(loop):08future = asyncio.run_coroutine_threadsafe(09save_data("消息"),10loop,11)1213print(future.result(timeout=2))1415async def main():16loop = asyncio.get_running_loop()17await asyncio.to_thread(worker, loop)1819asyncio.run(main())
它返回的是 concurrent.futures.Future,适合在提交协程的那个工作线程中使用同步 result(timeout=...) 等待。不要在事件循环线程中调用这个 Future 的阻塞 result(),否则事件循环无法推进刚刚提交的协程,可能形成死锁。
协程抛出的异常会保存到这个 concurrent Future 中;工作线程等待 result() 时会重新收到异常,也可以调用 Future 的 cancel() 请求取消事件循环中的 Task。
10. FastAPI
FastAPI 会根据路径函数和依赖使用 def 还是 async def 采取不同方式:
async def路径函数直接运行在异步调用链中,内部应使用可等待的非阻塞接口;- 普通
def路径函数由 FastAPI 放到外部线程池执行并等待; - 普通
def依赖也会在线程池中执行; - 在路径函数里直接调用的普通工具函数不受 FastAPI 自动处理,仍然在当前线程同步运行。
最后一条尤其容易踩坑:
01import time02from fastapi import FastAPI0304app = FastAPI()0506def load_document():07time.sleep(2)08return "文档内容"0910@app.get("/document")11async def get_document():12return load_document() # 会阻塞事件循环
如果暂时只有同步工具,可以显式移到线程:
01import asyncio02import time03from fastapi import FastAPI0405app = FastAPI()0607def load_document():08time.sleep(2)09return "文档内容"1011@app.get("/document")12async def get_document():13return await asyncio.to_thread(load_document)
也可以把整个路径函数声明为普通 def,让 FastAPI 在线程池中调用它。路径函数使用原生异步 SDK 时,则保持 async def 并直接 await。
FastAPI 的线程池由框架的异步基础设施管理,不应假设它与 asyncio 默认 Executor 使用完全相同的容量配置。大量同步路径、同步依赖和手动 to_thread() 调用都可能消耗线程资源,需要分别观察实际运行环境。
在 LangChain 代码中也遵循相同原则:模型、向量数据库或工具已经提供原生异步接口时,直接使用 ainvoke() 等异步方法;只有底层集成确实是同步阻塞调用时,才考虑线程封装。把原生异步调用再套进 to_thread() 不会提升并发能力。
11. 选择
遇到一个可能耗时的函数,可以按下面的顺序判断:
- 它是否提供真正的异步接口?如果有,直接
await; - 它是否主要等待同步 IO?如果是,可以使用
to_thread(); - 是否需要独立的线程数、名称和生命周期?如果需要,使用自定义 ThreadPoolExecutor;
- 它是否执行纯 Python CPU 计算?默认 CPython 下优先考虑进程池;
- 工作是否需要跨请求继续、失败重试或进程重启后保留?如果需要,使用任务队列和独立 worker。
线程池解决的是「不要让同步函数阻塞事件循环」,并不会自动解决超时、重试、限流、数据一致性和任务可靠性。每项能力仍然需要单独设计。