Python——Ray-多节点集群启动


整体说明

  • Ray 集群的启动有多种方法,本文简单总结这些方法
    • 注意:除了本文介绍的 ray start 方法外,还可以通过 从配置文件启动、Kubernetes 部署 等方法启动集群
  • 核心概念补充:
    • 头节点(Head Node) :集群的主节点,负责管理整个集群的资源、任务调度和元数据存储
    • 工作节点(Worker Node) :通过连接头节点加入集群,提供计算资源(CPU/GPU/内存等)

Python 程序内启动(单节点启动)

  • 场景一:单节点模拟(只有一个 Python 进程)
    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    12
    import ray

    ray.init(num_cpus=4, num_gpus=1, dashboard_host="0.0.0.0")
    print("节点数:", len(ray.nodes())) # 输出 1

    @ray.remote
    def f():
    return 1

    print(ray.get(f.remote()))
    input("按任意键关闭集群...") # 保持运行
    ray.shutdown()

Python 单节点多进程模拟

  • 场景二:外部客户端连接到已存在的头节点(需另开终端)

  • 终端1(头节点)

    1
    2
    3
    4
    import ray
    ray.init(num_cpus=4, port=6380, dashboard_host="0.0.0.0")
    print("Head started at", ray.get_runtime_context().gcs_address)
    input("Keep alive...") # 不要 shutdown
  • 终端2(客户端/所谓的“工作进程”)

    1
    2
    3
    4
    5
    6
    import ray
    # 假设头节点 IP 是 127.0.0.1
    ray.init(address="127.0.0.1:6380")

    # 此时该脚本只是一个客户端,不会增加节点数(仅是作为一个进程启动),可以使用集群资源
    print(ray.cluster_resources())

使用 ray start 命令启动(真正的多节点)

  • ray start 是启动 Ray 集群节点的常用命令
  • 可分别分别使用 ray start 用于初始化头节点(Head Node)和工作节点(Worker Node)
  • ray start 更多命令可参考:docs.ray.io/en/latest/cluster/cli.html#ray-start

