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

非阻塞线程池:如何实现结果就绪即通知主线程?

实现单任务完成即触发回调的多进程池方案

原始代码

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)

核心说明

  1. 两个方案均使用apply_async实现单任务回调能力,替代批量回调的map_async
  2. 进程安全的mp.Value和mp.Queue确保多进程间数据访问无冲突
  3. 主线程通过短休眠轮询实现非阻塞等待,避免占用过多CPU资源
  4. calc的动态变量通过apply_async的args参数传入,每次loop调用时更新即可

内容的提问来源于stack exchange,提问作者MirceaKitsune

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 06:23:12