创建时间: 2026-08-25最后更新: 2026-08-25

1. 阻塞

事件循环通过一个线程依次推进 Task。协程遇到尚未完成的异步 IO 时会主动暂停,事件循环便可以运行其他任务;普通同步函数并不知道这套协作规则,它会一直占用当前线程,直到函数返回。

blocking.py
01
import asyncio
02
import time
03
04
def read_data():
05
time.sleep(2)
06
return "数据"
07
08
async def heartbeat():
09
for _ in range(3):
10
print("心跳")
11
await asyncio.sleep(0.5)
12
13
async def main():
14
heartbeat_task = asyncio.create_task(heartbeat())
15
16
result = read_data()
17
print(result)
18
19
await heartbeat_task
20
21
asyncio.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() 接收一个普通同步函数,把它交给独立线程执行,并返回一个可以等待的协程:

to-thread.py
01
import asyncio
02
import threading
03
import time
04
05
def read_data(name):
06
print(f"{name} 运行于:{threading.current_thread().name}")
07
time.sleep(1)
08
return f"{name} 完成"
09
10
async def main():
11
print(f"事件循环运行于:{threading.current_thread().name}")
12
13
result = await asyncio.to_thread(read_data, "文件读取")
14
print(result)
15
16
asyncio.run(main())

可以把执行过程理解为:

  1. 当前 Task 调用 to_thread();
  2. 同步函数被提交给事件循环使用的默认线程池;
  3. 当前 Task 等待线程池返回的 Future,并暂停执行;
  4. 事件循环继续处理其他 Task;
  5. 工作线程完成后,事件循环恢复等待结果的 Task。

to_thread() 不会为每次调用永久创建一个新线程。默认情况下,它借助事件循环的线程池执行,同一个线程可以被后续工作重复使用。

函数的位置参数和关键字参数都可以直接传入:

to-thread-arguments.py
01
import asyncio
02
03
def format_message(message, *, prefix="[INFO]"):
04
return f"{prefix} {message}"
05
06
async def main():
07
result = await asyncio.to_thread(
08
format_message,
09
"处理完成",
10
prefix="[任务]",
11
)
12
print(result)
13
14
asyncio.run(main())

开始时机

to_thread() 返回的是协程对象。只调用它并不会保证同步函数已经开始,必须等待这个协程,或者把它创建成 Task:

to-thread-start.py
01
import asyncio
02
import time
03
04
def blocking_work():
05
print("同步函数开始")
06
time.sleep(0.5)
07
return "完成"
08
09
async def main():
10
coroutine = asyncio.to_thread(blocking_work)
11
12
print("只创建了协程")
13
await asyncio.sleep(0.1)
14
15
result = await coroutine
16
print(result)
17
18
asyncio.run(main())

如果希望多个同步函数在同一段时间内运行,可以把这些协程交给 gather() 或 TaskGroup:

concurrent-threads.py
01
import asyncio
02
import time
03
04
def blocking_work(name):
05
time.sleep(1)
06
return f"{name} 完成"
07
08
async def main():
09
results = await asyncio.gather(
10
asyncio.to_thread(blocking_work, "任务 A"),
11
asyncio.to_thread(blocking_work, "任务 B"),
12
)
13
14
print(results)
15
16
asyncio.run(main())

这里的等待时间可以重叠,但最终能同时运行多少函数仍受线程池容量限制。

3. 上下文

服务端经常使用 contextvars.ContextVar 保存请求编号、链路追踪编号或当前用户等上下文。to_thread() 会把调用方当前的 contextvars.Context 复制到工作线程:

thread-context.py
01
import asyncio
02
from contextvars import ContextVar
03
04
request_id = ContextVar("request_id", default="unknown")
05
06
def write_log():
07
print(f"request_id={request_id.get()}")
08
09
async def main():
10
token = request_id.set("req-1001")
11
12
try:
13
await asyncio.to_thread(write_log)
14
finally:
15
request_id.reset(token)
16
17
asyncio.run(main())

工作线程能读取提交任务时的上下文副本。在线程中重新设置 ContextVar,不会把修改反向同步回原 Task;它不是跨线程共享可变状态的工具。

普通全局变量和对象则仍然属于同一个进程,线程之间可以看到并修改它们。是否传播 ContextVar 与是否共享进程内存,是两个不同问题。