启动头节点(头结点机器上运行)

  • 头节点是集群的入口,必须先启动

  • 基本命令格式:

    1
    ray start --head [其他可选参数]
  • 关键参数包括:

    参数 说明 示例
    --head 声明当前节点为头节点(必选) -
    --port 指定 Ray 内部通信端口(默认 6379,若被占用会自动切换) --port=6380
    --dashboard-host 允许外部访问 Ray dashboard 的主机地址(默认仅本地访问) --dashboard-host=0.0.0.0 表示允许外部访问
    --dashboard-port Dashboard 端口(默认 8265) --dashboard-port=8266
    --num-cpus 手动指定该节点可用的 CPU 核心数(默认自动检测) --num-cpus=16
    --num-gpus 手动指定该节点可用的 GPU 数量(默认自动检测) --num-gpus=2
    --memory 限制节点可用内存(单位:字节,如 1000000000 表示 1GB) --memory=8000000000
    --object-store-memory 对象存储的内存上限(默认总内存的 30%) --object-store-memory=2000000000
    --block 启动后阻塞终端(不后台运行,便于调试) -
    --log-dir 指定日志目录(默认 ~/raylogs --log-dir=/path/to/logs
    --min-worker-port 指定当前机器上 Ray 工作进程可以绑定的 最低 端口号(默认选择随机可用端口) --min-worker-port=8080
    --max-worker-port 指定当前机器上 Ray 工作进程可以绑定的 最高 端口号(必须同时设置 --min-worker-port --max-worker-port=9090
    --worker-port-list 精细控制端口号 --worker-port-list=8001,8003,8008
  • 补充说明:

    • --min-worker-port--max-worker-port 这两个参数最主要的价值在于方便进行网络规划和防火墙设置
    • 在多节点集群中,为了安全,通常需要在节点间开放特定范围的端口用于通信
      • 通过固定这个端口范围,你只需在防火墙规则中开放这个范围,而不是开放所有端口或为每个新工作进程动态调整规则
    • 注:--worker-port-list 参数会覆盖 --min-worker-port--max-worker-port 的设置
  • 示例:启动一个允许外部访问 Dashboard、指定 CPU/GPU 资源的头节点:

    1
    2
    3
    4
    5
    6
    ray start --head \
    --port=6379 \
    --dashboard-host=0.0.0.0 \
    --dashboard-port=8265 \
    --num-cpus=12 \
    --num-gpus=1
  • 头节点启动成功后,终端会输出类似以下信息

    • 请记下头节点的地址(如 192.168.1.100:6379),工作节点连接时需要用到
    • 还可以记录 Dashboard,查看集群监控任务状态的地址和端口
  • ray start 命令的其他高级配置

    • 可通过 --runtime-env 指定环境配置文件(如依赖安装、环境变量等):

      1
      ray start --head --runtime-env=runtime_env.yaml
    • 可通过 --redis-password 设置密码,防止未授权节点加入:

      1
      2
      3
      4
      5
      # 头节点
      ray start --head --redis-password='mysecret'

      # 工作节点
      ray start --address=<IP:port> --redis-password='mysecret'
    • 可通过 --log-dir 指定日志目录(默认 ~/raylogs):

      1
      ray start --head --log-dir=/path/to/logs

启动工作节点(在其他机器上启动)

  • 工作节点需通过头节点的地址加入集群,命令格式:

    1
    ray start --address=$head_ip:$head_port [其他可选参数]
  • 关键参数说明:

    • --address:头节点的地址(必填,格式为 $head_ip:$head_port,即头节点启动时输出的地址)
    • 其他参数(如 --num-cpus--num-gpus 等)与头节点相同,用于限制工作节点的资源
      • 比如:上述 --min-worker-port--max-worker-port 等关键参数一样可以用于限制当前 Ray 工作节点的端口使用
  • 示例:连接到 IP 为 192.168.1.100、端口为 6379 的头节点,同时指定工作节点的资源:

    1
    2
    3
    ray start --address='192.168.1.100:6379' \
    --num-cpus=8 \
    --num-gpus=0

使用 ray 命令验证集群状态(头节点和工作节点均可)

  • 可在任意节点执行下面脚本查看节点列表

    1
    ray status
    • 输出会显示集群中的所有节点及资源使用情况
  • 通过头节点的 Dashboard http://<头节点IP>:8265(或其他指定端口,详情会在启动时输出日志) 查看集群监控、任务状态等

  • 通过执行 Python 脚本 也可以输出对应的集群状态

    1
    2
    3
    import ray
    ray.init(address="auto") # 自动发现本地集群
    print(ray.nodes()) # 查看所有节点信息
    • 建议通过 ray job submit 提交这个(后文会讲到)

      1
      ray job submit --address="http://<HEAD_NODE_IP>:8265" -- python my_check.py
    • 注:这段代码可能的执行效果:

      • 如果在头节点运行:ray.init(address="auto") 会连接到本地的 Ray 集群,打印节点信息
      • 如果在工作节点运行:同样会连接到集群(因为工作节点也有 Ray 进程),打印所有节点信息
      • 如果在外部机器运行(未启动 Ray 进程):address="auto" 会尝试连接本地默认端口,但可能失败,除非该机器能通过网络连接到头节点

使用 ray 命令停止某个节点(头节点和工作节点均可)

  • 停止单个节点(包括工作节点或头节点都可以):

    1
    ray stop
  • 特别注意:若头节点停止,整个集群会自动解散


通过 ray job submit 向已经启动的 Ray 集群提交任务

  • ray job submit 命令用于将作业提交到 Ray 集群
  • 可以在任何能访问头节点 Dashboard 的机器(包括你的本地开发机)上提交任务,不一定要在头节点上
    • --address 可指定 Ray 集群的地址,通常是集群头节点的地址和端口

基本语法

  • 用法说明:

    1
    ray job submit [options] -- <entrypoint> [<entrypoint_args>]
    • [options] 是命令的可选参数
    • -- 之后的 <entrypoint> 是要执行的入口点脚本或命令
      • 可以是 -- python my_script.pybash my_shell.sh
      • <entrypoint_args> 是传递给 <entrypoint> 的参数(my_script.py 等)

常用参数

  • --address:指定 Ray 集群的地址,通常是集群头节点的地址和端口

    • 优先级(从高到低):
      • 1)命令行参数 --address,显式指定,优先级最高,会覆盖所有环境变量
      • 2)环境变量 RAY_API_SERVER_ADDRESS(新版专用)
      • 3)环境变量 RAY_ADDRESS(旧版兼容,仅当前两者未设时生效)
      • 4)默认值 http://127.0.0.1:8265(连接本地集群)
  • --runtime-env:用于指定作业的运行时环境,可以是一个 JSON 格式的字符串或 YAML 文件路径

    • 注意:这个和 ray start 启动时的 runtime-env 不是完全等价的,只影响当前提交的 Python 任务
    • 可以通过该参数指定需要安装的 Python 包,如--runtime-env='{"pip": ["requests"]}'
      • 这会在所有节点上安装包,可用于临时任务添加 pip 包(一般建议提前装好,可以一直复用,不要用这个方法)
    • 还可以用于排除一些对运行没有用但是比较大的子文件夹(注意书写格式)
    • 示例(from verl):
      1
      2
      3
      4
      5
      6
      7
      working_dir: ./
      excludes: ["/.git/", "/wandb_log/", "/local_data/"]
      env_vars:
      TORCH_NCCL_AVOID_RECORD_STREAMS: "1"
      CUDA_DEVICE_MAX_CONNECTIONS: "1"
      HCCL_HOST_SOCKET_PORT_RANGE: "auto"
      HCCL_NPU_SOCKET_PORT_RANGE: "auto"
  • --working-dir:指定作业的工作目录

    • 该目录下的文件会被同步到集群节点上,默认为当前目录
    • 启动脚本会有一些文件依赖,这里上传所有文件可以保证本地能访问的文件,集群也能访问
    • 需要排除的文件在 --runtime-env 参数中配置,被排除的文件夹不用上传到各个节点上
  • --no-wait:提交作业后不等待作业完成,立即返回

    • 如果不指定该参数,命令会等待作业完成,并输出作业的日志和结果
  • --submission-id:指定作业的提交 ID

    • 如果不指定,Ray 会自动生成一个唯一的 ID

