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

1. 并发

并发不是让所有代码在同一时刻运行,而是让多个尚未完成的工作在同一段时间内交替推进。

假设程序要发出三个相互独立的网络请求,每个请求都需要等待一秒。顺序执行时,第二个请求必须等第一个完成后才开始:

sequential.py
01
import asyncio
02
import time
03
04
async def fetch(name):
05
await asyncio.sleep(1)
06
return f"{name} 完成"
07
08
async def main():
09
started_at = time.perf_counter()
10
11
first = await fetch("A")
12
second = await fetch("B")
13
third = await fetch("C")
14
15
print(first, second, third)
16
print(f"耗时:{time.perf_counter() - started_at:.1f} 秒")
17
18
asyncio.run(main())

总耗时接近三秒。三个请求之间没有依赖关系,所以可以先把它们都安排到事件循环,再一起等待:

concurrent.py
01
import asyncio
02
import time
03
04
async def fetch(name):
05
await asyncio.sleep(1)
06
return f"{name} 完成"
07
08
async def main():
09
started_at = time.perf_counter()
10
11
results = await asyncio.gather(
12
fetch("A"),
13
fetch("B"),
14
fetch("C"),
15
)
16
17
print(results)
18
print(f"耗时:{time.perf_counter() - started_at:.1f} 秒")
19
20
asyncio.run(main())

总耗时接近一秒,因为三个 Task 的等待时间发生了重叠。事件循环仍然一次只推进一个 Task:某个 Task 遇到尚未完成的 await 后暂停,事件循环再推进其他 Task。

需要区分三个容易混淆的概念:

概念含义
顺序一个工作结束后才开始下一个
并发多个工作在同一段时间内交替推进
并行多个工作在同一时刻分别占用不同 CPU 核心或执行单元

asyncio 主要解决 IO 并发,不会自动让 Python 计算并行。网络请求、数据库查询和模型调用在等待外部响应时可以让出事件循环;一个没有 await 的长计算仍然会占住事件循环线程。

并发也不是越多越快。创建更多任务只能增加正在等待或争抢资源的工作数量,无法突破对方服务限流、数据库连接池、网络带宽和本机 CPU 的上限。

2. gather

asyncio.gather() 接收多个可等待对象,并等待它们全部完成。传入协程对象时,gather() 会自动把它们安排为 Task:

gather.py
01
import asyncio
02
03
async def fetch(name, delay):
04
await asyncio.sleep(delay)
05
print(f"{name} 完成")
06
return name
07
08
async def main():
09
results = await asyncio.gather(
10
fetch("A", 1),
11
fetch("B", 0.2),
12
fetch("C", 0.5),
13
)
14
15
print(results) # ['A', 'B', 'C']
16
17
asyncio.run(main())

任务会按照 B、C、A 的顺序完成,但结果仍按照参数 A、B、C 的顺序排列。gather() 适合「提交一组工作,最后按输入顺序取得全部结果」的场景。

失败行为

默认情况下,某个子任务首次抛出异常时,等待 gather() 的调用方会立即收到这个异常。其余子任务不会因此自动取消,通常仍会继续运行:

gather-error.py
01
import asyncio
02
03
async def fail():
04
await asyncio.sleep(0.1)
05
raise ValueError("请求失败")
06
07
async def slow():
08
await asyncio.sleep(0.3)
09
print("慢请求仍然完成")
10
11
async def main():
12
slow_task = asyncio.create_task(slow())
13
14
try:
15
await asyncio.gather(fail(), slow_task)
16
except ValueError as error:
17
print(error)
18
19
await slow_task
20
21
asyncio.run(main())

这点很容易被误解:gather() 把首个异常传播给调用者,不等于替调用者关闭所有其他工作。

如果每项工作可以独立成功或失败,可以设置 return_exceptions=True:

