1. 并发
并发不是让所有代码在同一时刻运行,而是让多个尚未完成的工作在同一段时间内交替推进。
假设程序要发出三个相互独立的网络请求,每个请求都需要等待一秒。顺序执行时,第二个请求必须等第一个完成后才开始:
01import asyncio02import time0304async def fetch(name):05await asyncio.sleep(1)06return f"{name} 完成"0708async def main():09started_at = time.perf_counter()1011first = await fetch("A")12second = await fetch("B")13third = await fetch("C")1415print(first, second, third)16print(f"耗时:{time.perf_counter() - started_at:.1f} 秒")1718asyncio.run(main())
总耗时接近三秒。三个请求之间没有依赖关系,所以可以先把它们都安排到事件循环,再一起等待:
01import asyncio02import time0304async def fetch(name):05await asyncio.sleep(1)06return f"{name} 完成"0708async def main():09started_at = time.perf_counter()1011results = await asyncio.gather(12fetch("A"),13fetch("B"),14fetch("C"),15)1617print(results)18print(f"耗时:{time.perf_counter() - started_at:.1f} 秒")1920asyncio.run(main())
总耗时接近一秒,因为三个 Task 的等待时间发生了重叠。事件循环仍然一次只推进一个 Task:某个 Task 遇到尚未完成的 await 后暂停,事件循环再推进其他 Task。
需要区分三个容易混淆的概念:
| 概念 | 含义 |
|---|---|
| 顺序 | 一个工作结束后才开始下一个 |
| 并发 | 多个工作在同一段时间内交替推进 |
| 并行 | 多个工作在同一时刻分别占用不同 CPU 核心或执行单元 |
asyncio 主要解决 IO 并发,不会自动让 Python 计算并行。网络请求、数据库查询和模型调用在等待外部响应时可以让出事件循环;一个没有 await 的长计算仍然会占住事件循环线程。
并发也不是越多越快。创建更多任务只能增加正在等待或争抢资源的工作数量,无法突破对方服务限流、数据库连接池、网络带宽和本机 CPU 的上限。
2. gather
asyncio.gather() 接收多个可等待对象,并等待它们全部完成。传入协程对象时,gather() 会自动把它们安排为 Task:
01import asyncio0203async def fetch(name, delay):04await asyncio.sleep(delay)05print(f"{name} 完成")06return name0708async def main():09results = await asyncio.gather(10fetch("A", 1),11fetch("B", 0.2),12fetch("C", 0.5),13)1415print(results) # ['A', 'B', 'C']1617asyncio.run(main())
任务会按照 B、C、A 的顺序完成,但结果仍按照参数 A、B、C 的顺序排列。gather() 适合「提交一组工作,最后按输入顺序取得全部结果」的场景。
失败行为
默认情况下,某个子任务首次抛出异常时,等待 gather() 的调用方会立即收到这个异常。其余子任务不会因此自动取消,通常仍会继续运行:
01import asyncio0203async def fail():04await asyncio.sleep(0.1)05raise ValueError("请求失败")0607async def slow():08await asyncio.sleep(0.3)09print("慢请求仍然完成")1011async def main():12slow_task = asyncio.create_task(slow())1314try:15await asyncio.gather(fail(), slow_task)16except ValueError as error:17print(error)1819await slow_task2021asyncio.run(main())
这点很容易被误解:gather() 把首个异常传播给调用者,不等于替调用者关闭所有其他工作。
如果每项工作可以独立成功或失败,可以设置 return_exceptions=True:
01import asyncio0203async def success():04return "成功"0506async def fail():07raise ValueError("失败")0809async def main():10results = await asyncio.gather(11success(),12fail(),13return_exceptions=True,14)1516for result in results:17if isinstance(result, BaseException):18print(f"异常:{result}")19else:20print(f"结果:{result}")2122asyncio.run(main())
此时异常对象会和正常结果一起出现在列表中。调用者必须主动识别并处理异常,否则可能把失败误当作业务数据。
如果等待 gather() 的调用本身被取消,尚未完成的子任务会收到取消请求。反过来,单个子任务被取消不会直接取消其他子任务。取消行为关系到资源清理,协程仍然应该用 finally 关闭连接和临时资源。
3. TaskGroup
当一组任务共同完成一个业务目标,通常更适合使用 Python 3.11 提供的 TaskGroup:
01import asyncio0203async def fetch(name, delay):04await asyncio.sleep(delay)05return f"{name} 完成"0607async def main():08async with asyncio.TaskGroup() as group:09first = group.create_task(fetch("A", 0.3))10second = group.create_task(fetch("B", 0.1))1112print(first.result())13print(second.result())1415asyncio.run(main())
退出 async with 之前,TaskGroup 会等待所有子任务结束。任务的创建范围同时也是负责等待和清理它们的范围,这种组织方式叫结构化并发。
TaskGroup 与 gather() 最重要的差别在于失败处理:一个子任务异常失败后,TaskGroup 会取消其余尚未完成的子任务,等待它们完成清理,再抛出异常组。
01import asyncio0203async def fail():04await asyncio.sleep(0.1)05raise ValueError("查询失败")0607async def slow():08try:09await asyncio.sleep(10)10finally:11print("慢任务释放资源")1213async def main():14try:15async with asyncio.TaskGroup() as group:16group.create_task(fail())17group.create_task(slow())18except* ValueError as errors:19for error in errors.exceptions:20print(error)2122asyncio.run(main())
可以按下面的原则选择:
| 场景 | 更合适的方式 |
|---|---|
| 每项工作可以独立失败,需要收集全部结果 | gather(return_exceptions=True) |
| 一个任务失败后,其余工作已经没有意义 | TaskGroup |
| 需要最简单地按输入顺序得到全部结果 | gather() |
| 需要清晰管理一组子任务的生命周期 | TaskGroup |
不要为了“看起来更并发”而把有依赖关系的步骤放入同一个 TaskGroup。第二步必须使用第一步的结果时,它们本来就应该顺序执行。
4. 完成
gather() 会等所有任务结束后一次返回。如果希望谁先完成就先处理谁,可以使用 asyncio.as_completed()。
从 Python 3.13 开始,它支持异步迭代;传入 Task 后,异步迭代会按完成顺序产出原来的 Task 对象:
01import asyncio0203async def fetch(name, delay):04await asyncio.sleep(delay)05return f"{name} 的结果"0607async def main():08tasks = [09asyncio.create_task(fetch("慢任务", 1), name="slow"),10asyncio.create_task(fetch("快任务", 0.2), name="fast"),11asyncio.create_task(fetch("中等任务", 0.5), name="medium"),12]1314async for task in asyncio.as_completed(tasks):15result = await task16print(task.get_name(), result)1718asyncio.run(main())
输出顺序通常是 fast、medium、slow。因为取得的是原来的 Task,还能通过名称或映射判断结果属于哪个请求。
为了兼容 Python 3.12 及更早版本,也可以使用普通 for:
01import asyncio0203async def fetch(name, delay):04await asyncio.sleep(delay)05return name0607async def main():08tasks = [09asyncio.create_task(fetch("A", 0.5)),10asyncio.create_task(fetch("B", 0.1)),11]1213for next_result in asyncio.as_completed(tasks):14result = await next_result15print(result)1617asyncio.run(main())
普通迭代产出的是需要等待的新协程,不是原始 Task。两种写法都按完成顺序交付结果,但对象身份不同。
as_completed(tasks, timeout=...) 可以限制等待整组任务的最长时间。超时会抛出 TimeoutError,但不会自动取消剩余任务;调用者需要决定继续等待、取消还是把它们交给其他生命周期管理逻辑。
这种方式适合搜索多个数据源、并发下载、批量模型调用后逐个推送结果等场景。它改变的是结果处理顺序,不会自动限制并发数量。
5. 等待
asyncio.wait() 提供更底层的等待控制。它返回两个集合:已经完成的 done 和仍未完成的 pending。
01import asyncio0203async def fetch(name, delay):04await asyncio.sleep(delay)05return name0607async def main():08tasks = {09asyncio.create_task(fetch("A", 0.2)),10asyncio.create_task(fetch("B", 1)),11}1213done, pending = await asyncio.wait(14tasks,15return_when=asyncio.FIRST_COMPLETED,16)1718for task in done:19print(f"最先完成:{await task}")2021for task in pending:22task.cancel()2324await asyncio.gather(*pending, return_exceptions=True)2526asyncio.run(main())
return_when 支持三种条件:
| 条件 | 返回时机 |
|---|---|
FIRST_COMPLETED | 任意任务完成或取消 |
FIRST_EXCEPTION | 首次出现异常;如果始终没有异常,则等待全部完成 |
ALL_COMPLETED | 全部任务完成或取消,默认值 |
从 Python 3.11 开始,不能把裸协程对象直接传给 wait(),需要先使用 create_task() 创建 Task。
给 wait() 设置 timeout 时,它不会抛出 TimeoutError,也不会自动取消未完成任务,只会把它们放入 pending:
01import asyncio0203async def work(delay):04await asyncio.sleep(delay)05return delay0607async def main():08tasks = {09asyncio.create_task(work(0.1)),10asyncio.create_task(work(2)),11}1213done, pending = await asyncio.wait(tasks, timeout=0.5)1415print(f"已完成:{len(done)}")16print(f"未完成:{len(pending)}")1718for task in pending:19task.cancel()2021await asyncio.gather(*pending, return_exceptions=True)2223asyncio.run(main())
wait() 把后续策略交给调用者,所以比 gather() 更灵活,也更容易遗漏清理。使用它之后,必须明确处理 pending 中的每一个 Task。
6. 限制
批量数据有一千条,不等于应该同时发出一千个请求。过高并发可能耗尽连接池、触发模型服务限流、增加超时数量,并占用大量内存。
asyncio.Semaphore 使用一个计数器限制同时进入某段代码的 Task 数量:
01import asyncio0203async def fetch(item, limit):04async with limit:05print(f"开始:{item}")06await asyncio.sleep(0.5)07print(f"完成:{item}")08return item0910async def main():11limit = asyncio.Semaphore(2)1213results = await asyncio.gather(14*(fetch(item, limit) for item in range(5))15)1617print(results)1819asyncio.run(main())
初始值为 2,表示最多两个 Task 同时进入 async with limit。其他 Task 会在获取许可时暂停,离开代码块后许可自动归还。async with 能保证异常和取消发生时也执行释放,通常比手动调用 acquire() 和 release() 更安全。
信号量应该包住真正占用受限资源的部分,不要无意中把无关计算和退避等待也放在里面:
01import asyncio0203async def invoke_model(prompt, model, limit):04async with limit:05response = await model.ainvoke(prompt)0607return response0809async def main(model, prompts):10limit = asyncio.Semaphore(3)1112return 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 就有机会进入并修改同一状态。
01import asyncio0203balance = 1000405async def withdraw(amount):06global balance0708if balance < amount:09return False1011await asyncio.sleep(0)12balance -= amount13return True1415async def main():16results = await asyncio.gather(17withdraw(80),18withdraw(80),19)2021print(results) # 两次都可能显示成功22print(balance) # 可能变成 -602324asyncio.run(main())
两个 Task 都可能在余额为 100 时通过检查,然后分别扣除 80。问题不是它们同时执行了一条 Python 指令,而是完整业务操作被 await 分成了多个阶段。
同一事件循环中的内存状态可以使用 asyncio.Lock 保护:
01import asyncio0203balance = 1000405async def withdraw(amount, lock):06global balance0708async with lock:09if balance < amount:10return False1112await asyncio.sleep(0)13balance -= amount14return True1516async def main():17lock = asyncio.Lock()1819results = await asyncio.gather(20withdraw(80, lock),21withdraw(80, lock),22)2324print(results)25print(balance) # 202627asyncio.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() 后,所有等待者都会具备继续执行的条件:
01import asyncio0203async def worker(name, ready):04print(f"{name} 等待配置")05await ready.wait()06print(f"{name} 开始工作")0708async def main():09ready = asyncio.Event()10workers = [11asyncio.create_task(worker("任务 A", ready)),12asyncio.create_task(worker("任务 B", ready)),13]1415await asyncio.sleep(0.2)16print("配置加载完成")17ready.set()1819await asyncio.gather(*workers)2021asyncio.run(main())
set() 之后标记会一直保持为 True,后来的等待者可以直接通过,直到调用 clear() 再次把它设为 False。Event 只表达通知,不携带一批工作;需要逐项传递数据时应该使用 Queue。
Condition
Event 只能表达一个标记,Condition 可以等待更具体的共享状态。它把 Lock 与条件通知组合在一起:等待时暂时释放锁,被唤醒后重新取得锁并检查条件。
01import asyncio0203async def wait_until_ready(state, condition):04async with condition:05await condition.wait_for(lambda: state["ready"])06print("状态已经就绪")0708async def update_state(state, condition):09await asyncio.sleep(0.2)1011async with condition:12state["ready"] = True13condition.notify_all()1415async def main():16state = {"ready": False}17condition = asyncio.Condition()1819await asyncio.gather(20wait_until_ready(state, condition),21update_state(state, condition),22)2324asyncio.run(main())
通知只表示「状态可能变化了」,等待者仍应重新检查真实状态。condition.wait_for(predicate) 会替调用者重复检查条件,比只调用一次 wait() 更稳妥。调用 notify() 或 notify_all() 时必须持有 Condition 对应的锁。
Barrier 适合把多个 Task 对齐到同一阶段,例如三个并发步骤都完成准备后再进入下一轮。它在普通 Web 接口中不常见,只有业务确实存在固定参与者的阶段同步时才需要使用。
这些 asyncio 同步原语都不是线程安全或进程安全的,而且自身方法不直接接收超时参数。需要限制等待时间时,可以在外层使用 asyncio.timeout() 或 asyncio.wait_for()。
9. 队列
Semaphore 适合限制一批已经存在的工作,Queue 更适合工作持续产生、消费者按固定能力处理的场景。
01import asyncio0203STOP = object()0405async def producer(queue):06for item in range(6):07await queue.put(item)08print(f"生产:{item}")0910async def consumer(name, queue):11while True:12item = await queue.get()1314try:15if item is STOP:16return1718print(f"{name} 处理:{item}")19await asyncio.sleep(0.3)20finally:21queue.task_done()2223async def main():24queue = asyncio.Queue(maxsize=2)25workers = [26asyncio.create_task(consumer("消费者 A", queue)),27asyncio.create_task(consumer("消费者 B", queue)),28]2930await producer(queue)3132for _ in workers:33await queue.put(STOP)3435await queue.join()36await asyncio.gather(*workers)3738asyncio.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 |
开始并发之前,还要先确认工作之间是否真的独立。下面两步存在数据依赖,不能并发:
1memory = await search_memory(question)2answer = await generate_answer(question, memory)
而用户画像和长期记忆互不依赖,可以同时查询,再汇总给模型:
01import asyncio0203async def build_context(question, profile_store, memory_store):04async with asyncio.TaskGroup() as group:05profile_task = group.create_task(06profile_store.ainvoke(question)07)08memory_task = group.create_task(09memory_store.ainvoke(question)10)1112return {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 请求继续执行,就应把它交给拥有独立生命周期和持久化能力的后台任务系统。