用法示例

  • 假设要提交一个 Python 脚本 my_script.py 到 Ray 集群,指定集群地址为 http://127.0.0.1:8265,工作目录为当前目录:

    1
    RAY_ADDRESS='http://127.0.0.1:8265' ray job submit -- python my_script.py
  • 假设提交一个 Python 脚本 train.py,并传递参数 --epochs 10 --batch-size 32,同时指定运行时环境需要安装 torchnumpy 包:

    1
    ray job submit --address=http://your-ray-cluster-address:8265 --runtime-env='{"pip": ["torch", "numpy"]}' -- python train.py --epochs 10 --batch-size 32
  • 参数较多时的提交示例:

    1
    2
    3
    4
    5
    6
    7
    8
    ray job submit \
    --runtime-env=/mnt/data/runtime-env.yaml \
    --working-dir=/mnt/workspace/ \
    -- \
    python3 -m verl.trainer.main_ppo \
    data.train_files=/mnt/data/train.parquet \
    data.val_files=/mnt/data/test.parquet \
    data.train_batch_size=1024 \
    • 不指定 --address 时,默认使用 RAY_ADDRESS 环境变量,若没有则使用默认值:'http://127.0.0.1:8265'
    • --runtime-env--working-dir 都是环境参数,也可以使用类似 --working-dir /mnt/workspace/ 的方式提交
    • data.train_files 等是传入到 python 脚本 main_ppo.py 的参数
    • 注意在开始 python 调用之前需要 -- (包含空格)