这也是 to_thread() 相比手动提交线程池更方便的一点。底层 run_in_executor() 不会自动完成同样的上下文复制;需要时应手动使用 contextvars.copy_context(),或者优先使用 to_thread()。

4. Executor

asyncio.to_thread() 是高层快捷方式。如果需要指定线程池、使用进程池,或者更精细地控制提交方式,可以调用 loop.run_in_executor():

run-in-executor.py
01
import asyncio
02
import time
03
04
def blocking_work(name):
05
time.sleep(0.5)
06
return f"{name} 完成"
07
08
async def main():
09
loop = asyncio.get_running_loop()
10
11
future = loop.run_in_executor(
12
None,
13
blocking_work,
14
"默认线程池",
15
)
16
17
print(await future)
18
19
asyncio.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使用默认线程池可以显式指定线程池、进程池等

需要独立控制容量时,可以创建自己的线程池:

custom-thread-pool.py
01
import asyncio
02
import time
03
from concurrent.futures import ThreadPoolExecutor
04
05
def blocking_work(name):
06
time.sleep(0.5)
07
return name
08
09
async def main():
10
loop = asyncio.get_running_loop()
11
12
with ThreadPoolExecutor(
13
max_workers=2,
14
thread_name_prefix="document",
15
) as pool:
16
futures = [
17
loop.run_in_executor(pool, blocking_work, name)
18
for name in ["A", "B", "C"]
19
]
20
21
print(await asyncio.gather(*futures))
22
23
asyncio.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 建立反压:

limit-to-thread.py
01
import asyncio
02
import time
03
04
def blocking_work(item):
05
time.sleep(0.2)
06
return item
07
08
async def run_limited(item, limit):
09
async with limit:
10
return await asyncio.to_thread(blocking_work, item)
11
12
async def main():
13
limit = asyncio.Semaphore(4)
14
15
results = await asyncio.gather(
16
*(run_limited(item, limit) for item in range(10))
17
)
18
19
print(results)
20
21
asyncio.run(main())

这段代码把当前业务的线程并发限制为 4,不会改变默认线程池供其他功能使用的总容量。Semaphore 是当前事件循环和当前进程内的限制,多 worker 部署时每个进程都会有一份。

还要避免在线程池工作中同步等待同一个线程池的另一个 Future。如果池中所有线程都在互相等待,剩余工作没有线程可以执行,就会形成线程池死锁。工作之间存在依赖时,应由事件循环在池外组织依赖关系。

6. 结果

同步函数的返回值会成为 await asyncio.to_thread(...) 的结果,抛出的异常也会在等待处重新抛出:

thread-exception.py
01
import asyncio
02
03
def parse_data():
04
raise ValueError("数据格式错误")
05
06
async def main():
07
try:
08
await asyncio.to_thread(parse_data)
09
except ValueError as error:
10
print(error)
11
12
asyncio.run(main())

这使线程调用在协程看来仍然像普通异步函数,可以继续使用 try...except、gather() 和 TaskGroup 管理结果。

取消线程

取消等待 to_thread() 的 Task,不等于强制停止已经运行的线程函数。Python 没有一个安全通用的 API 可以从外部终止任意线程:

cancel-to-thread.py
01
import asyncio
02
import time
03
04
def slow_work():
05
print("线程开始")
06
time.sleep(1)
07
print("线程仍然完成")
08
09
async def main():
10
task = asyncio.create_task(asyncio.to_thread(slow_work))
11
12
await asyncio.sleep(0.1)
13
task.cancel()
14
15
try:
16
await task
17
except asyncio.CancelledError:
18
print("等待线程的 Task 已取消")
19
20
await asyncio.sleep(1)
21
22
asyncio.run(main())

如果工作还在 Executor 队列里、尚未开始,取消有机会阻止它执行;一旦底层调用已经进入运行状态,就不能通过 Future 取消。asyncio.timeout() 也只能停止当前协程继续等待,不能杀死已经运行的线程。

需要支持停止时,同步函数本身必须协作检查一个线程安全信号:

cooperative-stop.py
01
import asyncio
02
import threading
03
04
def work(stop_event):
05
while not stop_event.wait(0.1):
06
print("处理一批数据")
07
08
print("线程收到停止信号")
09
10
async def main():
11
stop_event = threading.Event()
12
task = asyncio.create_task(
13
asyncio.to_thread(work, stop_event)
14
)
15
16
await asyncio.sleep(0.35)
17
stop_event.set()
18
await task
19
20
asyncio.run(main())

