异步代码里怎么执行外部命令?asyncio 的子进程怎么用?
简化版
在协程里绝对不能直接用 subprocess.run()——它是同步阻塞的,会把整个事件循环卡住(外部命令跑 5 秒,你的服务就有 5 秒不处理任何请求)。asyncio 提供了两个对应的 API:asyncio.create_subprocess_exec(*args)(推荐,参数是列表,不经过 shell,没有注入风险)和 asyncio.create_subprocess_shell(cmd)(整条命令字符串交给 shell 执行,支持管道和重定向,但有命令注入风险)。它们返回一个 Process 对象,常用方法是 await proc.communicate(input=None)(一次性写入输入、读完全部输出、等待退出——最安全的用法)和 await proc.wait()(只等退出)。最经典的陷阱是管道死锁:如果你用 proc.stdout.read() 读输出却不读 stderr,一旦子进程往 stderr 写满了管道缓冲区(Linux 上通常 64KB),子进程会阻塞在写操作上,而你在等它的 stdout——双方互相等待,永久死锁。communicate() 内部会同时读两个管道,所以它是安全的;要流式处理输出时,必须为 stdout 和 stderr 各创建一个任务并发读取。超时处理要用 asyncio.wait_for(proc.communicate(), timeout=N),超时后必须 proc.kill() 并再 await proc.wait()——否则子进程会变成僵尸。另外 Windows 上必须使用 ProactorEventLoop(3.8 起是默认),SelectorEventLoop 不支持子进程。核心记忆:协程里用 create_subprocess_exec(不是 subprocess.run);优先用 communicate() 避免管道死锁;超时后要 kill + wait。
详细版
两个 API 与关键方法:
| API/方法 | 说明 |
|---|---|
create_subprocess_exec(*args) | 推荐,参数列表,无 shell 注入风险 |
create_subprocess_shell(cmd) | 走 shell,支持管道重定向,注意注入 |
await proc.communicate(input) | 一次性读完输出并等待退出(最安全) |
await proc.wait() | 只等退出(不读管道 → 可能死锁) |
proc.stdout / stderr | StreamReader,可流式读 |
proc.stdin | StreamWriter,写完要 drain() + close() |
proc.send_signal/terminate/kill | 发信号(之后仍要 wait()) |
proc.returncode / proc.pid | 退出码(负数=被信号杀死)/ 进程号 |
import asyncio, sys
# ① ★基本用法:exec(推荐)★
async def run(*args, timeout=30):
proc = await asyncio.create_subprocess_exec(
*args,
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE,
)
try:
out, err = await asyncio.wait_for(proc.communicate(), timeout)
except asyncio.TimeoutError:
proc.kill() # ★先杀★
await proc.wait() # ★★必须再 wait,否则变僵尸★★
raise
if proc.returncode != 0:
raise RuntimeError(f"命令失败 rc={proc.returncode}: {err.decode()}")
return out.decode()
data = await run("ls", "-l", "/tmp")
# ② shell 版本(★注意注入★)
proc = await asyncio.create_subprocess_shell(
"cat data.txt | grep ERROR | wc -l", # ★管道、重定向要用 shell★
stdout=asyncio.subprocess.PIPE)
# ✗ await asyncio.create_subprocess_shell(f"ls {user_input}") ★命令注入!★
# ✓ await asyncio.create_subprocess_exec("ls", user_input) ★参数不经 shell★
# ③ ★流式读取输出(长时间运行的命令)★
async def stream_output(*args):
proc = await asyncio.create_subprocess_exec(
*args, stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.STDOUT) # ★合并到 stdout,避免双管道死锁★
async for line in proc.stdout: # ★StreamReader 支持异步迭代★
print(line.decode().rstrip())
await proc.wait()
# ④ ★同时读 stdout 和 stderr(★不合并时必须并发读★)★
async def run_both(*args):
proc = await asyncio.create_subprocess_exec(
*args, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE)
async def drain(stream, tag):
async for line in stream:
logging.info("[%s] %s", tag, line.decode().rstrip())
async with asyncio.TaskGroup() as tg: # ★两个管道并发读★
tg.create_task(drain(proc.stdout, "out"))
tg.create_task(drain(proc.stderr, "err"))
return await proc.wait()
# ⑤ 向子进程写数据
proc = await asyncio.create_subprocess_exec(
"grep", "ERROR",
stdin=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE)
out, _ = await proc.communicate(input=b"line1\nERROR x\n") # ★input 是 bytes★
# ⑥ ★并发运行多个命令 + 限流★
async def run_many(cmds, concurrency=5):
sem = asyncio.Semaphore(concurrency) # ★限制同时运行的进程数★
async def one(cmd):
async with sem:
return await run(*cmd)
return await asyncio.gather(*(one(c) for c in cmds), return_exceptions=True)
# ⑦ ★Windows 注意★
# 3.8+ 默认就是 ProactorEventLoop,支持子进程
# 如果显式设了 SelectorEventLoop → ★NotImplementedError★
⚠️ 三个必须记住的点:①
subprocess.run()/Popen.wait()在协程里是灾难——它们是同步阻塞调用,会让整个事件循环停摆:其他请求不处理、超时不触发、心跳发不出去。协程里必须用asyncio.create_subprocess_*;如果非要用同步的subprocess(比如某个库封装了它),就await asyncio.to_thread(...)丢到线程里。② 管道死锁是最经典的陷阱:stdout=PIPE, stderr=PIPE之后只读其中一个,另一个的管道缓冲区(Linux 上通常 64KB)写满时子进程会阻塞在 write 上,而你在等它的输出——双方互相等待,永久卡住。三种正确做法:用communicate()(内部同时读两个管道)、把 stderr 合并进 stdout(stderr=asyncio.subprocess.STDOUT)、或者为两个管道各起一个任务并发读。③ 超时之后必须「kill + wait」两步:proc.kill()只是发送SIGKILL,子进程死后仍处于僵尸状态直到你await proc.wait()取走它的退出状态——只 kill 不 wait 会在长期运行的服务里累积僵尸进程,最终耗尽 PID。
完整版教学
一、为什么不能用 subprocess.run
★ 同步 subprocess 在协程里的后果:
async def handler():
result = subprocess.run(["ffmpeg", ...], capture_output=True) # ★5 秒★
return result.stdout
→ 这 5 秒里:
✗ 事件循环★完全停摆★
✗ 其他请求的协程一个都不执行
✗ ★asyncio.sleep 到期的任务不会被唤醒★(超时判断全部失真)
✗ 心跳/keepalive 发不出去 → ★连接被对端断开★
✗ 健康检查超时 → ★可能被编排系统重启★
★ 对比三种做法:
┌────────────────────────────────┬────────────────────────────┐
│ 写法 │ 后果 │
├────────────────────────────────┼────────────────────────────┤
│ subprocess.run(...) │ ★阻塞整个事件循环★ │
│ await to_thread(subprocess.run) │ 可用(★占一个线程池槽位★) │
│ ★await create_subprocess_exec★ │ ★原生异步,零线程开销★ │
└────────────────────────────────┴────────────────────────────┘
★ asyncio 子进程的原理(★理解它才知道边界★):
Unix:
- 创建子进程后,父进程需要知道"它什么时候退出"
- asyncio 通过 ★SIGCHLD 信号 + child watcher★ 来感知
(3.8+ 默认 ThreadedChildWatcher:★为每个子进程起一个线程等它★)
- 管道的读写通过 ★事件循环的 selector★ 注册为可读/可写事件
→ 所以 "等子进程" 和 "读管道" 都不会阻塞循环
Windows:
- ★必须用 ProactorEventLoop★(基于 IOCP)
- SelectorEventLoop ★不支持子进程★ → NotImplementedError
- 3.8 起 Windows 默认就是 Proactor,一般不用管
★ 什么时候仍然要用 to_thread(subprocess.run):
✓ 第三方库内部封装了同步 subprocess,你改不了
✓ 需要 subprocess 的某些高级参数(asyncio 版参数较少)
✓ ★Windows 上遇到 Proactor 的兼容问题★
★ 代价:占用线程池槽位(默认只有 min(32, cpu+4) 个)
★ 关键区别:asyncio 子进程"异步"在哪
★不是让命令本身跑得快★,而是"等待它"的过程不阻塞事件循环
→ 100 个 ffmpeg 并发:CPU 还是那么多,但你的服务仍能正常响应其他请求
→ ★所以还需要 Semaphore 限制同时运行的进程数★(CPU/内存是有限的)
协程里用 subprocess.run() 是灾难:那几秒内事件循环完全停摆——其他请求不处理、asyncio.sleep 到期的任务不被唤醒(超时判断全部失真)、心跳发不出去导致连接被断、健康检查超时甚至可能被编排系统重启。正确做法是 await asyncio.create_subprocess_exec(...)(原生异步、零线程开销),实在不行才 await asyncio.to_thread(subprocess.run, ...)(会占一个线程池槽位)。理解原理有助于知道边界:Unix 上 asyncio 通过 SIGCHLD 和 child watcher 感知子进程退出(3.8+ 默认的 ThreadedChildWatcher 会为每个子进程起一个线程等它),管道读写则注册到事件循环的 selector 上;Windows 必须用 ProactorEventLoop(SelectorEventLoop 不支持子进程,会抛 NotImplementedError,3.8 起默认已是 Proactor)。最后要澄清一个概念:异步子进程「异步」的是「等待它的过程」而不是「命令本身跑得更快」——CPU 就那么多,所以仍然需要 Semaphore 限制同时运行的进程数。
二、exec vs shell:安全与能力的取舍
★ create_subprocess_exec(*args)★ —— 首选
proc = await asyncio.create_subprocess_exec("ls", "-l", path)
→ 参数是★独立的列表项★,直接传给 execve,★不经过 shell★
→ ★没有命令注入风险★(用户输入里的 ; | $ 都只是普通字符)
→ 不支持:管道 |、重定向 >、通配符 *、环境变量展开 $VAR
★ create_subprocess_shell(cmd)★ —— 谨慎使用
proc = await asyncio.create_subprocess_shell("cat a.txt | grep x > out.txt")
→ 整条字符串交给 /bin/sh -c 执行
→ 支持 shell 的所有特性
→ ★命令注入风险★:
user = "a.txt; rm -rf /"
await create_subprocess_shell(f"cat {user}") # ★灾难★
✓ 如果必须用 shell + 用户输入 → ★shlex.quote() 转义★
import shlex
cmd = f"cat {shlex.quote(user)}"
✓ 更好的做法:★用 exec + 在 Python 里实现管道逻辑★
★ 什么时候真的需要 shell:
✗ "ls -l" → exec 就行(拆成两个参数)
✗ "ls *.txt" → 用 glob 模块自己展开
✗ "cmd1 | cmd2" → ★用两个 exec + 手动接管道★(见下)
✓ 复杂的一次性脚本、需要 shell 内建命令(cd、export、source)
✓ 用户自己写的完整命令行(如运维工具,★但要考虑权限边界★)
★ 用 exec 实现管道(不用 shell):
p1 = await asyncio.create_subprocess_exec(
"cat", "big.log", stdout=asyncio.subprocess.PIPE)
p2 = await asyncio.create_subprocess_exec(
"grep", "ERROR", stdin=p1.stdout, stdout=asyncio.subprocess.PIPE)
# ★把 p1 的 stdout 直接作为 p2 的 stdin★
out, _ = await p2.communicate()
await p1.wait()
★ 其他常用参数:
cwd="/path" 工作目录
env={...} ★环境变量(会完全替换,不是追加)★
→ 想追加:env={**os.environ, "KEY": "v"}
stdout/stderr/stdin:
asyncio.subprocess.PIPE 创建管道
asyncio.subprocess.DEVNULL 丢弃
asyncio.subprocess.STDOUT ★(仅 stderr)合并到 stdout★
None ★继承父进程的(默认)★
limit=65536 StreamReader 的缓冲上限(★读超长行时要调大★)
★ 安全清单:
□ ★默认用 exec,不用 shell★
□ 必须用 shell 时对用户输入 ★shlex.quote★
□ ★不要把用户输入拼进命令名★(哪怕用 exec)
□ 用绝对路径或校验过的白名单命令
□ 设置 cwd 和最小化的 env
□ ★限制并发数和超时★(防资源耗尽)
create_subprocess_exec 是首选:参数是独立的列表项、直接传给 execve、不经过 shell,因此没有命令注入风险(用户输入里的 ;、|、$ 都只是普通字符);代价是不支持管道、重定向、通配符。create_subprocess_shell 走 /bin/sh -c,支持 shell 的全部特性,但用户输入拼进去就是命令注入——必须用 shlex.quote() 转义,更好的做法是用 exec 并在 Python 里实现管道逻辑(把 p1.stdout 直接作为 p2.stdin 传给第二个进程)。参数上有两个易错点:env 是完全替换而不是追加(想追加要写 env={**os.environ, ...}),以及 limit 控制 StreamReader 的缓冲上限(读超长行时要调大,否则会抛 LimitOverrunError)。安全清单里最重要的三条:默认用 exec、必须用 shell 时 shlex.quote、不要把用户输入拼进命令名。
三、管道死锁:最经典的陷阱
★ 死锁是怎么发生的:
proc = await create_subprocess_exec(cmd,
stdout=PIPE, stderr=PIPE)
out = await proc.stdout.read() # ★只读 stdout★
await proc.wait()
时间线:
t0 子进程往 stdout 和 stderr 都在写
t1 你在读 stdout(OK)
t2 ★stderr 的管道缓冲区写满了(Linux 通常 64KB)★
t3 ★子进程阻塞在 write(stderr) 上★
t4 子进程不再往 stdout 写 → ★你的 read() 也等不到数据★
→ ★双方互相等待,永久死锁★
★ 同样的死锁也会发生在:
- 写 stdin 太多而不读 stdout(子进程的输出把管道写满了)
- 只 await proc.wait() 而不读管道(★子进程被输出堵死,永远不退出★)
★ 三种正确做法:
① ★communicate()(最省事)★
out, err = await proc.communicate()
→ 内部★同时并发读 stdout 和 stderr、写 stdin★,不会死锁
✗ 缺点:★把全部输出读进内存★(输出几个 GB 就 OOM)
② ★合并 stderr 到 stdout(流式处理时最简单)★
proc = await create_subprocess_exec(
*args, stdout=PIPE, stderr=asyncio.subprocess.STDOUT)
async for line in proc.stdout: # ★只有一个管道,不会死锁★
handle(line)
✗ 缺点:分不清哪些是错误输出
③ ★两个管道各起一个任务并发读★
async with asyncio.TaskGroup() as tg:
tg.create_task(read_stream(proc.stdout, "out"))
tg.create_task(read_stream(proc.stderr, "err"))
await proc.wait()
✓ 既流式又能区分 → ★大输出 + 需要区分时用这个★
★ 不需要的管道要明确处理:
stderr=asyncio.subprocess.DEVNULL ★不关心就丢弃★
stdout=None 继承父进程(★直接打到终端/日志★)
★ 不用的管道千万别设成 PIPE 又不读
★ StreamReader 的读取方式:
await proc.stdout.read() 读到 EOF(★全部进内存★)
await proc.stdout.read(n) 最多读 n 字节
await proc.stdout.readline() 读一行(★超过 limit 抛 LimitOverrunError★)
await proc.stdout.readexactly(n) 精确读 n 字节
async for line in proc.stdout: ★异步迭代,最常用★
★ 二进制输出用 read(n) 分块;文本输出用 async for
★ 写 stdin 的正确姿势:
proc.stdin.write(data)
await proc.stdin.drain() # ★等待缓冲区排空(背压)★
proc.stdin.close() # ★关闭 → 子进程收到 EOF★
await proc.stdin.wait_closed() # 3.7+
★ 不 close 的话,等待 EOF 的子进程(如 cat、grep)★永远不会退出★
管道死锁是最经典的陷阱:设了 stdout=PIPE, stderr=PIPE 却只读其中一个,另一个的管道缓冲区(Linux 上通常 64KB)写满时子进程会阻塞在 write 上,于是它不再往你正在读的那个管道写数据——双方互相等待、永久卡住。同样的死锁也发生在「只 await proc.wait() 而不读管道」(子进程被自己的输出堵死,永远不退出)。三种正确做法:① communicate()(内部并发读写,最省事,但会把全部输出读进内存);② 把 stderr 合并进 stdout(stderr=asyncio.subprocess.STDOUT,只有一个管道就不会死锁,但分不清错误输出);③ 为两个管道各起一个任务并发读(既流式又能区分,大输出场景用这个)。还有两个细节:不用的管道别设成 PIPE(用 DEVNULL 丢弃或用 None 继承父进程),以及写完 stdin 必须 close()——否则等待 EOF 的子进程(cat、grep)永远不会退出。
四、超时、终止与僵尸进程
★ 超时的正确写法(★三步缺一不可★):
try:
out, err = await asyncio.wait_for(proc.communicate(), timeout=30)
except asyncio.TimeoutError:
proc.kill() # ★① 发 SIGKILL★
await proc.wait() # ★② 必须 wait!否则是僵尸★
raise # ③ 向上传播
★ 为什么必须 wait:
kill() 只是★发信号★;子进程死后,内核保留它的退出状态等父进程来取
→ ★不 wait = 僵尸进程★(占 PID,累积到 pid_max 会导致整机 fork 失败)
★ 这和同步 subprocess 的规则完全一样
★ 优雅终止的三段式:
proc.terminate() # ★SIGTERM:给它机会清理★
try:
await asyncio.wait_for(proc.wait(), timeout=5)
except asyncio.TimeoutError:
proc.kill() # ★SIGKILL:兜底★
await proc.wait()
★ 终止整个进程组(★子进程又起了孙进程时★):
proc = await asyncio.create_subprocess_exec(
*args, preexec_fn=os.setsid) # ★新建会话/进程组(Unix)★
...
os.killpg(os.getpgid(proc.pid), signal.SIGTERM) # ★杀整个组★
★ 否则:kill 了 shell,它启动的实际命令★变成孤儿继续跑★
★ 注意 preexec_fn 在多线程程序里不安全(fork 安全问题),
3.11+ 可以用 process_group=0 参数(更安全)
★ 任务被取消时子进程会怎样:
task = asyncio.create_task(run_cmd())
task.cancel()
→ 协程在 await communicate() 处收到 CancelledError
→ ★但子进程不会自动被杀★!它会继续运行
✓ 必须自己处理:
try:
out, err = await proc.communicate()
except asyncio.CancelledError:
proc.kill()
await proc.wait() # ★清理★
raise
★ returncode 的含义:
0 成功
> 0 命令自己的退出码
★< 0★ ★被信号杀死,-N 表示信号 N★(-9 = SIGKILL,-15 = SIGTERM)
None ★还在运行★(wait/communicate 之前)
★ 完整的健壮封装:
async def run_safe(*args, timeout=30, input=None):
proc = await asyncio.create_subprocess_exec(
*args, stdin=asyncio.subprocess.PIPE if input else None,
stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE)
try:
out, err = await asyncio.wait_for(
proc.communicate(input), timeout=timeout)
except (asyncio.TimeoutError, asyncio.CancelledError):
proc.kill()
await proc.wait() # ★两种异常路径都要清理★
raise
if proc.returncode != 0:
raise RuntimeError(
f"{args[0]} 失败 rc={proc.returncode}: {err.decode(errors='replace')}")
return out
超时处理必须是三步:wait_for 捕获超时 → proc.kill() → await proc.wait()。第三步不能省——kill() 只是发信号,子进程死后内核保留退出状态等父进程来取,不 wait 就是僵尸进程(占 PID,累积到上限会导致整机 fork 失败)。更礼貌的是三段式终止:terminate()(SIGTERM,给它机会清理)→ wait_for(proc.wait(), 5) → 超时才 kill()。有两个进阶场景:子进程又起了孙进程时要终止整个进程组(创建时用 process_group=0(3.11+)或 preexec_fn=os.setsid,终止时 os.killpg)——否则杀掉 shell 后它启动的实际命令会变成孤儿继续跑;以及任务被取消时子进程不会自动被杀,必须在 except asyncio.CancelledError 里 kill + wait 再重新 raise。最后记住 returncode 为负数表示被信号杀死(-9 就是 SIGKILL)。
五、平台差异与常见问题
★ Windows:
① ★必须用 ProactorEventLoop★(基于 IOCP)
- 3.8 起 Windows 默认就是它 → 一般不用管
- ★如果显式 asyncio.set_event_loop_policy(WindowsSelectorEventLoopPolicy)
→ 子进程会 NotImplementedError★
- 有些库(如老版本 aiodns、Tornado 集成)会建议切到 Selector → ★冲突★
② 没有信号语义:terminate() 实际调 TerminateProcess(★等价于 kill★)
③ 没有进程组的 killpg → 要用 taskkill /T /F 或 job object
★ Unix 的 child watcher(★3.12 之前的复杂之处★):
asyncio 需要知道"子进程什么时候退出"→ 依赖 SIGCHLD
3.8~3.11 有多种 watcher:
ThreadedChildWatcher(默认) ★每个子进程一个线程★等它,最兼容
MultiLoopChildWatcher、SafeChildWatcher、FastChildWatcher…
★ 常见问题:在非主线程的事件循环里创建子进程 → 老版本可能报
"RuntimeError: Cannot add child handler, the child watcher does not
have a loop attached"
✓ 3.12 起大幅简化(默认使用基于 pidfd 的实现,不再需要手动配置 watcher)
✓ 实践建议:★在主线程的事件循环里创建子进程★
★ 编码问题:
asyncio 子进程的输出永远是 ★bytes★(没有 text=True 参数)
→ 自己 decode:out.decode("utf-8", errors="replace")
→ ★Windows 上命令输出常是 GBK/cp936★ → decode("gbk") 或用 locale
→ 想让子进程输出 UTF-8:env={**os.environ, "PYTHONIOENCODING": "utf-8"}
★ 常见报错速查:
NotImplementedError
→ Windows 上用了 SelectorEventLoop
LimitOverrunError / ValueError: Separator is not found, and chunk exceed the limit
→ ★单行超过了 StreamReader 的 limit(默认 64KB)★
✓ create_subprocess_exec(..., limit=1024*1024) 调大,或用 read(n) 分块
BrokenPipeError(写 stdin 时)
→ 子进程已经退出了
子进程一直不退出
→ ★① 没 close stdin(它在等 EOF)② 管道满了死锁★
僵尸进程堆积
→ ★kill 之后没有 wait★
★ 与同步 subprocess 的差异对照:
┌──────────────────┬──────────────────┬──────────────────┐
│ │ subprocess │ asyncio │
├──────────────────┼──────────────────┼──────────────────┤
│ text/encoding │ ★支持★ │ ★不支持,只有 bytes★│
│ timeout 参数 │ ★run/wait 支持★ │ ★用 wait_for 包★ │
│ check=True │ ★支持★ │ ★自己判 returncode★│
│ capture_output │ ★支持★ │ 手动设 PIPE │
│ 阻塞 │ ★阻塞调用线程★ │ ★不阻塞事件循环★ │
└──────────────────┴──────────────────┴──────────────────┘
→ asyncio 版参数更少、更底层 → ★通常要自己封装一层★
平台差异要注意两点。Windows 必须用 ProactorEventLoop(3.8 起默认),如果某个库建议你切到 SelectorEventLoop(老版本 aiodns 等),子进程功能就会 NotImplementedError;而且 Windows 没有信号语义(terminate() 等价于 kill)也没有进程组。Unix 上 3.12 之前有 child watcher 的复杂性(默认的 ThreadedChildWatcher 会为每个子进程起一个线程),在非主线程的事件循环里创建子进程可能报错——实践建议是在主线程的事件循环里创建子进程;3.12 起已大幅简化。另外要记住 asyncio 子进程的输出永远是 bytes(没有 text=True 参数),需要自己 decode(Windows 上命令输出常是 GBK)。常见报错里最容易困惑的是 LimitOverrunError(单行超过 StreamReader 的 64KB 上限,要调大 limit 或改用 read(n) 分块)和子进程一直不退出(多半是没 close stdin 或管道满了死锁)。
六、实践模式
★ 模式一:健壮的通用封装(★可直接抄★)
async def run_cmd(*args, timeout=30, input_=None, cwd=None, env=None,
check=True) -> tuple[int, bytes, bytes]:
proc = await asyncio.create_subprocess_exec(
*args,
stdin=asyncio.subprocess.PIPE if input_ is not None else None,
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE,
cwd=cwd, env=env)
try:
out, err = await asyncio.wait_for(
proc.communicate(input_), timeout=timeout)
except (asyncio.TimeoutError, asyncio.CancelledError):
proc.kill()
await proc.wait() # ★清理,防僵尸★
raise
if check and proc.returncode != 0:
raise RuntimeError(f"{args[0]} rc={proc.returncode}: "
f"{err.decode(errors='replace')[:500]}")
return proc.returncode, out, err
★ 模式二:流式处理长时间任务的输出(★进度上报★)
async def run_with_progress(*args, on_line):
proc = await asyncio.create_subprocess_exec(
*args, stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.STDOUT) # ★合并,避免死锁★
async for raw in proc.stdout:
line = raw.decode(errors="replace").rstrip()
await on_line(line) # ★实时上报进度★
return await proc.wait()
# 用于 ffmpeg 转码、大型构建、数据导出等
★ 模式三:并发执行 + 限流 + 超时
async def run_batch(cmds, concurrency=4, timeout=60):
sem = asyncio.Semaphore(concurrency) # ★进程是重资源,并发要小★
async def one(cmd):
async with sem:
try:
return await run_cmd(*cmd, timeout=timeout)
except Exception as e:
return e
return await asyncio.gather(*(one(c) for c in cmds))
★ 并发数怎么定:CPU 密集的命令 ≈ CPU 核数;IO 密集可以更高
★ 别忘了每个子进程都占内存和 fd
★ 模式四:把管道接起来(不用 shell)
async def pipeline(*stages):
procs = []
prev_out = None
for stage in stages:
p = await asyncio.create_subprocess_exec(
*stage, stdin=prev_out,
stdout=asyncio.subprocess.PIPE)
if prev_out:
prev_out.close() # ★父进程关掉自己那份 fd★
procs.append(p)
prev_out = p.stdout
out, _ = await procs[-1].communicate()
for p in procs[:-1]:
await p.wait()
return out
★ 检查清单:
□ ★协程里绝不用 subprocess.run★
□ ★默认 exec,shell 只在必要时用且转义★
□ ★用 communicate() 或并发读两个管道★(防死锁)
□ 不用的管道设 DEVNULL 或 None
□ ★写完 stdin 要 close★(否则子进程等 EOF 不退出)
□ ★超时后 kill + wait★(防僵尸)
□ ★捕获 CancelledError 也要 kill + wait★
□ 起了孙进程 → ★用进程组统一终止★
□ ★限制并发数★(进程是重资源)
□ 输出是 bytes,注意 decode 的编码和 errors 策略
□ 大输出用流式而不是 communicate(防 OOM)
四种实践模式覆盖了主要场景。通用封装要处理超时、取消、退出码检查、以及两种异常路径下的 kill + wait。流式处理用于长时间任务的进度上报(ffmpeg 转码、构建、数据导出),把 stderr 合并进 stdout 是最简单的防死锁方式。并发执行要限流——进程是重资源(每个都占内存和 fd),CPU 密集的命令并发数约等于核数。手动接管道可以完全避免 shell(注意父进程要 close() 掉自己那份中间管道的 fd)。检查清单里最容易漏的四条:写完 stdin 要 close()、超时和取消两条路径都要 kill + wait、起了孙进程要用进程组统一终止、大输出用流式而不是 communicate()(后者会把全部输出读进内存)。
记忆钩子:**「★协程里绝对不能用 subprocess.run★——它是同步阻塞的,会让整个事件循环停摆(其他请求不处理、超时不触发、心跳断连);要用 ★asyncio.create_subprocess_exec(args)★(参数列表、不经 shell、★无注入风险★)或 create_subprocess_shell(cmd)(走 /bin/sh、支持管道重定向、★但用户输入拼进去就是命令注入,必须 shlex.quote★)。★最经典的陷阱是管道死锁★:设了 stdout=PIPE, stderr=PIPE 却只读一个,另一个的★管道缓冲区(Linux 约 64KB)写满时子进程阻塞在 write 上★,于是它不再往你读的那个管道写 → 双方互相等待、永久卡住;三种解法:★communicate()(内部并发读写,最省事但会把全部输出读进内存)、把 stderr 合并进 stdout(stderr=asyncio.subprocess.STDOUT)、或为两个管道各起一个任务并发读★。同类问题:★只 await proc.wait() 不读管道★也会被输出堵死,★写完 stdin 不 close★会让等 EOF 的子进程(cat/grep)永远不退出。★超时必须三步:wait_for 捕获 → proc.kill() → await proc.wait()★——kill 只是发信号,★不 wait 就是僵尸进程★(占 PID);更礼貌是 terminate → 等 5 秒 → 再 kill。两个进阶点:★任务被 cancel 时子进程不会自动被杀★(要在 except CancelledError 里 kill+wait 再 raise)、★子进程又起了孙进程时要用进程组统一终止★(3.11+ 用 process_group=0,否则杀了 shell 而实际命令变孤儿继续跑)。平台:★Windows 必须用 ProactorEventLoop★(3.8+ 默认;切到 Selector 会 NotImplementedError);Unix 3.12 前有 child watcher 的复杂性,建议★在主线程的事件循环里创建子进程★。其他细节:★asyncio 版没有 text=True,输出永远是 bytes★(Windows 上常是 GBK);单行超过 64KB 会抛 LimitOverrunError(调大 limit);returncode ★负数表示被信号杀死★;进程是重资源,★并发要用 Semaphore 限制★。」*
七、常见误区与追问
- 误区:在
async def里调用subprocess.run()只是稍微慢一点,能跑通就行。 它会阻塞整个事件循环——外部命令跑多久,你的服务就有多久完全无响应:其他请求的协程一个都不执行、asyncio.sleep到期的任务不被唤醒(所有超时判断失真)、心跳和 keepalive 发不出去(连接被对端断开)、健康检查超时甚至可能被编排系统重启。一个 5 秒的 ffmpeg 调用,在 QPS 100 的服务上意味着 500 个请求同时被卡住。正确做法是用asyncio.create_subprocess_exec(原生异步、不占线程);如果是第三方库内部封装了同步subprocess而你改不了,就用await asyncio.to_thread(...)丢到线程池(代价是占一个线程槽位,默认只有min(32, cpu+4)个)。 - 误区:设置了
stdout=PIPE, stderr=PIPE之后,只读需要的那个流就行。 这是管道死锁的标准触发方式。管道的内核缓冲区是有限的(Linux 上通常 64KB),没有被读取的那个流一旦写满,子进程就会阻塞在write()系统调用上——它被卡住后自然也不会再往你正在读的那个流写数据,于是你在等它的输出、它在等你读走 stderr,双方永久互等。同样的死锁还会发生在「只await proc.wait()而完全不读管道」的情况(子进程被自己的输出堵死,永远不退出,你的wait()也就永远返回不了)。三种正确做法:用communicate()(内部并发读两个流并写 stdin)、把 stderr 合并进 stdout(stderr=asyncio.subprocess.STDOUT)、或为两个流各起一个任务并发读。用不到的流应该显式设成DEVNULL或保持None(继承父进程),绝不要设成 PIPE 却不读。 - 误区:
proc.kill()之后子进程就彻底清理干净了。kill()只是发送信号;子进程终止后,内核会保留它的退出状态(进程表项)直到父进程调用wait()取走——这段时间它就是僵尸进程。所以超时处理必须是proc.kill()加上await proc.wait()两步。在长期运行的服务里,只 kill 不 wait 会让僵尸持续累积,每个占用一个 PID,最终撞上pid_max或容器的pids.max限制,导致整个系统无法创建新进程。同样的规则也适用于任务被取消的路径:task.cancel()会让协程在await communicate()处抛CancelledError,但子进程不会自动被杀,必须在except asyncio.CancelledError:里做kill + wait再重新raise。 - 误区:
create_subprocess_shell用起来更方便,反正命令是我自己写的。 只要命令里拼接了任何外部输入(用户参数、文件名、数据库字段、环境变量),就是命令注入漏洞:f"cat {filename}"遇到filename = "a.txt; rm -rf /"就会执行删除操作。而且注入的形式很多(;、&&、|、反引号、$()),黑名单过滤几乎必然被绕过。默认应该用create_subprocess_exec——参数作为独立的列表项直接传给execve,根本不经过 shell 解析,用户输入里的特殊字符只是普通字符。确实需要 shell 特性(管道、重定向、通配符、shell 内建命令)时有两条路:用shlex.quote()对每个变量部分转义,或者用 exec 在 Python 里自己实现管道(把前一个进程的stdout作为后一个的stdin)。另外还要注意:即使用 exec,也不要把用户输入拼进「命令名」本身。 - 误区:子进程运行完了但程序一直卡着,肯定是命令本身有问题。 更常见的是你这边的用法有问题,三种典型情况:① 没有关闭 stdin —— 像
cat、grep、sort这类从标准输入读数据的命令会一直等待 EOF,你必须proc.stdin.close()(或用communicate(input),它会帮你关)才能让它们退出;② 管道死锁(见前面);③ 子进程自己又 fork 了孙进程并继承了管道 —— 即使子进程退出了,孙进程还持有管道的写端,所以你读不到 EOF(典型场景是通过 shell 启动的服务)。排查方法:先用ps看子进程还在不在(在 = 它没退出,查 ① ②;不在 = 查 ③),也可以在命令行手动跑一遍同样的命令看行为。 - 追问:asyncio 的子进程和同步
subprocess在 API 上有哪些差异要注意? 四个主要差异。① 没有text=True/encoding参数——asyncio 版输出永远是bytes,必须自己decode(推荐带errors="replace"防止编码问题导致异常,Windows 上命令输出常是 GBK/cp936 而不是 UTF-8)。② 没有timeout参数——要用asyncio.wait_for(proc.communicate(), timeout=N)包一层,且超时后要自己kill + wait。③ 没有check=True——要自己判断proc.returncode != 0并抛异常。④ 没有capture_output=True——要显式写stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE。此外还有一些同步版有而异步版没有的参数(text、errors、universal_newlines)。因为参数更少更底层,实践中通常要自己封装一层run_cmd()工具函数,把超时、取消清理、退出码检查、解码策略都固化进去,业务代码只调这个封装。 - 追问:怎么终止一个「自己还会启动其他进程」的子进程? 直接
proc.kill()只杀掉你直接创建的那一个进程——它启动的孙进程会变成孤儿被 PID 1 收养并继续运行。这在「通过 shell 启动一条命令」时特别常见:你杀掉的是sh,而真正干活的ffmpeg还在跑(而且它继续持有管道,导致你连 EOF 都读不到)。正确做法是用进程组统一终止:创建时让子进程成为新进程组的组长——Python 3.11+ 用process_group=0参数(更安全),旧版本用preexec_fn=os.setsid(注意preexec_fn在多线程程序里有 fork 安全隐患);终止时用os.killpg(os.getpgid(proc.pid), signal.SIGTERM)向整个进程组发信号,超时后再发SIGKILL。Windows 上没有进程组概念,要用taskkill /T /F /PID或 Job Object。另外别忘了最后仍然要await proc.wait()回收。 - 追问:需要处理大量输出(几个 GB)时该怎么写? 绝对不能用
communicate()—— 它会把全部输出读进内存,几个 GB 的输出直接 OOM。正确做法是流式读取并边读边处理:把stderr合并进stdout(stderr=asyncio.subprocess.STDOUT,避免双管道死锁),然后async for line in proc.stdout:逐行处理,或者用await proc.stdout.read(65536)分块读二进制数据。三个配套细节:① 注意LimitOverrunError——StreamReader的默认行缓冲上限是 64KB,如果输出中有超长行(比如一整个 JSON 打在一行里),readline/async for会抛异常,需要在创建时传limit=1024*1024调大,或改用read(n)分块读。② 处理逻辑本身不能太慢 —— 如果你的处理跟不上子进程的输出速度,管道缓冲区会满、子进程会被自动阻塞(这其实是天然的背压,通常是好事)。③ 如果只是想把输出写到文件,最高效的是根本不用管道:直接把stdout设成一个打开的文件对象(stdout=open("out.log", "wb")),让内核直接写,完全不经过 Python。
八、加强记忆
协程里绝对不能用 subprocess.run()——它是同步阻塞的,会让整个事件循环停摆(其他请求不处理、超时不触发、心跳断连、健康检查失败);要用 asyncio.create_subprocess_exec(*args)(参数列表、不经 shell、无注入风险)或 create_subprocess_shell(cmd)(走 /bin/sh、支持管道重定向,但用户输入拼进去就是命令注入,必须 shlex.quote)。最经典的陷阱是管道死锁:设了 stdout=PIPE, stderr=PIPE 却只读一个,另一个的管道缓冲区(Linux 约 64KB)写满时子进程会阻塞在 write 上,于是它不再往你读的那个管道写数据——双方互相等待、永久卡住;三种解法:communicate()(内部并发读写,最省事但会把全部输出读进内存)、把 stderr 合并进 stdout(stderr=asyncio.subprocess.STDOUT)、或为两个管道各起一个任务并发读。同类问题还有:只 await proc.wait() 而不读管道也会被输出堵死,写完 stdin 不 close() 会让等待 EOF 的子进程(cat/grep)永远不退出。超时必须三步:wait_for 捕获 → proc.kill() → await proc.wait()——kill 只是发信号,不 wait 就是僵尸进程(占 PID,累积会导致整机 fork 失败);更礼貌的是 terminate() → 等几秒 → 再 kill()。两个进阶点:任务被 cancel 时子进程不会自动被杀(要在 except CancelledError 里 kill + wait 再 raise)、子进程又起了孙进程时要用进程组统一终止(3.11+ 用 process_group=0,否则杀了 shell 而实际命令变成孤儿继续跑、还会堵着管道让你读不到 EOF)。平台方面:Windows 必须用 ProactorEventLoop(3.8+ 默认,切到 Selector 会 NotImplementedError);Unix 在 3.12 之前有 child watcher 的复杂性,建议在主线程的事件循环里创建子进程。其余细节:asyncio 版没有 text=True,输出永远是 bytes(Windows 上常是 GBK,decode 时带 errors="replace")、也没有 timeout/check/capture_output(要自己封装);单行超过 64KB 会抛 LimitOverrunError(调大 limit 或改用 read(n));returncode 为负数表示被信号杀死;进程是重资源,并发必须用 Semaphore 限制;大输出要流式处理,只想写文件时直接把 stdout 设成文件对象最高效。