Python异步无轮询等待可取消子进程 协程取消时await卡住问题
问题背景
需要在Python异步函数中以子进程形式运行同步函数,示例同步子进程代码如下:
def subProcess(): print(f"[{os.getpid()}] subProcess running", file=sys.stderr) try: time.sleep(5) finally: print(f"[{os.getpid()}] subProcess closing", file=sys.stderr)
需要实现两个核心能力:
- 可await子进程终止并获取退出码,区分优雅退出(包括SIGTERM、SIGINT触发的退出、退出码为0的场景)和其他异常退出场景,以此判断是否需要重启失败的子进程。
- 若启动子进程的协程被取消,能够自动终止对应子进程;且协程返回前可等待子进程完全退出,例如子进程占用某端口提供服务时,必须等待旧进程完全退出释放端口后,才能启动新进程,避免端口不可用问题。
现有方案存在明显局限:
- 无法使用
asyncio.get_event_loop().run_in_executor(executor=ProcessPool(), func=subProcess)实现需求,该方式没有提供终止子进程的途径。 - 无法在异步代码中直接使用
multiprocessing.Process。
原型实现逻辑
最初的实现思路如下:
- 协程中创建非异步的
multiprocessing.Process对象; - 创建工作线程执行Process的
.start()和.join()操作,等待进程退出并获取退出码,操作完成后将退出码或捕获的异常写入Future对象; - 原始协程等待该Future返回结果;
- 协程可能被取消,这种场景下协程会向子进程发送SIGINT信号终止进程;由于原本等待的Future已被取消,工作线程会将结果写入第二个Future,供协程后续等待。
故障现象
问题出在协程取消后的等待逻辑:
- 子进程被终止后,工作线程会向两个Future写入结果,调试日志已经打印
runSubprocess2:thread set2,可以确认result2.set_result方法已经被调用。 - 但协程等待第二个Future时,
await result2永远不会返回也不会抛出异常,卡在此处:
exitcode = await result2
完整可复现代码(已在Python3.8、3.11版本测试复现):
import asyncio import os import signal import sys import threading import time from multiprocessing import Process def subProcess(*args): print(f"[{os.getpid()}] subProcess{args} running", file=sys.stderr) try: time.sleep(5) finally: print(f"[{os.getpid()}] subProcess{args} done", file=sys.stderr) MISSING = object() async def runSubprocess2(target=subProcess, args=(), name=None): def task(): # Sync code to start, join the Process and write the # exitcode/exception to the outer cooroutine's futures. # This all works as expected. Both result2 is assigned to # (result1 was cancelled already) p.start() r = None try: p.join() r = p.exitcode if 1: print(f"runSubprocess2:thread set2") result2.set_result(r) if not result1.cancelled(): print(f"runSubprocess2:thread set1") result1.set_result(r) except BaseException as e: print(f"runSubprocess2:thread exception {type(e)}") if not result2.cancelled(): result2.set_exception(e) print("runSubprocess2:thread result2 set ex") if not result1.cancelled(): result1.set_exception(e) print("runSubprocess2:thread result1 set ex") else: print(f"runSubprocess2:thread terminated") print(f"runSubprocess2:thread exception OK") finally: print(f"runSubprocess2:thread terminated (exitcode={r})") result1 = asyncio.Future() result2 = asyncio.Future() # Used if result1 is cancelled p = Process(target=target, args=args, daemon=True, name=name) threading.Thread( target=task, daemon=True, name=name ).start() # task = asyncio.get_event_loop().run_in_executor(ProcessPoolExecutor(), subProcess) exitcode = MISSING try: exitcode = await result1 except asyncio.CancelledError as e: print(f"runSubprocess2:coro cancelled {type(e)}") except BaseException as e: print(f"runSubprocess2:coro exception {type(e)}") raise finally: # If the process was not killed, send a KeyboardInterrupt wait = not result1.done() if result1.cancelled(): # This path is visited in this example: wait = True print(f"runSubprocess2:coro send {p.pid} terminate ...") # p.terminate() os.kill(p.pid, signal.SIGINT) print("runSubprocess2:coro sent SIGINT ... OK") if wait: print(f"runSubprocess2:coro awaiting shutdown ... {result2}") try: print(f"runSubprocess2:coro awaiting 2") exitcode = await result2 # <<< STUCK HERE. What the hey! print(f"runSubprocess2:coro awaiting shutdown ... OK exitcode={exitcode}") except BaseException as e: print(f"runSubprocess2:coro awaiting shutdown ... {type(e)}") # ??? raise finally: print(f"runSubprocess2:coro awaiting shutdown ... DONE {result2}") else: print("runSubprocess2:coro terminated") assert exitcode is not MISSING # Handle exit code cases here. async def main(): print(f"[{os.getpid()}] main", file=sys.stderr) if 1: nTasks = 1 tasks = [asyncio.create_task(runSubprocess2(target=subProcess, args=(i,))) #asyncio.get_event_loop().run_in_executor(executor=None, func=runSubprocess2) for i in range(nTasks)] await asyncio.sleep(1) print("Cancelling...") for task in tasks: task.cancel() print("Cancelling... OK") try: await asyncio.wait(tasks) except asyncio.CancelledError: pass await asyncio.wait(tasks) print("Done", file=sys.stderr) if __name__ =='__main__': def _main(): try: asyncio.run(main()) finally: print(f"Main exit.")
故障原因
核心原因是asyncio.Future不是线程安全的。asyncio.Future设计上只允许在运行事件循环的线程内操作,你在独立工作线程中直接调用result2.set_result(r),虽然Future内部的值确实被修改了,但事件循环完全感知不到Future的状态变更,不会触发注册的回调唤醒等待的协程,await自然会永久挂起。跨线程直接操作asyncio Future属于未定义行为,哪怕日志显示set_result执行了,事件循环也不会调度后续协程逻辑。
修复方案
跨线程给asyncio Future设置结果,必须通过loop.call_soon_threadsafe()方法把设置结果的操作调度到事件循环线程执行,保证事件循环能感知到Future状态变化。
- 先在协程中获取当前运行的事件循环引用,传给工作线程
- 所有对
result1、result2的set_result/set_exception操作,都用loop.call_soon_threadsafe()包装
修改后的核心代码片段:
async def runSubprocess2(target=subProcess, args=(), name=None): # 获取当前事件循环引用 loop = asyncio.get_running_loop() def task(): p.start() r = None try: p.join() r = p.exitcode print(f"runSubprocess2:thread set2") # 跨线程调度set_result到事件循环线程执行 loop.call_soon_threadsafe(result2.set_result, r) if not result1.cancelled(): print(f"runSubprocess2:thread set1") loop.call_soon_threadsafe(result1.set_result, r) except BaseException as e: print(f"runSubprocess2:thread exception {type(e)}") if not result2.cancelled(): loop.call_soon_threadsafe(result2.set_exception, e) print("runSubprocess2:thread result2 set ex") if not result1.cancelled(): loop.call_soon_threadsafe(result1.set_exception, e) print("runSubprocess2:thread result1 set ex") else: print(f"runSubprocess2:thread terminated") print(f"runSubprocess2:thread exception OK") finally: print(f"runSubprocess2:thread terminated (exitcode={r})") # 其余原有逻辑保持不变即可
额外优化建议:不需要维护两个Future,只保留一个Future即可。协程取消时先给子进程发送终止信号,工作线程正常把退出码写入同一个Future,协程在CancelledError的finally块里直接await这个Future就行,逻辑更简洁,不容易出错。
内容的提问来源于stack exchange,提问作者user48956
相关产品推荐
相关产品推荐

