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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 07:35:25