基于Python编排多ksh脚本:并行执行、依赖调度与故障恢复需求
Python编排KSH脚本的流程优化方案
核心设计思路
- 用标记文件追踪每个KSH脚本的执行状态:成功完成则生成对应标记(如
script_1.done),失败则不生成 - 并行执行通过Python的
multiprocessing模块实现,串行依赖通过等待前序标记文件存在来触发 - 决策逻辑通过读取前序脚本的输出文件,结合条件判断决定后续执行路径
- 故障恢复时,先扫描所有标记文件,跳过已完成的脚本,从第一个未生成标记的后续节点开始执行
关键模块实现
1. 标记文件管理
- 生成标记:脚本执行成功后,创建以脚本名为前缀的
.done文件 - 检查标记:执行前先检查对应标记是否存在,存在则跳过该脚本
- 清理标记(可选):如需重新执行全流程,可批量删除所有
.done文件
2. 脚本执行器
- 同步执行:用
subprocess.run()执行单个脚本,捕获返回码判断是否成功 - 并行执行:用
multiprocessing.Pool提交多个脚本任务,等待所有并行任务完成后再执行后续串行脚本
3. 依赖与决策处理
- 串行依赖:执行脚本前循环等待前序脚本的标记文件出现,超时可抛出异常
- 决策逻辑:读取前序脚本输出的结果文件(如
script_3.output),根据内容判断执行script_5还是其他分支脚本
4. 故障处理
- 脚本执行失败时,立即终止所有后续任务,不生成当前脚本的标记文件
- 故障恢复时,遍历脚本执行序列,找到第一个未生成标记的脚本,从该节点开始执行后续流程
完整示例脚本
import os import subprocess import multiprocessing from time import sleep # 配置:脚本列表、依赖关系、决策规则,按实际流程图修改 CONFIG = { "scripts": [ {"name": "shell_script_1.ksh", "deps": [], "parallel": False}, {"name": "shell_script_2.ksh", "deps": [], "parallel": True}, {"name": "shell_script_3.ksh", "deps": [], "parallel": True}, {"name": "shell_script_4.ksh", "deps": ["shell_script_1.ksh"], "parallel": False}, {"name": "shell_script_5.ksh", "deps": ["shell_script_2.ksh", "shell_script_3.ksh"], "parallel": False}, {"name": "shell_script_6.ksh", "deps": ["shell_script_4.ksh"], "parallel": False, "decision": {"from": "shell_script_4.output", "condition": lambda x: x.strip() == "SUCCESS"}}, # 其余7个脚本按此格式补充 ], "marker_dir": "./script_markers", "output_dir": "./script_outputs" } def init_dirs(): """初始化标记文件和输出文件目录""" os.makedirs(CONFIG["marker_dir"], exist_ok=True) os.makedirs(CONFIG["output_dir"], exist_ok=True) def get_marker_path(script_name): """获取脚本对应的标记文件路径""" return os.path.join(CONFIG["marker_dir"], f"{os.path.splitext(script_name)[0]}.done") def get_output_path(script_name): """获取脚本输出文件路径""" return os.path.join(CONFIG["output_dir"], f"{os.path.splitext(script_name)[0]}.output") def run_script(script_name): """执行单个KSH脚本,返回执行结果""" marker_path = get_marker_path(script_name) if os.path.exists(marker_path): print(f"跳过已完成脚本: {script_name}") return True # 检查依赖是否完成 script_config = next(s for s in CONFIG["scripts"] if s["name"] == script_name) for dep in script_config["deps"]: dep_marker = get_marker_path(dep) while not os.path.exists(dep_marker): print(f"等待依赖脚本完成: {dep}") sleep(5) # 执行脚本 output_path = get_output_path(script_name) try: with open(output_path, "w") as f: subprocess.run( ["ksh", script_name], stdout=f, stderr=subprocess.STDOUT, check=True ) # 执行成功,生成标记文件 open(marker_path, "w").close() print(f"脚本执行成功: {script_name}") return True except subprocess.CalledProcessError as e: print(f"脚本执行失败: {script_name},返回码: {e.returncode}") return False def run_parallel_scripts(script_names): """并行执行多个脚本""" with multiprocessing.Pool(processes=len(script_names)) as pool: results = pool.map(run_script, script_names) # 并行任务有失败则终止流程 if not all(results): raise Exception("并行脚本执行失败,终止流程") def run_decision_script(script_config): """根据决策逻辑执行脚本""" script_name = script_config["name"] decision = script_config["decision"] output_path = get_output_path(decision["from"]) # 读取前序脚本输出 with open(output_path, "r") as f: output_content = f.read() if decision["condition"](output_content): print(f"满足决策条件,执行脚本: {script_name}") if not run_script(script_name): raise Exception(f"脚本{script_name}执行失败,终止流程") else: print(f"不满足决策条件,跳过脚本: {script_name}") def main(): init_dirs() # 按流程分组执行,可根据实际流程图调整顺序 # 第一组:无依赖的并行脚本 parallel_group_1 = [s["name"] for s in CONFIG["scripts"] if s["parallel"] and not s["deps"]] if parallel_group_1: run_parallel_scripts(parallel_group_1) # 第二组:无依赖的串行脚本 serial_group_1 = [s["name"] for s in CONFIG["scripts"] if not s["parallel"] and not s["deps"]] for script in serial_group_1: if not run_script(script): print("串行脚本执行失败,终止流程") return # 第三组:依赖script_1的串行脚本 serial_group_2 = [s["name"] for s in CONFIG["scripts"] if s["deps"] == ["shell_script_1.ksh"]] for script in serial_group_2: if not run_script(script): print("串行脚本执行失败,终止流程") return # 第四组:决策分支脚本 decision_group = [s for s in CONFIG["scripts"] if "decision" in s] for script_config in decision_group: run_decision_script(script_config) # 第五组:依赖script_2和script_3的串行脚本 serial_group_3 = [s["name"] for s in CONFIG["scripts"] if set(s["deps"]) == {"shell_script_2.ksh", "shell_script_3.ksh"}] for script in serial_group_3: if not run_script(script): print("串行脚本执行失败,终止流程") return # 其余脚本按实际依赖关系补充执行逻辑 print("所有脚本执行完成") if __name__ == "__main__": try: main() except Exception as e: print(f"流程终止,原因: {str(e)}")
使用说明
- 按照实际流程图修改
CONFIG中的脚本列表、依赖关系和决策规则 - 标记文件默认存储在
./script_markers目录,输出文件存储在./script_outputs目录,可根据需要修改路径 - 故障修复后,直接重新运行脚本,程序会自动跳过已生成标记的脚本,从第一个未完成的节点开始执行
内容的提问来源于stack exchange,提问作者Satyabrat Sahoo
相关产品推荐
相关产品推荐

