【发布时间】:2021-03-17 21:56:33
【问题描述】:
我正在构建一个简单的示例来了解 dask Distributed 如何在 HPC 集群上分发 python 脚本。 分发方法是一种基本操作,即在磁盘上写入文件。从命令行运行时,此脚本运行良好 ($python simple-function.py)
import os
import argparse
import time
def inc(x):
time.sleep(1)
return x + 1
def get_args():
"""
Get args
"""
parser = argparse.ArgumentParser()
parser.add_argument("x", help="Value")
args = parser.parse_args()
x = args.x
return int(x)
if __name__ == "__main__":
x = get_args()
print("{0} + 1 = {1}".format(x, inc(x)))
with open("results.txt", 'w') as file:
file.write(str(x) + '\n')
现在,我创建了另一个将分发此代码的 python 脚本。 这个想法是使用 subprocess 和 client.map 或 client.submit 来启动上述脚本的多个实例。 我遇到的问题是使用以下任何方法(client.map、client.submit 然后收集或计算或 .result())时未写入输出 .txt 文件。 也许我没有使用正确的方法?
import os
import time
import subprocess
import yaml
import argparse
import dask
from dask.distributed import Client
from dask_jobqueue import PBSCluster
def create_cluster():
cluster = PBSCluster(
cores=4,
memory="20GB",
interface="ib0",
queue="qdev",
processes=4,
nanny=True,
walltime="12:00:00",
shebang="#!/bin/bash",
local_directory="$TMPDIR"
)
cluster.scale(4)
time.sleep(10) # Wait for workers
return cluster
def get_args():
parser = argparse.ArgumentParser()
parser.add_argument("script", help="Script to distribute")
parser.add_argument("nodes", type=int, help="Number of nodes")
args = parser.parse_args()
script = args.script
n_nodes = args.nodes
return script, n_nodes
def close_all(client, cluster):
client.close()
cluster.close()
def methode(script, x):
subprocess.run(["python",
script,
x])
return None
if __name__ == "__main__":
cluster = create_cluster()
client = Client(cluster)
time.sleep(1)
script, n_nodes = get_args() #Get arguments
#With client.submit
futures = []
for n, o in enumerate(range(10)):
futures.append(client.submit(methode, *[script, str(o)], priority=-n))
[f.result() for f in futures]
#Or client.map
L = client.map(methode, *[script, str(range(10))])
client.compute(L)
client.gather(L)
time.sleep(20)
close_all(client, cluster)
附带说明,如果我执行以下代码: dask.compute(方法(args))
然后,.txt 输出文件被写入。
似乎只有不同的客户端方法不起作用
【问题讨论】:
标签: python dask dask-distributed