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,示例映射:Date JobRunId 2022-10-01 56728389 2022-10-02 56728390 2022-10-03 56728391 2022-10-04 56728392 - 现有问题:无法在第一批3个任务完成后输出对应JobRunId列表并执行状态校验,当前仅能在所有任务完成后拿到全量列表;需要实现分批执行逻辑:第一批任务完成后,校验每个JobRunId的状态(RUNNING/SUCCESS),再执行第二批任务,校验操作需在
run_job的time.sleep(60)后调用
问题根源
多进程模式下,子进程拥有独立内存空间,全局变量listJobRunIds在每个子进程中都是副本,子进程的append操作不会同步到主进程,导致主进程无法实时获取中间结果;同时pool.map会一次性提交所有任务,无法实现分批执行后的校验逻辑。
解决方案
核心思路
- 任务分批拆分:将任务列表按进程数(3个)分成若干批次
- 分批执行+校验:每批任务执行完成后,获取该批的JobRunId列表,执行状态校验,再启动下一批任务
- 修正
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
相关产品推荐
相关产品推荐

