03-asyncio 与并发模型盲区梳理
asyncio 与并发模型盲区梳理
概述
本篇梳理在学习 MiniClaude 项目过程中,围绕 asyncio、事件循环、协程、并发模型暴露的知识盲区。每条盲区包含:①原来的困惑/错误理解 ②正确解释 ③代码示例 ④延伸知识。
读者背景:有 Java 基础,第一次接触 Python asyncio。文中多处借用 Java 概念做类比(事件循环 ≈ 餐厅大堂经理、Future ≈ Java CompletableFuture、Task ≈ Thread 但更轻量)。
盲区清单(速查表)
| # | 盲区 | 关键词 | 出处 |
|---|---|---|---|
| 1 | 事件循环是什么 + asyncio.run() 三步流程 | 事件循环、调度器、run() | S0-Q7 |
| 2 | loop.add_signal_handler() 为何必须用 | 信号处理、线程安全、Ctrl+C | S0-Q6 |
| 3 | 协程优点(可读性、无锁、内存开销小) | 协程、线程、对比 | S0-Q8 |
| 4 | CPU 密集任务阻塞所有协程 + run_in_executor | 阻塞、线程池、让出控制权 | S0-Q9 |
| 5 | asyncio.start_server 自动创建协程 | accept、create_task、封装 | S2-Q9, Q14 |
| 6 | 单进程单线程多协程模型 | asyncio 核心模型、无锁 | S2-Q16 |
| 7 | 事件循环创建时机比主协程早 | asyncio.run、创建顺序 | S2-Q15 |
| 8 | ContextVar 协程局部变量 | ContextVar、threading.local、隐式传参 | S2-Q8, Q13, Q17 |
| 9 | Future + _pending 字典异步请求-响应配对 | Future、JSON-RPC、异步配对 | S2-Q5, Q28, Q32, Q38 |
| 10 | fut.done() 防御性检查 | done()、set_result、防重复 | S2-Q37, Q39, Q40 |
| 11 | create_task vs await 区别 | 串行 vs 并行、孤儿 Task | S2-Q29, S3-Q23 |
| 12 | 协程 vs Task 区别 | 菜谱 vs 厨师做菜 | S3-Q16 |
| 13 | queue.put/get/task_done/join 配对机制 | asyncio.Queue、计数器 | S3-Q9 |
| 14 | asyncio.create_subprocess_shell 启动子进程 | 子进程、超时、沙箱 | S3-Q21 |
| 15 | add_done_callback Task 完成回调 | 回调、自动清理 | S3-Q17 |
| 16 | 多 Run 并发 set[Task] + discard 自动清理 | 并发、set、discard | S3-Q3, Q17 |
| 17 | 进程/线程/协程的父子关系 | 父子关系、层级 | S3-Q23 |
逐条详解
1. 事件循环是什么 + asyncio.run() 三步流程
原来怎么理解的:看到 asyncio.run(CoreApp().run()) 这行,不清楚"事件循环"到底是个什么对象,也不明白 run() 之后程序是否一直重复跑。 正确解释:事件循环(Event Loop)是 asyncio 的核心调度器,像餐厅大堂经理——循环检查哪些协程准备好了(IO 完成、定时器到点),就把控制权交给谁。asyncio.run() 标准三步流程:
- 创建一个事件循环
- 把传入的协程(
CoreApp().run())设为主协程 - 启动调度,直到主协程结束
run() 函数本身只执行一次准备就结束,但事件循环在后台持续运行,调度所有被注册的协程/Task,直到主协程返回。 代码示例:
# 三步流程等价展开
import asyncio
async def main():
print("主协程运行中")
loop = asyncio.new_event_loop() # ① 创建事件循环
asyncio.set_event_loop(loop)
loop.run_until_complete(main()) # ②③ 丢入主协程并启动调度
# asyncio.run(main()) 是上面三步的封装
延伸:事件循环在 Python 中本质是一个 BaseEventLoop 对象,内部维护就绪队列 + 等待队列。Java NIO Selector 是相似概念——一个线程管理多个 Channel。
2. loop.add_signal_handler() 为何必须用
原来怎么理解的:以为信号处理(如 Ctrl+C)用标准库 signal.signal() 就够了,不理解为什么非得用 loop.add_signal_handler()。 正确解释:异步程序中事件循环占着主线程,传统 signal.signal() 注册的处理函数会在信号到达时被主线程中断执行,可能打断事件循环的原子操作,引发数据竞争。loop.add_signal_handler() 把信号处理融入事件循环——信号到达时只是给循环发一个"待办",等当前协程让出控制权后再调用处理函数,保证线程安全。 代码示例:
import asyncio
import signal
async def main():
loop = asyncio.get_running_loop()
# 把 SIGINT (Ctrl+C) 处理融入事件循环
loop.add_signal_handler(signal.SIGINT, lambda: print("收到 Ctrl+C,准备退出"))
await asyncio.Event().wait() # 永远等下去
asyncio.run(main())
延伸:signal.signal() 是同步信号处理(UNIX 信号机制),处理函数只能执行少量 async-safe 操作;add_signal_handler 内部用 loop.call_later 调度,可以在处理函数里安全地 await。
3. 协程优点(可读性、无锁、内存开销小)
原来怎么理解的:只知道"协程比线程轻",但说不清具体优势在哪。 正确解释:协程相对线程/进程有三大优势:
- 代码像同步:用
await写出顺序结构的代码,可读性高,没有回调地狱。 - 无锁编程:单线程内多协程通过
await主动让出控制权,不存在抢占,共享变量不需要锁。 - 内存开销小:协程只几十字节到几 KB,线程通常要几 MB 栈空间,能开数量差 3 个量级。 代码示例:
import asyncio
# 顺序结构,可读性高
async def fetch_and_process():
data = await fetch_data() # 像同步代码
result = await process(data)
return result
# 无锁共享 counter(单线程,无竞争)
counter = 0
async def increment():
global counter
for _ in range(1000):
counter += 1 # 不需要 threading.Lock
延伸:协程不是银弹——CPU 密集任务、调用阻塞 IO 库(如同步 requests)会让整个事件循环卡死。Java Project Loom 的虚拟线程是相似思路。
4. CPU 密集任务阻塞所有协程 + run_in_executor 解决
原来怎么理解的:以为"协程是高并发银弹",没意识到单线程内的 CPU 密集代码会让所有协程一起堵死。 正确解释:协程通过 await 让出控制权,但纯 CPU 计算不调 await 就不会让出——它会霸占整个线程,事件循环转不起来,其他协程全部饿死。解决方法:把 CPU 密集任务丢进 run_in_executor() 让线程池执行,事件循环继续转。 代码示例:
import asyncio
import time
# 反例:CPU 密集不 await,所有协程卡死
def heavy_compute(n):
s = 0
for i in range(n):
s += i
return s
async def bad():
# 直接调用,阻塞整个事件循环
result = heavy_compute(10**8)
return result
async def good():
# 正解:丢进线程池
loop = asyncio.get_running_loop()
result = await loop.run_in_executor(None, heavy_compute, 10**8)
return result
延伸:run_in_executor(None, fn, *args) 的 None 表示用默认 ThreadPoolExecutor;CPU 密集真正想要多核并行要用 ProcessPoolExecutor。Java 中对应 ExecutorService.submit(Callable)。
5. asyncio.start_server 自动创建协程
原来怎么理解的:以为每接受一个连接就要手写 while True: accept() 循环 + 手动 create_task,疑惑框架到底帮做了什么。 正确解释:asyncio.start_server 内部实现:每 accept 一个 TCP 连接,就 create_task 一个新协程来跑 _handle_connection 回调。这是 asyncio.start_server 的基础行为——程序员只需要提供 client_connected_cb 回调函数,框架自动处理并发,根本不需要写 while True: accept() 循环。 代码示例:
import asyncio
async def handle_connection(reader, writer):
data = await reader.read(100)
writer.write(data)
await writer.drain()
writer.close()
# 框架自动: 每来一个连接 → create_task(handle_connection(reader, writer))
asyncio.run(asyncio.start_server(handle_connection, "127.0.0.1", 8888))
延伸:Java NIO 的 ServerSocketChannel.accept() 需要手写 Selector 循环;Netty 的 ChannelInboundHandler 才是类似的"框架帮你并发"模式。asyncio.start_server 的并发能力来自事件循环,不是线程池。
6. 单进程单线程多协程模型
原来怎么理解的:以为 MiniClaude daemon 跑那么多协程,背后肯定有线程池或者多个工作线程。 正确解释:这就是 asyncio 的核心模型:单进程、单线程、事件循环 + 多协程。没有 ThreadPool,没有额外线程。多个协程通过 await 主动让出控制权,事件循环轮流调度它们。所有协程都在同一个线程里,不存在 CPU 并发竞争。 代码示例:
import asyncio
async def task_a():
await asyncio.sleep(1) # 让出控制权
print("A done")
async def task_b():
await asyncio.sleep(1)
print("B done")
async def main():
# 两个协程跑在同一个线程,并发不并行
await asyncio.gather(task_a(), task_b())
asyncio.run(main())
# 输出(约 1 秒后同时):
# A done
# B done
延伸:"并发(concurrent)" 不等于 "并行(parallel)"。并发是同时处理多件事(一个线程切换),并行是同时执行多件事(多核多线程)。asyncio 是并发不是并行。
7. 事件循环创建时机比主协程早
原来怎么理解的:以为 asyncio.run(CoreApp().run()) 是"主协程先跑,事件循环跟着启动"。 正确解释:调用链:run() → asyncio.run(CoreApp().run())。asyncio.run() 内部先创建事件循环,再把它设为当前线程的事件循环,然后才把 CoreApp().run() 作为主协程丢进去跑。所以事件循环的创建比主协程早一步。 代码示例:
import asyncio
async def main():
# 主协程里能直接拿当前事件循环,因为循环已经存在
loop = asyncio.get_running_loop()
print(f"事件循环: {loop}")
# asyncio.run 内部顺序:
# 1. loop = asyncio.new_event_loop() ← 先创建
# 2. asyncio.set_event_loop(loop) ← 设置为当前
# 3. loop.run_until_complete(main()) ← 最后才跑主协程
asyncio.run(main())
延伸:Python 3.10+ 推荐用 asyncio.run() 而不是手动 get_event_loop()。手动获取循环在协程外会报警告,因为 Python 在逐步淘汰"隐式事件循环"概念。
8. ContextVar 协程局部变量
原来怎么理解的:第一次见 ContextVar,不知道是什么,也不知道为什么不用函数参数显式传递。 正确解释:ContextVar 是 Python 3.7+ 的协程局部变量。在 asyncio 中,多个协程共享同一个线程,没法用 threading.local(它按线程隔离)。ContextVar 让每个协程有自己的"独立副本"。MiniClaude 中 _writer_var 用来隐式传递"当前请求对应的客户端 writer",这样 handler 不需要显式接收 writer 参数,调 get_connection_writer() 就能拿到。 代码示例:
import asyncio
from contextvars import ContextVar
# 1. 定义
_writer_var: ContextVar[str] = ContextVar("connection_writer")
async def handle_client(writer_id: str):
# 2. 设置(在连接协程中)
_writer_var.set(writer_id)
await process_request()
async def process_request():
# 3. 读取(在 handler 中)
writer = _writer_var.get()
print(f"当前客户端: {writer}")
async def main():
await asyncio.gather(
handle_client("client-A"),
handle_client("client-B"),
)
asyncio.run(main())
延伸:ContextVar 内部维护"协程上下文 → 值"的映射,每次 set() 只影响当前协程及其子协程,get() 只返回当前协程的值。Java 的 ThreadLocal 是按线程隔离,对应 Python 的 threading.local();ContextVar 是 asyncio 时代的协程级 ThreadLocal。
9. Future + _pending 字典异步请求-响应配对
原来怎么理解的:不清楚 _pending 字典为什么是 dict[str, asyncio.Future],也不知道它解决了什么核心问题。 正确解释:SocketClient._pending 用于解决 "异步请求-响应时间差" 的问题。客户端 send_command 发 JSON-RPC 请求后,不能阻塞整个客户端等待这一条响应(因为事件推送也在同一 TCP 连接上),所以用 Future 让 send_command 异步等待自己那条响应:
send_command创建一个 Future,存入_pending[req.id]await fut阻塞当前协程- 服务端返回响应时,
_dispatch根据msg["id"]找到对应的 Future fut.set_result(msg.get("result"))把结果放进去send_command被唤醒拿到结果
代码示例:
import asyncio
import uuid
class SocketClient:
def __init__(self, host: str, port: int):
self._host = host
self._port = port
self._reader = None
self._writer = None
self._pending: dict[str, asyncio.Future] = {}
async def send_command(self, method: str, params: dict):
req_id = str(uuid.uuid4()) # 用 UUID 而非自增 id
request = JsonRpcRequest(id=req_id, method=method, params=params)
fut = asyncio.get_running_loop().create_future()
self._pending[req_id] = fut # ① 存入
self._writer.write(request.model_dump_json().encode() + b"\n")
await self._writer.drain()
return await fut # ② 阻塞等待
async def run_event_loop(self):
while line := await self._reader.readline():
if not line:
break
await self._dispatch(line)
async def _dispatch(self, line: bytes):
msg = json.loads(line)
if "jsonrpc" in msg and msg.get("id") in self._pending:
fut = self._pending.pop(msg["id"])
if not fut.done():
fut.set_result(msg.get("result") or {}) # ③ 唤醒 send_command
延伸:Java CompletableFuture 是几乎一样的抽象——一个"未来的值",可被异步完成。Java NIO 的 SelectionKey attachment 也是相同思路:把请求上下文挂在通道上,等响应回来时取出来配对。
10. fut.done() 防御性检查
原来怎么理解的:看到 if not fut.done(): fut.set_result(...) 觉得多此一举,Future 不是只会被完成一次吗? 正确解释:同一个请求的处理逻辑可能被多次触发(比如异常情况下重复读、超时和正常响应同时到达)。fut.done() 检查 Future 是否已完成(有结果或异常)。set_result 对已完成的 Future 会抛 InvalidStateError。配合 not fut.done() 就是"还没完成才设置结果",防止重复 set_result 抛异常——属于防御性编程。 代码示例:
import asyncio
async def main():
fut = asyncio.get_running_loop().create_future()
# 第一次设置:成功
if not fut.done():
fut.set_result("第一次结果")
print(fut.result()) # 第一次结果
# 第二次设置:done() 防御
if not fut.done(): # 已完成,跳过
fut.set_result("第二次结果") # 不会执行,避免 InvalidStateError
asyncio.run(main())
延伸:_pending.pop(req_id, None) 也是同样的防御思路——用 pop 而不是 del,key 不存在(被重复处理过)也不报错。防御性编程在异步代码里特别重要,因为执行顺序不再线性。
11. create_task vs await 区别
原来怎么理解的:把 create_task 和 await 当成"两种调用协程的方式",分不清什么时候用哪个。 正确解释:
await coro:是子调用,串行执行。当前协程暂停,等被 await 的协程跑完才继续。asyncio.create_task(coro):把协程包装成 Task 丢到事件循环后台运行,立即返回 Task 对象,当前协程不等待,并行执行。
MiniClaude 中 agent.run handler 用 create_task 把 self._sessions.send_message(...) 丢后台(agent 可能跑几百秒,不能阻塞客户端的 send_command),而 ping/subscribe 这种毫秒级操作直接 await。注意 runner 在 SessionManager 内部创建,handler 只调 send_message。 代码示例:
import asyncio
async def slow_task(name, seconds):
print(f"{name} 开始")
await asyncio.sleep(seconds)
print(f"{name} 结束")
async def main():
# ① await:串行,总耗时 3 秒
await slow_task("A", 1)
await slow_task("B", 2)
# ② create_task:并行,总耗时 2 秒
t1 = asyncio.create_task(slow_task("C", 1))
t2 = asyncio.create_task(slow_task("D", 2))
await t1
await t2
asyncio.run(main())
延伸:create_task 创建的 Task 如果创建者不持有引用,会变成"孤儿 Task"——还在跑但没人管。Java ExecutorService.submit(Runnable) 返回的 Future 如果不持有,也是同样的"发射后不管"。
12. 协程 vs Task 区别
原来怎么理解的:把"协程"和"Task"混为一谈,看到文档里两个词交替用很晕。 正确解释:
- 协程(Coroutine):
async def定义的函数,调用返回协程对象——这是"菜谱",还没开始做。 - Task:被事件循环调度的协程,用
create_task()创建——这是"厨师正在按菜谱做菜"。
协程对象本身不会跑,必须 await 它或者 create_task 包装成 Task 才会被事件循环调度。 代码示例:
import asyncio
async def my_coro():
await asyncio.sleep(1)
return "done"
async def main():
coro = my_coro() # 协程对象,还没跑
print(type(coro)) # <class 'coroutine'>
task = asyncio.create_task(coro) # 包装成 Task,开始调度
print(type(task)) # <class 'asyncio.Task'>
result = await task # 等它完成
print(result) # done
asyncio.run(main())
延伸:Java 没有直接对应——Runnable/Callable 是接口对象(≈ 协程对象),Thread 或 Future 才是被调度的实体(≈ Task)。Python 协程是语言级特性(基于生成器),Task 是 asyncio 框架的调度单元。
13. queue.put/get/task_done/join 配对机制
原来怎么理解的:看到 TraceWriter 用 queue.put + 后台 task_done + queue.join,不清楚这一套配对到底在干什么。 正确解释:asyncio.Queue 的内置方法,配对使用机制:
| 方法 | 作用 | 内部计数 |
|---|---|---|
queue.put(item) | 入队 | +1 |
queue.get() | 出队(取走但未处理完) | 不变 |
queue.task_done() | 告诉队列"处理完了" | -1 |
queue.join() | 阻塞直到计数归零 | 等待 |
类比:餐厅里服务员贴订单(put),厨师取订单(get),做完划掉(task_done),经理等所有订单划掉才关门(join)。
MiniClaude TraceWriter 的设计:emit() 不直接写文件,扔进队列立即返回(不阻塞事件循环);后台 _drain() task 慢慢从队列取数据写文件;stop() 时 await queue.join() 等队列清空,再取消 drain task,避免丢数据。 代码示例:
import asyncio
async def producer(q):
for i in range(3):
await q.put(f"item-{i}") # 入队
print(f"生产 item-{i}")
async def consumer(q):
while True:
item = await q.get() # 出队
await asyncio.sleep(0.1) # 模拟处理
q.task_done() # 标记完成
print(f"消费 {item}")
async def main():
q = asyncio.Queue()
cons = asyncio.create_task(consumer(q))
await producer(q)
await q.join() # 等队列清空
cons.cancel() # 关闭消费者
asyncio.run(main())
延伸:Java BlockingQueue 没有 task_done 配对——它用 take()/put() 阻塞语义直接控制流量。Python 的 join/task_done 设计更接近 Go 的 WaitGroup:Add(1) ≈ put,Done() ≈ task_done,Wait() ≈ join。
14. asyncio.create_subprocess_shell 启动子进程
原来怎么理解的:以为 BashTool 执行 shell 命令是 subprocess.run() 同步调用,不清楚它和 asyncio 怎么配合。 正确解释:MiniClaude BashTool 用 asyncio.create_subprocess_shell() 启动子进程——这是 asyncio 版的 subprocess.Popen,异步等待子进程,不会阻塞事件循环。配合 asyncio.wait_for 实现超时保护,超时则 kill 子进程。 代码示例:
import asyncio
async def run_shell(command: str, timeout: int = 60):
# 启动子进程(异步,不阻塞事件循环)
proc = await asyncio.create_subprocess_shell(
command,
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.STDOUT, # 合并错误流和输出流
)
try:
stdout, _ = await asyncio.wait_for(proc.communicate(), timeout=timeout)
return stdout.decode()
except asyncio.TimeoutError:
proc.kill()
await proc.wait()
return f"超时(>{timeout}s)"
asyncio.run(run_shell("ls -la"))
延伸:stderr=STDOUT 合并错误流和输出流,避免分别读两个流时的死锁问题(管道满了一边另一边写不进去)。没有真正的沙箱——BashTool 直接在宿主机执行命令,没有容器隔离、没有权限限制,这是 MiniClaude 作为课程项目的简化处理。
15. add_done_callback Task 完成回调
原来怎么理解的:看到 run_task.add_done_callback(self._running_runs.discard) 这行,不懂"回调"是什么时候被谁调用的。 正确解释:add_done_callback 是 asyncio 的规定:Task 完成后自动调用 回调函数(task自身)。这是事件循环在 Task 终止时主动触发的,不需要开发者轮询 done()。
逐层拆解:
self._running_runs— 一个set[asyncio.Task],存着正在跑的 Agent Run Taskself._running_runs.discard— set 的 discard 方法(安全删除,元素不存在也不报错)add_done_callback— Task 完成后自动调用
为什么用 discard 不用 remove:边界情况——关闭时主协程可能已经手动取消了 task,如果再用 remove 会抛 KeyError,discard 则静默忽略。 代码示例:
import asyncio
async def my_task():
await asyncio.sleep(1)
return "done"
running = set()
async def main():
t = asyncio.create_task(my_task())
running.add(t)
t.add_done_callback(running.discard) # 完成后自动从 set 删除
print(f"运行中: {len(running)}") # 1
await t
print(f"运行中: {len(running)}") # 0(被回调自动清理)
asyncio.run(main())
延伸:Java Guava 的 ListenableFuture.addListener(Runnable, Executor) 是直接对应物。回调在事件循环线程同步执行,所以回调函数内不能写阻塞代码——如果要 await 别的协程,得在回调里再 create_task。
16. 多 Run 并发 set[Task] + discard 自动清理
原来怎么理解的:以为 MiniClaude 像 S2 那样一次只能跑一个 agent run,不清楚 S3 怎么改成并发的。 正确解释:S2 用单变量 _current_run_task + RuntimeError("a run is already in progress") 互斥;S3 改成 set[asyncio.Task[Any]] 集合支持多 Run 并发——每个 agent.run 命令创建一个 Task 加入集合,handler 立即返回 run_id,agent 在后台跑,进度通过 EventBus 推送。Task 完成后 add_done_callback(discard) 自动从集合移除,防止内存泄漏和 set 无限膨胀。 代码示例:
import asyncio
from typing import Any
class CoreApp:
def __init__(self):
# S3: 集合替代单变量
self._running_runs: set[asyncio.Task[Any]] = set()
self._sessions: SessionManager | None = None
async def _agent_run_handler(self, params: dict[str, Any]) -> AgentRunResult:
assert self._sessions is not None
cmd = AgentRunCommand.model_validate(params)
session = await self._sessions.create(mode="one_shot", title=cmd.goal[:40])
run_id = new_run_id()
# runner 在 SessionManager 内部创建,handler 只调 send_message
run_task = asyncio.create_task(
self._sessions.send_message(session.id, cmd.goal, run_id=run_id)
)
self._running_runs.add(run_task) # ← 加入集合
run_task.add_done_callback(self._running_runs.discard) # ← 完成后自动移除
return AgentRunResult(run_id=run_id) # ← 立即返回,不拒绝
async def shutdown(self):
for t in list(self._running_runs):
t.cancel()
if self._running_runs:
await asyncio.gather(*self._running_runs, return_exceptions=True)
延伸:S3 引入 Task 系统,用户可能同时创建多个任务,每个任务触发独立的 Agent Run,所以必须支持并发。但并发的只是 Run(Agent 执行实例),不是 LLM 调用——每个 Run 有自己的 AgentLoop、ExecutionContext、TaskManager,互不干扰,共享的只有 EventBus 和 TraceWriter(队列是协程安全的)。
17. 进程/线程/协程的父子关系
原来怎么理解的:以为像进程有父子一样,"父协程"可以管理"子协程","父线程"可以 join "子线程"。 正确解释:三种并发单元的父子关系并不对等:
| 概念 | 有没有"父子"关系 | 原因 |
|---|---|---|
| 进程 | 有 | 操作系统维护,父进程可以 kill/等待子进程 |
| 线程 | 没有 | 全部平级,没有"父线程"概念 |
| 协程/Task | 逻辑上有,机制上没有 | 创建者可以 cancel Task,但 asyncio 不维护层级 |
- 进程:父进程创建子进程,父死子受影响(变孤儿)。
- 线程:所有线程平级,共享内存,任何线程可 join 任何线程。
- Task:创建者持有引用可以 cancel,但如果创建者不持有引用,Task 就是"孤儿 Task"——还在跑,但没人管。
await是子调用(串行),create_task才是创建并行 Task。
代码示例:
import asyncio
async def child():
try:
await asyncio.sleep(10)
except asyncio.CancelledError:
print("子 Task 被取消")
raise
async def main():
# 持有引用:可管理(取消/等待)
t = asyncio.create_task(child())
await asyncio.sleep(1)
t.cancel() # 父子关系:逻辑上"父"取消"子"
await t
# 不持有引用:孤儿 Task,没人管
asyncio.create_task(child()) # 跑着,但 cancel 不了
# main 结束后事件循环关闭,孤儿 Task 被强制取消
asyncio.run(main())
延伸:Python 3.11+ 引入 asyncio.TaskGroup,提供结构化并发(structured concurrency)——子 Task 必须在父作用域结束前完成,类似 Go 的 errgroup 或 Rust 的 tokio::task::JoinSet。Java 21 虚拟线程的 StructuredTaskScope 也是同一思路。
复习自检
- 能用一句话解释事件循环是什么,并说出
asyncio.run()三步流程 - 知道为什么信号处理必须用
loop.add_signal_handler()而不是signal.signal() - 能列举协程相对线程的 3 个优势
- 能解释为什么 CPU 密集任务会阻塞所有协程,并用
run_in_executor解决 - 知道
asyncio.start_server自动帮程序员做了什么 - 能说出 asyncio 的核心并发模型(单进程/单线程/多协程)
- 知道事件循环和主协程的创建先后顺序
- 能用
ContextVar实现协程局部变量并解释为什么不用threading.local - 能画出
_pending字典 + Future 的请求-响应配对时序图 - 知道
fut.done()防御性检查防的是什么异常 - 能区分
await coro和create_task(coro)的执行语义 - 能用"菜谱 vs 厨师做菜"类比解释协程和 Task 的区别
- 能说清楚
queue.put/get/task_done/join四个方法的作用和计数变化 - 知道
asyncio.create_subprocess_shell与subprocess.run的区别 - 能解释
add_done_callback的触发时机和回调函数的入参 - 能用
set[Task]+discard写一个支持并发的 task 管理器 - 能列出进程/线程/协程三者父子关系的差异
易错点总结
- 事件循环不是循环跑主协程:它持续调度所有注册的协程/Task,主协程只是其中一个。
- 协程 ≠ Task:协程是函数对象(菜谱),Task 是被事件循环调度的实体(厨师正在做)。
await≠create_task:await串行子调用,create_task创建并行 Task。- 单线程 ≠ 不能并发:asyncio 通过
await让出控制权实现并发,但不是并行(不能利用多核)。 - CPU 密集 ≠ 协程问题:CPU 密集不调 await 不让出,会卡死整个事件循环,必须丢线程池/进程池。
ContextVar≠threading.local:前者按协程隔离,后者按线程隔离,asyncio 必须用前者。- Future 完成 ≠ 一次性事件:必须用
not fut.done()防御,避免重复set_result抛InvalidStateError。 task_done()不是必须:但如果不调,queue.join()会永远等不到计数归零。add_done_callback回调里不能 await:回调是同步执行的,要异步操作得在回调里再create_task。- Task 没有真正的"父":
create_task创建者持有引用可 cancel,不持有就是孤儿 Task。 asyncio.run()创建的事件循环在主协程之前:所以主协程内可以get_running_loop()直接拿到。asyncio.start_server不需要手写 accept 循环:框架每接受一个连接自动create_task。