非阻塞线程池:如何实现结果就绪即通知主线程?
实现单任务完成即触发回调的多进程池方案
原始代码
import multiprocessing as mp pool = mp.Pool() def calc(i): return i * 2 def done(results): for result in results: print(result) def loop(): pool.map_async(calc, [0, 1, 2, 3], callback = done) while True: loop()
需求说明
当前配置下,map_async会在所有任务完成后统一触发done回调,传入完整结果列表。期望实现:
- 单个任务完成后立即触发回调或更新主线程可访问的结果容器
- 批次任务全部完成后停止回调,等待主循环再次触发后重复流程
- 主线程可随时获取已完成结果,不能使用阻塞式的
result.get() calc函数依赖的动态变量需在每次主循环中由主线程更新,[0,1,2,3]为固定任务列表
方案一:单任务独立回调+进程安全计数器
用apply_async替代map_async,为每个任务单独设置回调,同时用进程安全计数器跟踪批次完成状态。
import multiprocessing as mp import time pool = mp.Pool(4) # 进程数与任务数匹配 batch_counter = mp.Value('i', 0) # 进程安全计数器 TASK_LIST = [0, 1, 2, 3] # 固定任务列表 def calc(i, current_var): # 使用主线程传入的动态变量current_var time.sleep(0.5) # 模拟计算耗时 return i * 2 + current_var def single_done(result): # 单个任务完成后的回调逻辑 print(f"单任务完成:{result}") # 更新计数器 with batch_counter.get_lock(): batch_counter.value += 1 def loop(): # 重置批次计数器 with batch_counter.get_lock(): batch_counter.value = 0 # 主线程本次循环更新的动态变量 current_var = int(time.time()) % 10 # 逐个提交任务并绑定单任务回调 for i in TASK_LIST: pool.apply_async(calc, args=(i, current_var), callback=single_done) # 非阻塞轮询,等待当前批次全部完成 while True: with batch_counter.get_lock(): completed = batch_counter.value if completed == len(TASK_LIST): print("当前批次任务全部完成,等待下一轮循环") break time.sleep(0.1) # 短休眠避免CPU空转 if __name__ == "__main__": while True: loop() time.sleep(2) # 模拟主循环间隔
方案二:进程安全队列收集结果
用multiprocessing.Queue作为结果容器,任务完成后将结果存入队列,主线程可随时非阻塞获取已完成结果。
import multiprocessing as mp import time pool = mp.Pool(4) result_queue = mp.Queue() TASK_LIST = [0, 1, 2, 3] def calc(i, current_var): time.sleep(0.5) return i * 2 + current_var def single_done(result): # 将结果存入进程安全队列 result_queue.put(result) def loop(): current_var = int(time.time()) % 10 # 提交所有任务 for i in TASK_LIST: pool.apply_async(calc, args=(i, current_var), callback=single_done) # 主线程非阻塞获取结果,直到批次完成 completed_count = 0 while completed_count < len(TASK_LIST): try: # 非阻塞取结果,超时时间短避免阻塞 result = result_queue.get(block=False) print(f"获取到结果:{result}") completed_count += 1 # 此处可直接处理结果后丢弃 except mp.queues.Empty: time.sleep(0.1) print("当前批次任务全部完成") if __name__ == "__main__": while True: loop() time.sleep(2)
核心说明
- 两个方案均使用
apply_async实现单任务回调能力,替代批量回调的map_async - 进程安全的
mp.Value和mp.Queue确保多进程间数据访问无冲突 - 主线程通过短休眠轮询实现非阻塞等待,避免占用过多CPU资源
calc的动态变量通过apply_async的args参数传入,每次loop调用时更新即可
内容的提问来源于stack exchange,提问作者MirceaKitsune
相关产品推荐
相关产品推荐

