使用asyncio调用子进程需按序执行并满足条件时终止任务
AsyncIO子进程有序执行与终止优化方案
核心问题
- 子进程执行顺序混乱,无法满足基准测试的顺序输出要求
- 执行大量子进程(如240个)时,触发
ValueError: I/O operation on closed pipe错误 - 找到正确数字后,无法安全且及时终止剩余任务
现有Python代码
import asyncio import time async def call_subprocess(num: int): cmd_go = "guess_the_number_go.exe" # cmd_c = "guess_the_number_c.exe" process = await asyncio.create_subprocess_exec( cmd_go, stdin=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE) user_input = f"{num}\n".encode('utf-8') # encode to bytes process.stdin.write(user_input) bytes_reply, trash = await process.communicate() str_data = bytes_reply.decode("utf-8") cleaned_reply = str_data.strip("\n").strip("\r").split("\n")[-1] BAD_MESSAGE = "Sorry but your number is wrong, please try again." if cleaned_reply == BAD_MESSAGE: print(f"You entered {num}, it is the wrong number please try again!") print("---------------------------------------------------------") else: print(f"\nGreat your guess was right! The correct number is: {num}\n") print("\n*********************************************************") # cancel all the other tasks, since we have found the correct number for task in asyncio.Task.all_tasks(): task.cancel() async def main(number: int): start_time = time.time() tasks = [asyncio.create_task(call_subprocess(i)) for i in range(number)] try: _ = await asyncio.gather(*tasks, return_exceptions=True) except asyncio.CancelledError: print("Remaining tasks were cancelled!") end_time = time.time() elapsed_time = end_time - start_time print("--------------------------------------------------") print(f"Execution took {round(elapsed_time, 2)} seconds.") asyncio.run(main(130))
输出对比
当前无序输出
You entered 0, it is the wrong number please try again!
You entered 4, it is the wrong number please try again!
You entered 2, it is the wrong number please try again!
...
You entered 124, it is the wrong number please try again!
Great your guess was right! The correct number is: 123
期望有序输出
You entered 0, it is the wrong number please try again!
You entered 1, it is the wrong number please try again!
You entered 2, it is the wrong number please try again!
...
You entered 122, it is the wrong number please try again!
Great your guess was right! The correct number is: 123
子进程代码
Go代码(guess_the_number.go)
package main // compile with: >go build guess_the_number.go import "fmt" func main() { var secretNumber = 123 var userInput int fmt.Printf("Please input a number: ") _, err := fmt.Scanf("%d", &userInput) if err != nil { fmt.Println("Something went wrong, only input numbers!") panic("Closing program!") } fmt.Printf("Your chosen number was: %d\n", userInput) if (userInput == secretNumber) { fmt.Printf("Congratulation you guessed the correct number!\n") } else { fmt.Printf("Sorry but your number is wrong, please try again.\n") } }
C代码(guess_the_number.c)
#include <stdio.h> // compile with: >gcc -o guess_the_number_c.exe guess_the_number.c int main(void) { int secretNumber = 123; int userInput; printf( "Please input a number: "); scanf("%i", &userInput); printf("Your entered number was: %i\n", userInput); if (userInput == secretNumber) { printf("Congratulation you guessed the correct number!\n"); } else { printf("Sorry but your number is wrong, please try again.\n"); }; return (0); }
优化后的解决方案代码
import asyncio import time async def call_subprocess(num: int, stop_event: asyncio.Event): # 检查终止信号,若已触发则直接返回 if stop_event.is_set(): return cmd_go = "guess_the_number_go.exe" # cmd_c = "guess_the_number_c.exe" process = await asyncio.create_subprocess_exec( cmd_go, stdin=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE) try: user_input = f"{num}\n".encode('utf-8') await process.stdin.write(user_input) await process.stdin.drain() # 确保数据完全写入管道,避免滞留 bytes_reply, _ = await process.communicate() str_data = bytes_reply.decode("utf-8") cleaned_reply = str_data.strip("\n").strip("\r").split("\n")[-1] BAD_MESSAGE = "Sorry but your number is wrong, please try again." if cleaned_reply == BAD_MESSAGE: print(f"You entered {num}, it is the wrong number please try again!") print("---------------------------------------------------------") else: print(f"\nGreat your guess was right! The correct number is: {num}\n") print("\n*********************************************************") stop_event.set() # 设置全局终止信号,停止后续任务 finally: # 确保子进程被终止,释放系统资源 if process.returncode is None: process.terminate() await process.wait() async def main(number: int): start_time = time.time() stop_event = asyncio.Event() # 逐个执行子进程,确保顺序,同时检查终止信号 for i in range(number): if stop_event.is_set(): break await call_subprocess(i, stop_event) end_time = time.time() elapsed_time = end_time - start_time print("--------------------------------------------------") print(f"Execution took {round(elapsed_time, 2)} seconds.") asyncio.run(main(130))
关键修改说明
- 有序执行控制:放弃批量创建所有任务的方式,改为循环逐个
await子进程任务,严格保证0到n的执行顺序,完全匹配期望输出。 - 安全终止机制:使用
asyncio.Event作为全局终止信号,找到正确数字后设置Event,后续任务启动前直接返回,已启动的任务通过finally块确保资源释放,避免强制取消任务导致的管道错误。 - 管道操作修复:添加
await process.stdin.drain()确保写入子进程的数据被完全发送,解决大数量任务时的数据滞留问题;finally块中主动终止未结束的子进程,释放管道资源,消除ValueError: I/O operation on closed pipe错误。 - 避免全局任务取消:原代码中
asyncio.Task.all_tasks()会取消包括main在内的所有任务,导致计时代码无法执行,优化后通过Event让任务自行终止,保证基准测试的计时逻辑完整执行。
内容的提问来源于stack exchange,提问作者Gambrinos
相关产品推荐
相关产品推荐

