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

多进程编程:如何基于条件触发回调函数?

解决方案:按批次顺序多进程处理,触发条件立即终止

我参考了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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:38:54