提交脚本后发生的行为

  • 在作业运行期间,ray job submit 命令默认会开启 --wait 模式,将 Driver 和各个 Worker 产生的日志实时地流式传输到终端的 stdout,方便监控进度

  • ray job submit 命令会一直等待,直到作业完成

    • 如果 Driver 进程成功退出(退出码为0),ray job submit 命令也会以 0 退出,表示作业成功
    • 如果作业失败(例如脚本报错),命令会以 1 退出
  • 如果 通过 ray submit 提交下面的脚本

    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    12
    13
    14
    15
    16
    17
    18
    19
    20
    21
    22
    23
    24
    25
    26
    27
    28
    29
    import ray # head 节点

    ray.init(address="auto") # head 节点

    # 定义阶段:
    ## head 节点(作业运行环境)装饰器 @ray.remote 将函数标记为远程任务;
    ## head 不执行函数体,仅注册到 Ray 的任务系统中; 函数定义存储在 Head 节点的内存中

    # 执行阶段:
    # 每个任务(square.remote(n) 提交的任务)在分配的 Worker or Head 节点上执行
    @ray.remote
    def square(x):
    import os
    import socket
    print(f"计算 {x} 的平方,运行在节点: {socket.gethostname()}")
    return x * x

    def main(): # Head 节点
    print(f"主程序运行在节点: {socket.gethostname()}")

    # 创建 20 个任务
    numbers = list(range(20))
    futures = [square.remote(n) for n in numbers] # head 节点创建任务并回收结果,任务执行会提交到 worker
    results = ray.get(futures)

    print(f"结果: {results}")

    if __name__ == "__main__": # Head 节点
    main() # Head 节点
  • 具体执行情况总结

    • main() 中的普通代码(打印主程序、提交任务、获取结果)
      • 仅 Head 节点 执行
      • 日志打印到 Head 节点上
      • 日志还会显示到提交终端(实时回显)
    • square.remote(n) 远程任务(20个)
      • Head + 所有 Worker 节点随机分配
      • 日志打印在 各自执行节点
      • Ray Job Server 会统一聚合 Head 节点与所有 Worker 节点的输出日志
      • 注:不管是主进程 main() 里的打印,还是远程 Actor/Task 内部的 print,都会实时流式回显到执行 ray job submit 的终端
        • 只有交互式直接连集群(ray.init(address=”auto”) 本地运行脚本直连集群),未经过 Job Server 托管时,Worker 任务打印才默认只保存在对应节点本地,不会自动回传到本机终端
  • 执行详情:

    • Head 节点(确定运行)
      • Driver 进程 执行内容:
        • ray.init()
        • 顶层的 if __name__ == "__main__" 以及 main() 内部的普通 Python 逻辑
    • Head 节点 + 所有 Worker 节点(由调度器决定)
      • 远程任务 square.remote(n) 提交的这 20 个任务不会 全部在 Head 上运行,也不会均匀固定分配
      • Ray 的默认调度器(DEFAULT 策略)会综合考虑集群所有节点(包括 Head 和 Worker)当前的 CPU 负载和可用资源,动态地将这 20 个任务分散到各个节点上
      • 核心:
        • Head 节点如果还有空闲 CPU 资源,也会被调度执行 一部分 square 任务
        • 如果 Worker 节点资源充足,大部分任务会倾向于分发给 Worker 节点以分担压力
  • 如何快速查看这 20 个任务的分布与日志

    • 由于日志分散在各节点,有三种方式查看:
      • 1)通过 Ray Dashboard(最推荐) :访问 Head 节点的 8265 端口,进入 Jobs -> 点击该作业 -> 查看 Logs
        • Dashboard 会自动聚合显示 Driver 日志和所有 Worker 的任务日志,并标注每条日志来自哪个节点(IP/主机名)
      • 2)通过 ray job logs :在客户端执行 ray job logs <submission_id>(需配合 Ray Job Server),它会自动从所有节点拉取并合并该作业相关的日志输出
      • 3)手动 SSH 登录 :分别登录 Head 和每个 Worker 节点,进入 /tmp/ray/session_latest/logs/,用 grep "计算.*的平方" worker-*.out 来逐一查看
补充:直接提交与 submit 提交的区别
  • 核心区别:两行代码完全可以一模一样,执行方式不一样,决定了会不会经过 Job Server,都包含 ray.init(address="auto")
    • ray job submit 启动 :由JobServer接管日志,全集群日志统一回显到提交终端
    • 直接 python xxx.py 运行 : 进程直连集群,worker 打印默认只留在对应节点本地,不会回传