gather-return-exceptions.py
01
import asyncio
02
03
async def success():
04
return "成功"
05
06
async def fail():
07
raise ValueError("失败")
08
09
async def main():
10
results = await asyncio.gather(
11
success(),
12
fail(),
13
return_exceptions=True,
14
)
15
16
for result in results:
17
if isinstance(result, BaseException):
18
print(f"异常:{result}")
19
else:
20
print(f"结果:{result}")
21
22
asyncio.run(main())

此时异常对象会和正常结果一起出现在列表中。调用者必须主动识别并处理异常,否则可能把失败误当作业务数据。

如果等待 gather() 的调用本身被取消,尚未完成的子任务会收到取消请求。反过来,单个子任务被取消不会直接取消其他子任务。取消行为关系到资源清理,协程仍然应该用 finally 关闭连接和临时资源。

3. TaskGroup

当一组任务共同完成一个业务目标,通常更适合使用 Python 3.11 提供的 TaskGroup:

task-group.py
01
import asyncio
02
03
async def fetch(name, delay):
04
await asyncio.sleep(delay)
05
return f"{name} 完成"
06
07
async def main():
08
async with asyncio.TaskGroup() as group:
09
first = group.create_task(fetch("A", 0.3))
10
second = group.create_task(fetch("B", 0.1))
11
12
print(first.result())
13
print(second.result())
14
15
asyncio.run(main())

退出 async with 之前,TaskGroup 会等待所有子任务结束。任务的创建范围同时也是负责等待和清理它们的范围,这种组织方式叫结构化并发。

TaskGroup 与 gather() 最重要的差别在于失败处理:一个子任务异常失败后,TaskGroup 会取消其余尚未完成的子任务,等待它们完成清理,再抛出异常组。

task-group-error.py
01
import asyncio
02
03
async def fail():
04
await asyncio.sleep(0.1)
05
raise ValueError("查询失败")
06
07
async def slow():
08
try:
09
await asyncio.sleep(10)
10
finally:
11
print("慢任务释放资源")
12
13
async def main():
14
try:
15
async with asyncio.TaskGroup() as group:
16
group.create_task(fail())
17
group.create_task(slow())
18
except* ValueError as errors:
19
for error in errors.exceptions:
20
print(error)
21
22
asyncio.run(main())

可以按下面的原则选择:

场景更合适的方式
每项工作可以独立失败,需要收集全部结果gather(return_exceptions=True)
一个任务失败后,其余工作已经没有意义TaskGroup
需要最简单地按输入顺序得到全部结果gather()
需要清晰管理一组子任务的生命周期TaskGroup

不要为了“看起来更并发”而把有依赖关系的步骤放入同一个 TaskGroup。第二步必须使用第一步的结果时,它们本来就应该顺序执行。

4. 完成

gather() 会等所有任务结束后一次返回。如果希望谁先完成就先处理谁,可以使用 asyncio.as_completed()。

从 Python 3.13 开始,它支持异步迭代;传入 Task 后,异步迭代会按完成顺序产出原来的 Task 对象:

as-completed.py
01
import asyncio
02
03
async def fetch(name, delay):
04
await asyncio.sleep(delay)
05
return f"{name} 的结果"
06
07
async def main():
08
tasks = [
09
asyncio.create_task(fetch("慢任务", 1), name="slow"),
10
asyncio.create_task(fetch("快任务", 0.2), name="fast"),
11
asyncio.create_task(fetch("中等任务", 0.5), name="medium"),
12
]
13
14
async for task in asyncio.as_completed(tasks):
15
result = await task
16
print(task.get_name(), result)
17
18
asyncio.run(main())

输出顺序通常是 fast、medium、slow。因为取得的是原来的 Task,还能通过名称或映射判断结果属于哪个请求。

为了兼容 Python 3.12 及更早版本,也可以使用普通 for:

