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

如何使用已部署的Dask Distributed集群并行运行Shell脚本?

嘿,你的思路完全没问题,用Dask Distributed来并行跑shell脚本是个很棒的选择,我之前也做过类似的事情,给你完善下代码和几个关键注意点,让整个流程更健壮:

用Dask Distributed并行运行Shell脚本

先修正你现有代码的小问题

你原来的代码里列表推导式没写完,应该遍历commands列表而不是传入Client对象,先把基础版本跑通:

import dask
import distributed
import os
from glob import glob

# 批量获取文件夹下所有.sh脚本(不用手动一个个写命令)
commands = glob("./your_script_folder/*.sh")

# 用delayed包装执行命令的函数
@dask.delayed
def run_script(cmd):
    # os.system返回0表示执行成功,非0为失败
    return os.system(cmd)

# 连接到已搭建的集群
client = distributed.Client('my_server:8786')

# 提交所有任务到集群
futures = client.compute([run_script(cmd) for cmd in commands])

# 等待所有任务完成并收集结果
results = client.gather(futures)

# 可以打印每个脚本的执行结果
for cmd, code in zip(commands, results):
    print(f"脚本 {cmd} 执行返回码: {code}")

更健壮的优化方案

os.system虽然简单,但没法捕获脚本的输出和错误日志,排查问题很麻烦。推荐用subprocess模块替代,这样能完整获取执行过程中的信息:

import dask
import distributed
import subprocess
from glob import glob

@dask.delayed
def run_script(cmd):
    try:
        # 捕获标准输出、错误,设置超时时间(根据你的脚本时长调整)
        result = subprocess.run(
            cmd,
            shell=True,
            check=True,
            stdout=subprocess.PIPE,
            stderr=subprocess.PIPE,
            text=True,
            timeout=3600  # 1小时超时,防止脚本挂起
        )
        return {
            "cmd": cmd,
            "returncode": result.returncode,
            "stdout": result.stdout,
            "stderr": ""
        }
    except subprocess.CalledProcessError as e:
        # 脚本执行返回非0码时触发
        return {
            "cmd": cmd,
            "returncode": e.returncode,
            "stdout": e.stdout,
            "stderr": e.stderr
        }
    except subprocess.TimeoutExpired as e:
        # 脚本超时触发
        return {
            "cmd": cmd,
            "returncode": -1,
            "stdout": e.stdout,
            "stderr": "脚本执行超时"
        }

# 批量获取脚本命令
commands = glob("./your_script_folder/*.sh")

# 连接集群
client = distributed.Client('my_server:8786')

# 用client.map更简洁地批量提交任务
futures = client.map(run_script, commands)

# 收集所有任务结果
results = client.gather(futures)

# 批量处理结果,比如筛选出执行失败的脚本
failed_jobs = [res for res in results if res["returncode"] != 0]
if failed_jobs:
    print("以下脚本执行失败:")
    for res in failed_jobs:
        print(f"脚本: {res['cmd']}")
        print(f"返回码: {res['returncode']}")
        print(f"错误信息: {res['stderr']}\n")
else:
    print("所有脚本都执行成功啦!")

必须注意的几个关键点

  • 脚本路径可达性:所有Worker节点必须能访问到这些.sh脚本!如果你的脚本在本地机器,Worker是远程节点的话,要么把脚本同步到每个Worker的相同路径,要么用NFS这类共享存储挂载到所有节点的同一个目录下。
  • 执行权限:确保Worker节点的运行用户有脚本的执行权限,提前在Worker上执行chmod +x /path/to/your/scripts/*.sh来添加执行权限。
  • 资源限制:如果你的脚本很占用CPU或内存,可以在提交任务时指定资源要求,避免集群过载:
    # 每个任务占用1核CPU和2GB内存,根据实际情况调整
    futures = client.map(run_script, commands, resources={"cpu": 1, "memory": "2GB"})
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 06:51:37