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

Python多进程分批执行:如何获取每批JobRunId并校验任务状态

Python多进程分批执行并校验任务状态的实现方案

现有代码

import multiprocessing
import time

# 假设START_JOB_RUN和execute是已定义的变量/函数
START_JOB_RUN = "some_command {dataset_date}"
def execute(command):
    # 模拟返回JobRunId
    return "56728389"

listJobRunIds = []
    
def run_job(dataset_date):
    command = START_JOB_RUN.format(dataset_date=dataset_date)
    jobRunId =  execute(command=command)
    print("jobRunId:"+jobRunId)
    listJobRunIds.append(jobRunId)
    print(listJobRunIds)
    time.sleep(60)
    return jobRunId

pool = multiprocessing.Pool()
pool = multiprocessing.Pool(processes=3)

listJobRunDates = ['2022-10-01','2022-10-02','2022-10-03','2022-10-04']
jobRunId = pool.map(run_job, listJobRunDates)

print(jobRunId)

需求说明

  • 仅允许3个进程并行执行
  • execute函数执行任务并返回JobRunId,示例映射:
    DateJobRunId
    2022-10-0156728389
    2022-10-0256728390
    2022-10-0356728391
    2022-10-0456728392
  • 现有问题:无法在第一批3个任务完成后输出对应JobRunId列表并执行状态校验,当前仅能在所有任务完成后拿到全量列表;需要实现分批执行逻辑:第一批任务完成后,校验每个JobRunId的状态(RUNNING/SUCCESS),再执行第二批任务,校验操作需在run_job的time.sleep(60)后调用

问题根源

多进程模式下,子进程拥有独立内存空间,全局变量listJobRunIds在每个子进程中都是副本,子进程的append操作不会同步到主进程,导致主进程无法实时获取中间结果;同时pool.map会一次性提交所有任务,无法实现分批执行后的校验逻辑。

解决方案

核心思路

  1. 任务分批拆分:将任务列表按进程数(3个)分成若干批次
  2. 分批执行+校验:每批任务执行完成后,获取该批的JobRunId列表,执行状态校验,再启动下一批任务
  3. 修正run_job函数:在time.sleep(60)后加入状态校验逻辑

修改后的完整代码

import multiprocessing
import time
from itertools import islice

# 假设START_JOB_RUN、execute、check_job_status是已定义的变量/函数
START_JOB_RUN = "some_command {dataset_date}"
def execute(command):
    # 模拟根据日期返回对应JobRunId
    date_to_id = {
        '2022-10-01': '56728389',
        '2022-10-02': '56728390',
        '2022-10-03': '56728391',
        '2022-10-04': '56728392'
    }
    return date_to_id.get(command.split()[-1], "unknown_id")

def check_job_status(job_run_id):
    # 模拟校验任务状态,需根据实际业务实现
    print(f"Checking status for JobRunId: {job_run_id} -> SUCCESS")
    return "SUCCESS"

def run_job(dataset_date):
    command = START_JOB_RUN.format(dataset_date=dataset_date)
    job_run_id = execute(command=command)
    print(f"jobRunId: {job_run_id}")
    time.sleep(60)
    # sleep后执行状态校验
    check_job_status(job_run_id)
    return job_run_id

def batch_process(jobs, batch_size=3):
    pool = multiprocessing.Pool(processes=batch_size)
    # 分批迭代任务列表
    it = iter(jobs)
    while True:
        batch = list(islice(it, batch_size))
        if not batch:
            break
        # 执行当前批次任务
        print(f"\nStarting batch: {batch}")
        batch_job_ids = pool.map(run_job, batch)
        # 输出当前批次的JobRunId列表
        print(f"Batch completed, JobRunIds: {batch_job_ids}")
    pool.close()
    pool.join()

if __name__ == "__main__":
    listJobRunDates = ['2022-10-01','2022-10-02','2022-10-03','2022-10-04']
    batch_process(listJobRunDates)

关键说明

  • 用itertools.islice实现任务列表的分批拆分,每次取3个任务作为一批
  • batch_process函数负责分批提交任务,每批执行完成后直接获取该批的JobRunId列表并输出
  • run_job函数中在time.sleep(60)后调用check_job_status完成单个任务的状态校验
  • 添加if __name__ == "__main__":确保多进程代码在Windows系统下正常运行(Windows多进程的必要条件)
  • 主进程通过pool.map的返回值获取批次任务的JobRunId,无需依赖全局变量,彻底解决多进程内存隔离导致的变量不同步问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 03:40:48