as-completed-compatible.py
01
import asyncio
02
03
async def fetch(name, delay):
04
await asyncio.sleep(delay)
05
return name
06
07
async def main():
08
tasks = [
09
asyncio.create_task(fetch("A", 0.5)),
10
asyncio.create_task(fetch("B", 0.1)),
11
]
12
13
for next_result in asyncio.as_completed(tasks):
14
result = await next_result
15
print(result)
16
17
asyncio.run(main())

普通迭代产出的是需要等待的新协程,不是原始 Task。两种写法都按完成顺序交付结果,但对象身份不同。

as_completed(tasks, timeout=...) 可以限制等待整组任务的最长时间。超时会抛出 TimeoutError,但不会自动取消剩余任务;调用者需要决定继续等待、取消还是把它们交给其他生命周期管理逻辑。

这种方式适合搜索多个数据源、并发下载、批量模型调用后逐个推送结果等场景。它改变的是结果处理顺序,不会自动限制并发数量。

5. 等待

asyncio.wait() 提供更底层的等待控制。它返回两个集合:已经完成的 done 和仍未完成的 pending。

wait.py
01
import asyncio
02
03
async def fetch(name, delay):
04
await asyncio.sleep(delay)
05
return name
06
07
async def main():
08
tasks = {
09
asyncio.create_task(fetch("A", 0.2)),
10
asyncio.create_task(fetch("B", 1)),
11
}
12
13
done, pending = await asyncio.wait(
14
tasks,
15
return_when=asyncio.FIRST_COMPLETED,
16
)
17
18
for task in done:
19
print(f"最先完成:{await task}")
20
21
for task in pending:
22
task.cancel()
23
24
await asyncio.gather(*pending, return_exceptions=True)
25
26
asyncio.run(main())

return_when 支持三种条件:

条件返回时机
FIRST_COMPLETED任意任务完成或取消
FIRST_EXCEPTION首次出现异常;如果始终没有异常,则等待全部完成
ALL_COMPLETED全部任务完成或取消,默认值

从 Python 3.11 开始,不能把裸协程对象直接传给 wait(),需要先使用 create_task() 创建 Task。

给 wait() 设置 timeout 时,它不会抛出 TimeoutError,也不会自动取消未完成任务,只会把它们放入 pending:

wait-timeout.py
01
import asyncio
02
03
async def work(delay):
04
await asyncio.sleep(delay)
05
return delay
06
07
async def main():
08
tasks = {
09
asyncio.create_task(work(0.1)),
10
asyncio.create_task(work(2)),
11
}
12
13
done, pending = await asyncio.wait(tasks, timeout=0.5)
14
15
print(f"已完成:{len(done)}")
16
print(f"未完成:{len(pending)}")
17
18
for task in pending:
19
task.cancel()
20
21
await asyncio.gather(*pending, return_exceptions=True)
22
23
asyncio.run(main())

wait() 把后续策略交给调用者,所以比 gather() 更灵活,也更容易遗漏清理。使用它之后,必须明确处理 pending 中的每一个 Task。

6. 限制

批量数据有一千条,不等于应该同时发出一千个请求。过高并发可能耗尽连接池、触发模型服务限流、增加超时数量,并占用大量内存。

asyncio.Semaphore 使用一个计数器限制同时进入某段代码的 Task 数量:

semaphore.py
01
import asyncio
02
03
async def fetch(item, limit):
04
async with limit:
05
print(f"开始:{item}")
06
await asyncio.sleep(0.5)
07
print(f"完成:{item}")
08
return item
09
10
async def main():
11
limit = asyncio.Semaphore(2)
12
13
results = await asyncio.gather(
14
*(fetch(item, limit) for item in range(5))
15
)
16
17
print(results)
18
19
asyncio.run(main())

初始值为 2,表示最多两个 Task 同时进入 async with limit。其他 Task 会在获取许可时暂停,离开代码块后许可自动归还。async with 能保证异常和取消发生时也执行释放,通常比手动调用 acquire() 和 release() 更安全。

