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

如何在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节点可以直接访问该路径下的文件,无需额外传输。这是大文件场景下最推荐的方案。

步骤:

  1. 在AzureML工作区中创建并挂载共享存储(如Azure File Share)到集群的所有节点,挂载路径例如/mnt/azureml-share。
  2. 修改远程函数,将文件写入共享路径:

修改后的代码示例:

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 00:13:16