如何使用已部署的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
相关产品推荐
相关产品推荐

