You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Python类中如何在主循环内非阻塞运行并行异步任务?

Windows 11下Python 3.11非阻塞并行进程实现问题

我在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.14 10:23:12