Python多进程池提交SGE qsub任务大规模数据失败问题求助
问题分析与解决方案
核心问题
- SGE任务未等待完成:
qsub命令默认仅提交任务到队列就返回,代码提交后立刻检查CSV文件,此时任务可能还在排队或运行中,导致文件不存在。小数据量时SGE队列无积压,任务快速完成,所以未暴露问题;数据量大时队列阻塞,就会触发报错。 - 进程池使用不当:
multiprocessing.Pool(1000)设置的并发数过高,会耗尽系统资源;手动调用close()和join()时机错误,导致"Pool not running"报错。
解决方案
1. 让Python等待SGE任务完成
修改qsub命令,添加-sync y参数,强制qsub等待任务执行完成后再返回:
cmd = ["qsub","-S", "/usr/bin/bash", "-sync", "y", "run_list_segment.sh", "run_list_segment.py", temp_json]
注:如果你的SGE版本不支持
-sync y,可以通过qsub返回的任务ID,循环调用qstat检查状态,直到任务结束:# 提交任务并获取任务ID run_command = subprocess.run(cmd, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True) run_command.check_returncode() job_id = run_command.stdout.strip().split('.')[0] # 循环检查任务状态 while True: qstat_cmd = ["qstat", job_id] result = subprocess.run(qstat_cmd, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True) # 任务完成时qstat会返回非0状态码,或输出中无该任务ID if result.returncode != 0: break time.sleep(5) # 每5秒检查一次
2. 优化进程池使用
- 降低并发数:1000的进程数远超系统承载能力,建议根据CPU核心数或SGE允许的并发任务数调整(比如20-50)。
- 使用
with语句管理进程池,自动处理关闭和等待,避免手动调用close()/join()的错误:
list_argument = [{dict1}, {dict2}, ...] # 调整合理的进程池大小 pool_size = 20 with multiprocessing.Pool(pool_size) as pool: output_data = pool.map(process_list, list_argument) # with块结束后自动关闭并join进程池,无需手动操作
3. 额外优化建议
- 临时文件唯一化:确保
temp_json的文件名全局唯一,避免多进程写入冲突(比如用uuid生成文件名)。 - 添加重试机制:即使使用
-sync y,磁盘IO可能有延迟,检查CSV文件时可增加重试:
import time pandas_csv = process_dict['name'].csv max_retries = 3 retry_delay = 2 for _ in range(max_retries): if os.path.isfile(pandas_csv): df = pandas.read_csv(pandas_csv, dtype={'CHROM': 'str'}) break time.sleep(retry_delay) else: raise ValueError(f"{pandas_csv} wasn't created after {max_retries} retries")
- 清理临时文件:任务完成后删除
temp_json,避免磁盘空间浪费:
os.unlink(temp_json)
内容的提问来源于stack exchange,提问作者Norman Kuo
相关产品推荐
相关产品推荐