信号量应该包住真正占用受限资源的部分,不要无意中把无关计算和退避等待也放在里面:

semaphore-scope.py
01
import asyncio
02
03
async def invoke_model(prompt, model, limit):
04
async with limit:
05
response = await model.ainvoke(prompt)
06
07
return response
08
09
async def main(model, prompts):
10
limit = asyncio.Semaphore(3)
11
12
return await asyncio.gather(
13
*(invoke_model(prompt, model, limit) for prompt in prompts)
14
)

这个模式可以限制单个进程同时进行的 LangChain 模型调用数量。

并发与频率

Semaphore 限制的是同时处于受保护区域的数量,不是每秒请求次数。例如,每个请求只用 10 毫秒,即使并发上限是 2,一秒内仍可能发出很多请求。供应商限制为「每分钟 60 次」时,还需要令牌桶等速率限制器,不能只依赖 Semaphore。

并发上限也不应该脱离其他配置单独决定。通常需要一起考虑:

  • HTTP 或数据库连接池大小;
  • 上游接口的并发和频率限制;
  • 单次请求占用的内存;
  • 超时、重试和退避策略;
  • 服务进程与 worker 数量。

还要注意,Semaphore 只限制进入临界区域的数量。上面的 gather() 仍然会为所有输入创建工作;如果输入有数百万条,即使并发上限很小,也可能因 Task 和协程对象过多占用大量内存。持续的大批量工作更适合使用固定数量消费者配合 Queue。

7. 竞态

事件循环通常在一个线程中运行,但单线程不等于没有竞态。只要一段逻辑在读取和写入共享状态之间发生 await,其他 Task 就有机会进入并修改同一状态。

race.py
01
import asyncio
02
03
balance = 100
04
05
async def withdraw(amount):
06
global balance
07
08
if balance < amount:
09
return False
10
11
await asyncio.sleep(0)
12
balance -= amount
13
return True
14
15
async def main():
16
results = await asyncio.gather(
17
withdraw(80),
18
withdraw(80),
19
)
20
21
print(results) # 两次都可能显示成功
22
print(balance) # 可能变成 -60
23
24
asyncio.run(main())

两个 Task 都可能在余额为 100 时通过检查,然后分别扣除 80。问题不是它们同时执行了一条 Python 指令,而是完整业务操作被 await 分成了多个阶段。

同一事件循环中的内存状态可以使用 asyncio.Lock 保护:

lock.py
01
import asyncio
02
03
balance = 100
04
05
async def withdraw(amount, lock):
06
global balance
07
08
async with lock:
09
if balance < amount:
10
return False
11
12
await asyncio.sleep(0)
13
balance -= amount
14
return True
15
16
async def main():
17
lock = asyncio.Lock()
18
19
results = await asyncio.gather(
20
withdraw(80, lock),
21
withdraw(80, lock),
22
)
23
24
print(results)
25
print(balance) # 20
26
27
asyncio.run(main())

锁保护的不是某个变量本身,而是「检查余额并完成扣款」这一整段不可交错的业务操作。

临界区应该尽量短。如果在锁内等待一个很慢的外部接口,其他需要同一把锁的 Task 都会被阻塞。也不要在已经持有 asyncio.Lock 时再次获取同一把锁,它不是可重入锁,这样会让当前 Task 一直等待自己释放锁。

asyncio.Lock 不是线程锁,也不是分布式锁,只能协调同一事件循环中的 Task。FastAPI 部署多个 worker 或多台服务器后,每个进程都有独立的变量和锁。余额、库存、订单状态等持久化数据,应该使用数据库事务、条件更新、唯一约束、行锁或经过验证的分布式协调方案保证一致性。

8. 同步

Lock 和 Semaphore 都属于 asyncio 同步原语。标准库还提供 Event、Condition、BoundedSemaphore 和 Barrier,用于表达不同的任务协作关系:

