Python类中如何在主循环内非阻塞运行并行异步任务?
我在Windows 11系统的Python 3.11环境中开发了一个类,其中包含永久运行的主循环loop A。需求是:当loop A内满足指定条件时,启动N个并行异步的loop B进程,且不能阻塞loop A的执行;之后每次loop A迭代时,检查loop B的结果是否可用,若可用则执行后续处理,否则继续loop A的逻辑。
我的初始代码框架如下:
class MyClass: def __init__(self, ): pass def loop_A(self): while True: ... do some stuff ... if condition: res_B = launch_pool_of_B_in_parallel(n_processes=N, args=([some_args[i]) for i in range(N)) if res_B is done: post_process(res_B) else: pass # go back to loop A
我尝试用asyncio实现非阻塞但始终失败;用ThreadPoolExecutor实现线程并发时,不仅无法使用进程,还会阻塞loop A。我写的启动loop B的辅助函数如下:
def run_parallel_loop(self, func, args_): with concurrent.futures.ThreadPoolExecutor(n_jobs) as executor: loop = asyncio.get_event_loop() tasks = [loop.run_in_executor(executor, func, *args_[i]) for i in range(self.n_jobs)] results = await asyncio.gather(*tasks) return results
在loop A中启动loop B的代码:
if not task_running: async_jobs = asyncio.create_task(self.run_parallel_jobs()) results = []
检查结果的逻辑:
if len(results) < self.n_jobs: try: results_ = await asyncio.wait_for(async_jobs, timeout=0.1) results.extend(results_) except asyncio.TimeoutError: pass if len(results) == self.n_jobs: res_tmp = results # ready to post-process task_running = False # ready to launch again async processes
目前遇到多种错误(如asyncio.exceptions.CancelledError、无法获取loop B结果、阻塞loop A等),方法明显有问题,请求帮助。
另外,使用ProcessPoolExecutor时出现TypeError: cannot pickle '_io.BufferedReader' object错误,推测是进程间内存不共享导致,是否需要创建LoopB类来解决?
解决方案
1. 核心问题修正:同步loop A与并行进程的非阻塞协作
你的loop A是同步永久循环,不能直接用asyncio的await(会阻塞循环),正确的做法是用concurrent.futures.ProcessPoolExecutor提交任务,通过检查Future对象的状态实现非阻塞监控:
import concurrent.futures import time class MyClass: def __init__(self, n_processes): self.n_processes = n_processes # 初始化进程池(全局复用,避免反复创建销毁) self.executor = concurrent.futures.ProcessPoolExecutor(max_workers=n_processes) self.running_futures = [] self.task_running = False self.collected_results = [] def loop_B_func(self, args): # 这里是你的loop B逻辑,注意参数必须可序列化 time.sleep(2) # 模拟耗时操作 return args * 2 def loop_A(self): while True: # ... 执行loop A的常规逻辑 ... print("Running loop A iteration...") time.sleep(0.5) # 满足条件且无任务在运行时,启动并行loop B condition = True # 替换为你的实际条件 if condition and not self.task_running: print("Launching parallel loop B tasks...") # 提交N个任务到进程池 self.running_futures = [ self.executor.submit(self.loop_B_func, i) for i in range(self.n_processes) ] self.task_running = True self.collected_results = [] # 检查并行任务状态,非阻塞 if self.task_running: # 遍历所有未完成的future,收集已完成的结果 completed, pending = concurrent.futures.wait( self.running_futures, timeout=0, return_when=concurrent.futures.FIRST_COMPLETED ) for future in completed: self.collected_results.append(future.result()) # 更新未完成的任务列表 self.running_futures = list(pending) # 所有任务完成,执行后续处理 if not self.running_futures: print(f"All loop B tasks done, results: {self.collected_results}") self.post_process(self.collected_results) self.task_running = False def post_process(self, results): # 你的结果处理逻辑 pass if __name__ == "__main__": obj = MyClass(n_processes=3) obj.loop_A()
2. 解决ProcessPoolExecutor的pickle错误
出现cannot pickle '_io.BufferedReader' object是因为你传递给子进程的参数包含不可序列化的对象(比如文件句柄),和是否创建LoopB类无关。解决方法:
- 不要传递文件对象,而是传递文件路径,让子进程自己打开文件读取内容;
- 如果是自定义对象,确保实现了
__reduce__方法或使用cloudpickle替代默认pickle(但优先尽量用原生可序列化类型)。
3. 原代码的错误原因
- 同步
loop A中使用await会阻塞整个循环,违背非阻塞需求; ThreadPoolExecutor受GIL限制,仅适合IO密集型任务,且你用asyncio.run_in_executor配合await的方式依然会阻塞同步循环;- 用
asyncio.create_task但未在事件循环中运行,导致任务无法正确执行,进而出现CancelledError。
内容的提问来源于stack exchange,提问作者Yoda