线程中使用的是 threading.Event,不是 asyncio.Event。如果第三方同步函数没有超时参数或停止机制,提交后就只能等待它自行返回,因此调用同步网络库时尤其要配置库自身的连接和读取超时。

7. GIL

GIL 是 CPython 中的全局解释器锁。在默认 CPython 构建中,同一进程通常只有一个线程可以在某个时刻执行 Python 字节码。因此,把纯 Python CPU 计算交给多个线程,通常不能获得多核并行加速,反而可能增加切换开销。

cpu-bound.py
1
def calculate():
2
total = 0
3
4
for number in range(20_000_000):
5
total += number * number
6
7
return total

线程仍然适合阻塞 IO,因为线程等待网络或磁盘时不会持续执行 Python 字节码。某些使用 C、C++ 或 Rust 实现的扩展库也会在耗时计算期间主动释放 GIL,这时线程可能获得并行收益,但要以对应库的文档为准。

从 Python 3.13 开始,CPython 提供可以禁用 GIL 的 free-threaded 构建,使线程能够并行执行 Python 代码,但它不是默认构建。不能仅凭 Python 版本号就假设生产环境已经无 GIL,还要确认解释器构建、依赖兼容性和实际性能。

默认 CPython 中的纯 Python CPU 计算,通常使用 ProcessPoolExecutor:

process-pool.py
01
import asyncio
02
from concurrent.futures import ProcessPoolExecutor
03
04
def calculate(limit):
05
return sum(number * number for number in range(limit))
06
07
async def main():
08
loop = asyncio.get_running_loop()
09
10
with ProcessPoolExecutor(max_workers=2) as pool:
11
results = await asyncio.gather(
12
loop.run_in_executor(pool, calculate, 10_000_000),
13
loop.run_in_executor(pool, calculate, 10_000_000),
14
)
15
16
print(results)
17
18
if __name__ == "__main__":
19
asyncio.run(main())

进程池中的函数、参数和返回值通常需要能够被 pickle 序列化,启动进程和传输数据也有额外成本。if __name__ == "__main__" 入口保护对于安全启动子进程非常重要。

Python 3.14 还提供了 InterpreterPoolExecutor。每个工作线程拥有隔离的解释器和独立 GIL,因此可以获得多核并行,但可变对象不能直接共享,调用和结果同样需要跨解释器传递。它是新的高级选项,不应在没有测试依赖兼容性和数据传输成本时直接替代进程池。

工作类型优先选择
原生异步 HTTP、数据库或模型 SDK直接 await 异步接口
暂时无法替换的同步 IOasyncio.to_thread()
需要独立线程池配置的同步工作run_in_executor() 加 ThreadPoolExecutor
默认 CPython 中的纯 Python CPU 计算ProcessPoolExecutor 或独立 worker
长时间、需要重试和持久化的工作消息队列与独立任务进程

8. 共享

线程共享同一个进程的列表、字典和对象,所以会产生数据竞态。GIL 不能把多条 Python 语句组成的业务操作自动变成原子事务。

thread-race.py
01
import asyncio
02
import time
03
04
counter = 0
05
06
def increase():
07
global counter
08
09
current = counter
10
time.sleep(0)
11
counter = current + 1
12
13
async def main():
14
await asyncio.gather(
15
*(asyncio.to_thread(increase) for _ in range(100))
16
)
17
18
print(counter) # 可能小于 100
19
20
asyncio.run(main())

多个线程可能读到相同的旧值,再分别写回同一个新值,导致部分更新丢失。线程共享状态应使用 threading.Lock:

thread-lock.py
01
import asyncio
02
import threading
03
04
counter = 0
05
lock = threading.Lock()
06
07
def increase():
08
global counter
09
10
with lock:
11
counter += 1
12
13
async def main():
14
await asyncio.gather(
15
*(asyncio.to_thread(increase) for _ in range(100))
16
)
17
18
print(counter) # 100
19
20
asyncio.run(main())

两类锁不能互换:

锁协调对象等待方式
asyncio.Lock同一事件循环中的 Taskawait,不会阻塞事件循环线程
threading.Lock同一进程中的线程同步阻塞当前线程