原语解决的问题
Lock同一时间只允许一个 Task 进入临界区
Event通知多个 Task:某件事已经发生
Condition等待共享状态满足条件,并在检查状态时持有锁
Semaphore同时允许固定数量的 Task 进入
BoundedSemaphore限制并发,并检测释放次数超过初始值的错误
Barrier等到固定数量的 Task 都到达同一阶段后再一起继续

Event

asyncio.Event 内部保存一个布尔标记。标记初始为 False,Task 调用 wait() 时会暂停;其他代码调用 set() 后,所有等待者都会具备继续执行的条件:

event.py
01
import asyncio
02
03
async def worker(name, ready):
04
print(f"{name} 等待配置")
05
await ready.wait()
06
print(f"{name} 开始工作")
07
08
async def main():
09
ready = asyncio.Event()
10
workers = [
11
asyncio.create_task(worker("任务 A", ready)),
12
asyncio.create_task(worker("任务 B", ready)),
13
]
14
15
await asyncio.sleep(0.2)
16
print("配置加载完成")
17
ready.set()
18
19
await asyncio.gather(*workers)
20
21
asyncio.run(main())

set() 之后标记会一直保持为 True,后来的等待者可以直接通过,直到调用 clear() 再次把它设为 False。Event 只表达通知,不携带一批工作;需要逐项传递数据时应该使用 Queue。

Condition

Event 只能表达一个标记,Condition 可以等待更具体的共享状态。它把 Lock 与条件通知组合在一起:等待时暂时释放锁,被唤醒后重新取得锁并检查条件。

condition.py
01
import asyncio
02
03
async def wait_until_ready(state, condition):
04
async with condition:
05
await condition.wait_for(lambda: state["ready"])
06
print("状态已经就绪")
07
08
async def update_state(state, condition):
09
await asyncio.sleep(0.2)
10
11
async with condition:
12
state["ready"] = True
13
condition.notify_all()
14
15
async def main():
16
state = {"ready": False}
17
condition = asyncio.Condition()
18
19
await asyncio.gather(
20
wait_until_ready(state, condition),
21
update_state(state, condition),
22
)
23
24
asyncio.run(main())

通知只表示「状态可能变化了」,等待者仍应重新检查真实状态。condition.wait_for(predicate) 会替调用者重复检查条件,比只调用一次 wait() 更稳妥。调用 notify() 或 notify_all() 时必须持有 Condition 对应的锁。

Barrier 适合把多个 Task 对齐到同一阶段,例如三个并发步骤都完成准备后再进入下一轮。它在普通 Web 接口中不常见,只有业务确实存在固定参与者的阶段同步时才需要使用。

这些 asyncio 同步原语都不是线程安全或进程安全的,而且自身方法不直接接收超时参数。需要限制等待时间时,可以在外层使用 asyncio.timeout() 或 asyncio.wait_for()。

9. 队列

Semaphore 适合限制一批已经存在的工作,Queue 更适合工作持续产生、消费者按固定能力处理的场景。

queue.py
01
import asyncio
02
03
STOP = object()
04
05
async def producer(queue):
06
for item in range(6):
07
await queue.put(item)
08
print(f"生产:{item}")
09
10
async def consumer(name, queue):
11
while True:
12
item = await queue.get()
13
14
try:
15
if item is STOP:
16
return
17
18
print(f"{name} 处理:{item}")
19
await asyncio.sleep(0.3)
20
finally:
21
queue.task_done()
22
23
async def main():
24
queue = asyncio.Queue(maxsize=2)
25
workers = [
26
asyncio.create_task(consumer("消费者 A", queue)),
27
asyncio.create_task(consumer("消费者 B", queue)),
28
]
29
30
await producer(queue)
31
32
for _ in workers:
33
await queue.put(STOP)
34
35
await queue.join()
36
await asyncio.gather(*workers)
37
38
asyncio.run(main())