方式一:ray job submit 方式
  • 命令:

    1
    ray job submit -- python script.py
  • 脚本内部依然会写:

    1
    ray.init(address="auto")
  • 说明:

    • 整个脚本是由 Ray Job Server 拉起运行
    • Job Server 会拦截全集群所有 worker 的 stdout/stderr,把所有节点的打印统一收集,实时推送到你执行命令的终端
    • 哪怕远程 worker 节点上的 print,也会回显到提交命令的窗口
方式二:本地直接运行脚本方式
  • 在自己电脑/head节点终端直接执行:

    1
    python script.py
  • 脚本里同样也是:

    1
    ray.init(address="auto")
  • 此时没有经过 Job Server,只是当前进程直连 Ray 集群

  • 默认情况下:

    • 主进程 print:在当前终端输出
    • 远程 Worker 进程 print:只写入对应机器本地日志文件,不会自动回传到当前的终端

其他高级配置

  • py_modules :指定需要导入的自定义 Python 模块路径,支持将本地模块添加到 Python 路径(sys.path),示例如下:

    1
    ray.init(runtime_env={"py_modules": ["./my_utils"]}) # 同步 my_utils 模块并添加到路径
    • py_modules 用于将本地开发的自定义 Python 模块分发到Ray集群的所有节点上
    • 当任务或 Actor 依赖了项目内的本地模块(而非通过 pip 安装的公开包)时,这个配置就非常有用
    • 还可以通过下面的示例上传一些导入的包/模块:
      1
      2
      3
      4
      5
      6
      7
      8
      9
      10
      11
      12
      13
      import ray
      import DataTransformerUtils # 你的本地模块
      import CustomLoggingUtils # 你的本地模块

      ray.init(runtime_env={
      "py_modules": [DataTransformerUtils, CustomLoggingUtils]
      })

      @ray.remote
      def test_my_module():
      # 无需在函数内再导入,可以直接使用
      DataTransformerUtils.f()
      CustomLoggingUtils.g()
  • env :指定预定义的环境名称(如 Ray 集群中已配置的共享环境),避免重复配置,示例如下:

    1
    ray.init(runtime_env={"env": "shared-training-env"})  # 使用集群中预定义的环境
    • 工作原理:集群管理员可以预先在集群上创建并“命名”一个包含特定依赖(如 pip 包、working_dir 等)的环境
      • 用户在提交任务时,只需通过 {“env”: “环境名称”} 引用即可
    • 这个 env 参数与在 runtime_env 中使用的 env_vars(用于设置环境变量)是两个完全不同的概念,不要混淆
    • 核心优势:避免重复配置
  • 特别说明:runtime_env 还支持其他非官方扩展,比如 verl 库中就为英伟达显卡配置了 nsight 工具的参数传入

    1
    2
    ./verl/trainer/main_ppo.py
    runner = TaskRunner.options(runtime_env={"nsight": nsight_options}).remote()
    • nsight 是 NVIDIA Nsight 系列性能分析工具的配置参数

      • runtime_env 中传入 nsight 配置,可以让 Ray 在执行任务时自动挂载性能分析工具,用于 GPU 调试和优化
    • 工作原理与示例:

      • 1)框架集成 :像 vLLMverl 这样的高级框架,会在其内部代码中根据条件动态构建包含 nsight 字段的 runtime_env
      • 2)配置传递 :用户在调用 TaskRunner 等 Actor 时,将 nsight 配置作为 runtime_env 的一部分传入
      • 3)自动挂载 :Ray 在启动该 Actor 或任务时,会识别 nsight 字段,并按照配置(如追踪的API、输出文件名等)启动 nsys 等分析工具
    • 一个具体的 nsight 配置示例可能如下:

      1
      2
      3
      4
      5
      6
      7
      8
      9
      nsight_options = {
      "trace": "cuda,nvtx,cublas,ucx",
      "cuda-memory-usage": "true",
      "cuda-graph-trace": "graph",
      }

      # 在verl或vLLM等框架中,可能这样使用(注意是 cuda 可用时才开启)
      if is_cuda_available:
      runner = TaskRunner.options(runtime_env={"nsight": nsight_options}).remote()
    • 注:如果不传入 nsight 参数,Ray 任务/Actor 将不会启动任何性能分析工具(如 NVIDIA Nsight Systems),程序会正常执行,只是不会产生性能追踪文件