如何在AzureML中使用Ray将Worker Node文件传输至Head Node
在AzureML的Ray集群中将Worker节点文件传输到Head节点的解决方案
针对你在AzureML上使用Ray时,Worker节点生成的文件无法同步到Head节点的问题,以下是几种可行的解决方法:
方法一:直接返回文件内容到Head节点(适合小文件)
修改远程函数,将生成的文件内容作为返回值的一部分,在Head节点接收后写入本地文件。这种方式无需额外存储,适合数据量较小的场景。
修改后的代码示例:
import os import sys import ray import time from pathlib import Path from ray_on_aml.core import Ray_On_AML import numpy as np @ray.remote def run_script(size): value = np.square(size) # 先将内容存入变量,同时可选择在Worker节点保留文件 content = str(value) np.savetxt(f"output_{size}.txt", value) return {"matrix_size": size, "file_content": content} if __name__ == "__main__": BASE_DIR = Path(__file__).resolve().parent sys.path.append(BASE_DIR.as_posix()) ray_on_aml = Ray_On_AML() ray = ray_on_aml.getRay() if ray: ray.init(address="auto") print("head node detected") print(ray.cluster_resources()) start = time.time() mat_size = [10, 20, 30, 40, 50, 60, 70, 80, 90, 100] script_handles = [] for i in range(5): subset = mat_size[:2] mat_size = mat_size[2:] script_handles.extend([run_script.remote(size) for size in subset]) results = ray.get(script_handles) # 在Head节点写入文件 for res in results: size = res["matrix_size"] content = res["file_content"] with open(f"output_{size}.txt", "w") as f: f.write(content) end = time.time() dur = end - start print(f"tests took {dur} seconds") print("The current files in the directory are", os.listdir(os.getcwd())) ray.shutdown() else: print("in worker node, do nothing")
方法二:使用Ray分布式对象存储(适合中等大小文件)
通过ray.put将文件内容存入Ray的分布式对象存储,Head节点通过对象引用获取内容后写入文件。Ray会自动处理跨节点的数据传输,比直接返回更高效。
修改后的代码示例:
import os import sys import ray import time from pathlib import Path from ray_on_aml.core import Ray_On_AML import numpy as np @ray.remote def run_script(size): value = np.square(size) np.savetxt(f"output_{size}.txt", value) # 读取文件内容并存入Ray对象存储,返回对象引用 with open(f"output_{size}.txt", "r") as f: content = f.read() content_ref = ray.put(content) return {"matrix_size": size, "content_ref": content_ref} if __name__ == "__main__": BASE_DIR = Path(__file__).resolve().parent sys.path.append(BASE_DIR.as_posix()) ray_on_aml = Ray_On_AML() ray = ray_on_aml.getRay() if ray: ray.init(address="auto") print("head node detected") print(ray.cluster_resources()) start = time.time() mat_size = [10, 20, 30, 40, 50, 60, 70, 80, 90, 100] script_handles = [] for i in range(5): subset = mat_size[:2] mat_size = mat_size[2:] script_handles.extend([run_script.remote(size) for size in subset]) results = ray.get(script_handles) # 在Head节点通过对象引用获取内容并写入文件 for res in results: size = res["matrix_size"] content = ray.get(res["content_ref"]) with open(f"output_{size}.txt", "w") as f: f.write(content) end = time.time() dur = end - start print(f"tests took {dur} seconds") print("The current files in the directory are", os.listdir(os.getcwd())) ray.shutdown() else: print("in worker node, do nothing")
方法三:使用AzureML共享存储(适合大文件)
AzureML支持挂载Blob存储或File Share到集群的所有节点,Worker节点直接将文件写入共享路径,Head节点可以直接访问该路径下的文件,无需额外传输。这是大文件场景下最推荐的方案。
步骤:
- 在AzureML工作区中创建并挂载共享存储(如Azure File Share)到集群的所有节点,挂载路径例如
/mnt/azureml-share。 - 修改远程函数,将文件写入共享路径:
修改后的代码示例:
import os import sys import ray import time from pathlib import Path from ray_on_aml.core import Ray_On_AML import numpy as np # 定义共享存储路径(根据实际挂载路径调整) SHARED_DIR = "/mnt/azureml-share" os.makedirs(SHARED_DIR, exist_ok=True) @ray.remote def run_script(size): value = np.square(size) # 写入共享存储路径 file_path = os.path.join(SHARED_DIR, f"output_{size}.txt") np.savetxt(file_path, value) return {"matrix_size": size} if __name__ == "__main__": BASE_DIR = Path(__file__).resolve().parent sys.path.append(BASE_DIR.as_posix()) ray_on_aml = Ray_On_AML() ray = ray_on_aml.getRay() if ray: ray.init(address="auto") print("head node detected") print(ray.cluster_resources()) start = time.time() mat_size = [10, 20, 30, 40, 50, 60, 70, 80, 90, 100] script_handles = [] for i in range(5): subset = mat_size[:2] mat_size = mat_size[2:] script_handles.extend([run_script.remote(size) for size in subset]) results = ray.get(script_handles) end = time.time() dur = end - start print(f"tests took {dur} seconds") # 直接访问共享路径下的文件 print("The current files in shared directory are", os.listdir(SHARED_DIR)) ray.shutdown() else: print("in worker node, do nothing")
内容的提问来源于stack exchange,提问作者bharat
相关产品推荐
相关产品推荐