这个例子有两个消费者,所以同一时间最多处理两个工作。Queue 的 maxsize=2 表示队列最多积压两个尚未取出的项目。队列装满后,生产者会暂停在 await queue.put(item),直到消费者取走项目,这就是反压。

需要区分两个数字:

配置控制内容
消费者 Task 数量同时处理多少项工作
Queue(maxsize=...)最多允许积压多少项等待处理的工作

消费者每次 get() 后都要调用一次 task_done()。它表示这个队列项目已经处理完,并不是在完成某个 asyncio Task。queue.join() 会等待所有放入队列的项目都有对应的 task_done();漏掉一次就可能永远无法返回,多调用一次则会抛出 ValueError。

示例使用特殊的 STOP 对象结束消费者,有几个消费者就要放入几个结束标记。从 Python 3.13 开始,Queue 也提供 shutdown() 和 QueueShutDown,可以停止继续放入项目,并让消费者处理完已有项目后退出。结束标记兼容更早的 Python 版本,也更容易展示完整的数据流。

asyncio.Queue 同样不是线程安全或进程安全的。跨线程使用标准库 queue.Queue,跨进程或跨机器传递可靠任务则需要 Redis、RabbitMQ、Kafka 等外部系统。

10. 选择

不同并发工具解决的是不同问题:

目标工具
按输入顺序得到全部结果gather()
共同成功或失败,并管理子任务生命周期TaskGroup
按完成顺序逐个处理结果as_completed()
等到部分任务完成,并自行处理剩余任务wait()
限制同时访问某项资源的数量Semaphore
保护单事件循环内的共享状态Lock
通知一个或多个 Task 状态已经变化Event、Condition
解耦持续生产与固定数量消费Queue

开始并发之前,还要先确认工作之间是否真的独立。下面两步存在数据依赖,不能并发:

dependent-steps.py
1
memory = await search_memory(question)
2
answer = await generate_answer(question, memory)

而用户画像和长期记忆互不依赖,可以同时查询,再汇总给模型:

independent-steps.py
01
import asyncio
02
03
async def build_context(question, profile_store, memory_store):
04
async with asyncio.TaskGroup() as group:
05
profile_task = group.create_task(
06
profile_store.ainvoke(question)
07
)
08
memory_task = group.create_task(
09
memory_store.ainvoke(question)
10
)
11
12
return {
13
"profile": profile_task.result(),
14
"memory": memory_task.result(),
15
}

真正的优化不是把每个函数都包装成 Task,而是识别调用关系中相互独立的等待,并给并发设置合理边界。

11. 边界

在 FastAPI 中,一个 worker 通常运行一个事件循环。多个请求可以在同一个事件循环中并发推进,但每段同步 Python 代码仍然要在这个线程上依次执行。

假设同时有 100 个请求,每个请求又创建 5 个模型调用 Task,理论上可能迅速积累 500 个上游请求。并发设计必须从整个进程甚至整个服务观察,不能只看单个接口内部是否运行正确。

进程内同步工具也有明确边界:

  • 一个 Semaphore 只限制当前进程,四个 worker 各有一份时,总并发上限可能变成四倍;
  • 一个 Lock 只能保护当前进程的内存,无法保护其他 worker 修改数据库;
  • 一个 Queue 只保存在当前进程内,服务重启后未处理项目会丢失;
  • 上游 API 的全局频率限制,需要所有实例共享的限流状态或网关支持。

遇到同步阻塞 IO 时,可以使用原生异步驱动,或者用 asyncio.to_thread() 避免阻塞事件循环。CPU 密集型 Python 计算通常不会因为创建更多 Task 而加速,应考虑进程池、独立 worker 或专门的计算服务。

还要给并发操作设置超时、取消和错误记录。如果请求已经取消,仍让几十个子任务在后台消耗模型额度,通常不是合理行为;如果确实要让工作脱离 HTTP 请求继续执行,就应把它交给拥有独立生命周期和持久化能力的后台任务系统。

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