如何让Python multiprocessing进程运行超60秒后自动终止?
实现多进程任务超时终止的修改方案
你当前使用的pool.starmap()是批量阻塞执行任务的方式,没有内置超时控制。要实现单个任务超过60秒就终止,有几种可行的修改方案,以下是具体实现:
方案一:使用concurrent.futures.ProcessPoolExecutor(推荐)
这个模块是对multiprocessing的高层封装,原生支持任务超时处理,代码更简洁:
import concurrent.futures # 保持原arguments不变 arguments = [(self.meshgrid_lat, self.meshgrid_lon, self.dbzArray, 0, 100, 15), (self.meshgrid_lat, self.meshgrid_lon, self.dbzArray, 15, 100, 28), (self.meshgrid_lat, self.meshgrid_lon, self.dbzArray, 28, 100, 42), (self.meshgrid_lat, self.meshgrid_lon, self.dbzArray, 42, 100, 55), (self.meshgrid_lat, self.meshgrid_lon, self.dbzArray, 55, 100, 75), (self.meshgrid_lat, self.meshgrid_lon, self.dbzArray, 75, 100, 85)] self.dbz_features = [] # 创建进程池,最大工作进程数保持为2 with concurrent.futures.ProcessPoolExecutor(max_workers=2) as executor: # 提交所有任务,获取Future对象列表 futures = [executor.submit(dbz_processing, *args) for args in arguments] for future in concurrent.futures.as_completed(futures): try: # 给单个任务设置60秒超时 result = future.result(timeout=60) self.dbz_features.append(result) except concurrent.futures.TimeoutError: # 超时处理:可以记录日志或标记结果为None print("任务执行超时,已终止") self.dbz_features.append(None)
方案二:基于原multiprocessing.Pool修改
如果不想切换模块,可以用apply_async逐个提交任务,再通过get(timeout)实现超时控制:
import multiprocessing as mp arguments = [(self.meshgrid_lat, self.meshgrid_lon, self.dbzArray, 0, 100, 15), (self.meshgrid_lat, self.meshgrid_lon, self.dbzArray, 15, 100, 28), (self.meshgrid_lat, self.meshgrid_lon, self.dbzArray, 28, 100, 42), (self.meshgrid_lat, self.meshgrid_lon, self.dbzArray, 42, 100, 55), (self.meshgrid_lat, self.meshgrid_lon, self.dbzArray, 55, 100, 75), (self.meshgrid_lat, self.meshgrid_lon, self.dbzArray, 75, 100, 85)] self.dbz_features = [] pool = mp.Pool(processes=2) # 逐个提交任务,保存异步结果对象 async_results = [pool.apply_async(dbz_processing, args=args) for args in arguments] for res in async_results: try: # 等待60秒获取结果,超时抛出TimeoutError result = res.get(timeout=60) self.dbz_features.append(result) except mp.TimeoutError: print("任务执行超时") self.dbz_features.append(None) pool.close() pool.join()
方案三:手动创建单个进程(支持强制终止超时进程)
如果必须强制终止超时的进程(避免进程占用资源),可以直接用Process类手动管理每个任务:
import multiprocessing as mp def run_with_timeout(func, args, timeout): """包装带超时的任务执行""" result_queue = mp.Queue() # 创建子进程执行任务 proc = mp.Process(target=lambda q: q.put(func(*args)), args=(result_queue,)) proc.start() # 等待超时时间 proc.join(timeout) if proc.is_alive(): # 进程仍在运行,强制终止 proc.terminate() proc.join() print(f"任务超时,已终止进程 {proc.pid}") return None else: # 获取任务结果 try: return result_queue.get_nowait() except mp.queues.Empty: return None # 执行所有任务 self.dbz_features = [] for args in arguments: task_result = run_with_timeout(dbz_processing, args, 60) self.dbz_features.append(task_result)
注意事项
- 方案一和方案二中,进程池会复用进程,超时的任务可能还会在后台运行(直到进程处理完当前任务),如果必须立即终止,优先选方案三。
- 终止进程可能导致资源泄漏(比如未关闭的文件、网络连接),如果
dbz_processing涉及资源操作,建议在函数内部添加异常处理或资源释放逻辑。
内容的提问来源于stack exchange,提问作者Woodz
相关产品推荐
相关产品推荐

