多进程编程:如何基于条件触发回调函数?
解决方案:按批次顺序多进程处理,触发条件立即终止
我参考了Stack Overflow上的"KILLING IT"实现思路,针对你的需求(按25条一组顺序处理、禁止异步乱序、触发条件就终止所有进程)调整了代码,已经测试验证可行,修改后的代码如下:
import random from time import sleep import multiprocessing as mp def worker(i, data_item): print(f"{i} started") # 替换为你的实际计算逻辑,Calcs需提前定义 x = data_item * Calcs # ListOfDataRow请根据实际业务逻辑赋值当前处理的数据行 if x > 0.95: return data_item, True else: return data_item, False # 仅在主进程运行的回调函数,负责触发全局终止 def quit(result): if result[1] == True: pool.terminate() # 立即终止所有进程池工作进程 if __name__ == "__main__": # 假设ListOfData是已定义的原始数据集 batch_size = 25 total_batches = len(ListOfData) // batch_size pool = mp.Pool() should_stop = False for batch_idx in range(total_batches): if should_stop: break # 截取当前批次的数据集 start = batch_idx * batch_size end = start + batch_size current_batch = ListOfData[start:end] # 提交当前批次任务并绑定回调 results = [] for idx, item in enumerate(current_batch): task = pool.apply_async(worker, args=(start + idx, item), callback=quit) results.append(task) # 等待当前批次任务完成,检查是否触发终止条件 for task in results: res = task.get() if res[1]: should_stop = True break # 确保进程池正确关闭 if not pool._closed: pool.close() pool.join()
关键细节说明:
- 严格顺序处理:按25条一组的批次依次提交任务,必须等当前批次所有任务处理完成,且未触发终止条件时,才会处理下一批次,彻底避免异步乱序返回的问题
- 即时终止机制:通过回调函数
quit监听每一个任务的结果,一旦有任务返回True,立刻终止所有工作进程,同时标记停止后续批次的处理 - 语法与逻辑优化:修复了原代码中的语法错误(Python3的
print格式、循环冒号缺失、变量名拼写不一致等),优化了批次截取逻辑,让代码更健壮
内容的提问来源于stack exchange,提问作者sethmac86
相关产品推荐
相关产品推荐