不要从工作线程直接操作 asyncio.Lock,它不是线程安全对象;也不要在事件循环线程中同步等待竞争激烈的 threading.Lock,这会阻塞整个事件循环。让线程内同步逻辑使用 threading 原语,让协程侧的共享状态使用 asyncio 原语。

多进程 worker 不共享普通全局变量,进程内 threading.Lock 也无法保护其他进程。数据库状态仍然需要事务、约束、条件更新或其他跨进程协调机制。

9. 回调

线程中的函数不能随意调用事件循环对象。需要从其他线程安排普通回调时,使用 loop.call_soon_threadsafe():

call-soon-threadsafe.py
01
import asyncio
02
import threading
03
import time
04
05
def complete(future, result):
06
if not future.done():
07
future.set_result(result)
08
09
def worker(loop, future):
10
time.sleep(0.5)
11
loop.call_soon_threadsafe(
12
complete,
13
future,
14
"线程结果",
15
)
16
17
async def main():
18
loop = asyncio.get_running_loop()
19
future = loop.create_future()
20
21
thread = threading.Thread(
22
target=worker,
23
args=(loop, future),
24
)
25
thread.start()
26
27
print(await future)
28
thread.join()
29
30
asyncio.run(main())

工作线程没有直接调用 future.set_result(),而是让事件循环在线程安全入口中执行完成操作。取消 Task、设置 asyncio Future 结果或修改事件循环管理的状态时,也应遵循这个边界。

如果工作线程需要提交一整个协程,可以使用 asyncio.run_coroutine_threadsafe():

run-coroutine-threadsafe.py
01
import asyncio
02
03
async def save_data(value):
04
await asyncio.sleep(0.2)
05
return f"已保存:{value}"
06
07
def worker(loop):
08
future = asyncio.run_coroutine_threadsafe(
09
save_data("消息"),
10
loop,
11
)
12
13
print(future.result(timeout=2))
14
15
async def main():
16
loop = asyncio.get_running_loop()
17
await asyncio.to_thread(worker, loop)
18
19
asyncio.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 自动处理,仍然在当前线程同步运行。

最后一条尤其容易踩坑:

fastapi-blocking-helper.py
01
import time
02
from fastapi import FastAPI
03
04
app = FastAPI()
05
06
def load_document():
07
time.sleep(2)
08
return "文档内容"
09
10
@app.get("/document")
11
async def get_document():
12
return load_document() # 会阻塞事件循环

如果暂时只有同步工具,可以显式移到线程:

fastapi-to-thread.py
01
import asyncio
02
import time
03
from fastapi import FastAPI
04
05
app = FastAPI()
06
07
def load_document():
08
time.sleep(2)
09
return "文档内容"
10
11
@app.get("/document")
12
async def get_document():
13
return await asyncio.to_thread(load_document)

也可以把整个路径函数声明为普通 def,让 FastAPI 在线程池中调用它。路径函数使用原生异步 SDK 时,则保持 async def 并直接 await。

FastAPI 的线程池由框架的异步基础设施管理,不应假设它与 asyncio 默认 Executor 使用完全相同的容量配置。大量同步路径、同步依赖和手动 to_thread() 调用都可能消耗线程资源,需要分别观察实际运行环境。

在 LangChain 代码中也遵循相同原则:模型、向量数据库或工具已经提供原生异步接口时,直接使用 ainvoke() 等异步方法;只有底层集成确实是同步阻塞调用时,才考虑线程封装。把原生异步调用再套进 to_thread() 不会提升并发能力。

11. 选择

遇到一个可能耗时的函数,可以按下面的顺序判断:

  1. 它是否提供真正的异步接口?如果有,直接 await;
  2. 它是否主要等待同步 IO?如果是,可以使用 to_thread();
  3. 是否需要独立的线程数、名称和生命周期?如果需要,使用自定义 ThreadPoolExecutor;
  4. 它是否执行纯 Python CPU 计算?默认 CPython 下优先考虑进程池;
  5. 工作是否需要跨请求继续、失败重试或进程重启后保留?如果需要,使用任务队列和独立 worker。

线程池解决的是「不要让同步函数阻塞事件循环」,并不会自动解决超时、重试、限流、数据一致性和任务可靠性。每项能力仍然需要单独设计。

正在验证登录状态
请稍候,验证完成后将继续显示文章内容