Hexo

凡事预则立,不预则废


  • Home

  • Tags

  • Archives

  • Navigation

  • Search

Python——Ray-资源管理详解


前置补充:Raylet 和 GCS

Raylet(运行于集群的每个节点上)

  • 主要功能包括:
    • 本地资源管理 :监控并分配所在节点的 CPU、GPU 及内存资源
    • 工作进程管理 :维护一个工作进程池(Worker Pool),负责任务执行进程的分配、回收与异常处理
    • 本地对象存储 :管理节点上的共享内存存储(Object Store),处理本地数据的读写与内存回收
    • 信息上报 :定期向 GCS 发送心跳,上报本节点的资源状态、任务执行情况及对象存储位置
  • Raylet 是无状态的
    • 自身状态存储在 GCS 中,故障后可从 GCS 恢复

GCS(运行于集群的头节点)

  • GCS(Global Control Store)是 Ray 的集中式元数据管理组件
  • 主要功能包括:
    • 元数据存储 :作为中心化键值存储(后端常为 Redis 或 RocksDB),保存全局状态
      • 包括节点拓扑、Actor 位置、分布式对象的位置映射、Job 及 Task 的执行记录
    • 集群协调 :处理节点注册,通过心跳检测节点存活状态,触发故障节点的标记与恢复流程
    • 全局服务 :提供 Actor 全局命名与定位、Placement Group 管理、分布式锁等原子化服务
  • GCS 是有状态的
    • 需配置高可用后端以支持容错

Raylet 与 GCS 的协作

  • 状态汇总与调度 :各 Raylet 上报状态至 GCS,GCS 构建全局资源视图
    • 调度器查询该视图选定目标节点,并向该节点的 Raylet 下发任务执行指令
  • 数据定位与传输 :工作进程读取远程对象时,本地 Raylet 向 GCS 查询该对象的存储节点列表,随后通过点对点网络直接从目标节点的 Raylet 关联存储中拉取数据
  • 故障处理 :GCS 检测到某 Raylet 失联后,停止向该节点分配任务,并触发其上任务和 Actor 的重调度
    • 若 GCS 自身故障,各 Raylet 依据高可用配置尝试重连或选举新实例

资源定义(节点侧容量声明)

  • Ray 的资源模型基于键值对(Key-Value)
    • Key 为资源名称(字符串)
    • Value 为 容量(浮点数)

Ray 资源分类

  • 内置资源(Built-in) :CPU、GPU、memory、object_store_memory
    • 这四个名称被 Ray 调度器(GCS,Global Control Store)硬编码识别
  • 自定义资源(Custom) :用户任意命名的资源,例如inference_license、cache_slots 等自定义名称
    • 注:自定义资源名称不能与内置资源(CPU/GPU/memory/object_store_memory)重名,否则会被当作内置资源处理或报错

节点启动时的容量声明

  • 在启动 Ray 节点进程时,通过命令行参数固定该节点的资源总容量,这些数值注册到全局状态存储(即存储到 GCS 中)后不可动态修改

    1
    2
    # 声明该节点拥有 8 个 CPU 单位、4 个 GPU 单位、1 个自定义资源单位
    ray start --num-cpus=8 --num-gpus=4 --resources='{"inference_license": 1}'
    • 这里的声明是逻辑资源,所以 --num-cpus=1000 也可以,完全可以超过真实 cpu 核数
      • 但这个节点总资源是用于限制 Ray 提交太多任务到同一个节点的,声明过多逻辑资源容易造成资源抢占甚至崩溃
    • --num-cpus 默认值物理 CPU 核心数(占用全部物理 CPU 核心),对于头节点(包含 --head 参数)的节点,建议设置为 0

资源请求与限制(Task/Actor 侧语法)

  • 用户在提交计算任务时,通过 ray.remote 装饰器或 .options() 方法声明资源需求量
    • 提交任务时的需求量是硬性约束 ,调度器仅在有节点剩余容量满足该数值时才会进行调度(对应前面节点的逻辑资源申明)
    • 调用时的任务的 .options() 参数可以覆盖定义时的 ray.remote 参数

Ray 中的任务资源申请

  • Ray 资源是逻辑资源,非物理隔离,更多是用于调度时的准入控制

    • Ray 不会将一个物理 CPU 核心绑定并独占给这个任务
    • Ray 甚至不会阻止一个 num_cpus=1 的任务在内部启动多个线程
      • num_cpus=1 下内部也可以使用多个线程使用多个物理 CPU 核心
      • 实际的 CPU 使用由操作系统调度
  • Task :默认需求为 num_cpus=1

    • 若不显式设置,每个提交的任务都会消耗 1 单位 CPU 逻辑配额
  • Actor :默认需求为 num_cpus=0

    • 这意味着若不手动设置,调度器不会因 CPU 资源不足而阻止 Actor 创建,可能导致单节点上同时运行大量 Actor,引发物理内存和 CPU 的过度争抢
      • 强烈建议始终为 Actor 显式设置 num_cpus

        Task/Actor 资源申请示例

  • Task(无状态函数)资源示例:

    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    import ray
    ray.init()

    # 定义阶段:装饰器作用于函数
    @ray.remote(num_cpus=2, num_gpus=1)
    def my_function(x):
    return x + 1

    # 提交阶段:使用 .options() 动态覆盖资源需求
    # .remote() 执行后返回 ObjectRef(对象引用)
    task_ref = my_function.options(num_cpus=0.5, resources={"inference_license": 1}).remote(10)
    • my_function 是 Ray 封装的远程可调用对象
    • my_function.options(...) 生成一个携带新资源配置的副本
    • .remote(10) 将任务下发到集群,返回占位符 ObjectRef
  • Actor(有状态类)资源申请示例

    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    12
    13
    14
    15
    16
    17
    import ray
    ray.init()

    # 定义阶段:装饰器作用于类
    @ray.remote(num_cpus=2, num_gpus=1, memory=1024 * 1024 * 1024)
    class MyActor:
    def __init__(self, val):
    self.val = val
    def get(self):
    return self.val

    # 创建阶段:使用 .options() 动态覆盖资源需求
    # .remote(init_args) 在远端创建实例,返回 ActorHandle(句柄)
    actor_handle = MyActor.options(num_cpus=0.5, resources={"inference_license": 1}).remote(100)

    # 通过句柄调用方法,返回 ObjectRef
    result_ref = actor_handle.get.remote()
    • MyActor 是Ray封装的远程Actor构造器
    • MyActor.options(...) 生成一个携带新资源配置的构造器副本
    • .remote(100) 在远端节点实例化对象,返回 ActorHandle(非本地类实例)
    • 设置 num_cpus=0.5 表明该任务声明需要 0.5 个单位的逻辑 CPU 配额
      • 调度准入时仅占用 0.5 cpu
      • 物理执行时理论上可以多用
      • 环境变量设置(软性建议)
        • Ray 为执行该任务的工作进程设置环境变量 OMP_NUM_THREADS,规则为将num_cpus向上取整 :
          • 0.5向上取整为1,故 OMP_NUM_THREADS=1
          • Ray 确实会设置 OMP_NUM_THREADS,但并非所有库都遵守(如 NumPy 默认不限制)
        • 科学计算库(OpenBLAS、MKL、PyTorch)在初始化时读取此变量作为线程池上限,从而软性建议该任务减少内部并行度,避免挤占其他任务的物理 CPU 时间

任务提交后的调度器的资源分配流程

  • 当任务提交后,Ray 的 GCS 调度器执行以下确定性流程:
  • 步骤一:可行性检查(Feasibility)
    • 遍历所有存活节点,剔除两类节点:
      • 1)资源类型缺失 :例如在无 GPU 节点上请求 num_gpus=1(注意:这里是提交任务时,不是初始化节点进程时)
      • 2)剩余容量不足 :节点的某种资源剩余量 < 任务请求量
  • 步骤二:节点排序(Scoring)
    • 默认调度策略(DEFAULT)对候选节点进行综合打分:
      • 数据本地性优先 :若任务的输入数据(通过 ray.put 存储)已存在于某节点的对象存储中,该节点得分大幅提升
      • 负载均衡优先 :计算每个节点资源利用率(已占用/总容量)
        • 节点剩余可用资源越多,得分越高
  • 步骤三:Top-K 随机选取
    • 为避免所有任务涌入得分最高的单一节点,调度器从 排序前 K 个 节点中随机抽取一个
    • K 值默认为集群总节点数的 20%(由环境变量 RAY_SCHEDULER_TOP_K_AWARE=1 启用,RAY_SCHEDULER_TOP_K_AWARE_FRACTION 控制)

高级资源分配机制

自定义资源(数值型准入)

  • 自定义资源仅用于数值累加准入控制
    • 例如节点有{"inference_license": 1.0},两个任务分别请求 0.5 可同时运行,各请求 1.0 则只能串行
  • 注:不推荐将自定义资源用于节点属性筛选(如 GPU 型号),因其无法表达 不等于 或 包含 等逻辑

标签选择器(属性型调度,Ray 2.49+)

  • 官方推荐使用标签选择器处理非数值型调度约束 ,其机制独立于资源配额:

    • 节点打标签 :

      1
      ray start --labels='{"region": "us-west", "instance_family": "spot"}'
    • 任务声明约束 :在 .options() 中传入 NodeLabelSchedulingStrategy,支持 In、NotIn、Exists 操作符

      1
      2
      3
      4
      5
      6
      from ray.util.scheduling_strategies import NodeLabelSchedulingStrategy

      strategy = NodeLabelSchedulingStrategy(
      label_selector={"region": "us-west", "instance_family": "spot"}
      )
      task_ref = my_function.options(scheduling_strategy=strategy).remote()
  • 调度器在可行性检查阶段直接剔除标签不匹配的节点,再执行后续排序与随机选取

放置组(原子性跨节点资源预留)

  • 当任务需要 同时成功预留跨多个节点的多份资源(如分布式训练:1个PS + 4个Worker),默认调度可能因部分资源预留成功、部分失败而陷入死锁
  • 机制 :调用 placement_group()一次性向 GCS 声明资源包(Bundles)列表
    • GCS 执行原子调度,要么全部预留成功,要么全部回滚
  • 绑定 :后续 Task/Actor 通过 PlacementGroupSchedulingStrategy 绑定到特定 Bundle 上执行

核心配置建议

  • Head 节点的 CPU :
    • 建议设置 --num-cpus=0
    • Head 节点仅负责任务调度和 GCS 状态维护,若允许用户任务占用其逻辑 CPU,可能导致系统组件响应延迟
    • 直接使用 ray start 启动时,默认值为物理 CPU 核心数
  • Actor 的 CPU :
    • 建议显式设置 num_cpus(如num_cpus=1)
    • 默认值为 0 会导致单节点可无限创建 Actor,耗尽物理内存
  • 内存资源 :
    • 务必为任务设置 memory 参数
    • 若不设置,Ray 无法进行内存准入控制,高并发下极易触发节点 OOM 导致进程崩溃
  • 浮点数精度 :
    • 调度器进行浮点数累加
    • 若使用 0.1 等非精确二进制小数,多次加减后可能导致 剩余容量 与 请求量 的误差大于阈值,造成调度误判
      • 建议使用0.5或0.25
  • 资源与属性分离 :
    • 节点型号、区域等属性型约束请使用标签选择器 ,不要使用自定义资源(自定义资源仅保留给数值型并发控制)

补充:指定任务执行 Worker 的方法

  • Ray 中有多种方法可实现

方法一:使用节点亲和性策略 (NodeAffinitySchedulingStrategy)

  • 最直接的方法,通过节点的唯一 ID 将任务强制绑定 到特定节点

  • 在提交任务或创建 Actor 时,通过 scheduling_strategy 参数指定 NodeAffinitySchedulingStrategy,并将 node_id 设为刚才获取的 ID

    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    12
    13
    14
    from ray.util.scheduling_strategies import NodeAffinitySchedulingStrategy

    # 假设这是你想要运行任务的目标节点ID
    target_node_id = "YOUR_TARGET_NODE_ID"

    @ray.remote
    def my_task():
    return "This task runs on a specific node."

    # 将任务调度到目标节点
    # soft=False 表示强制调度,如果目标节点不可用,任务将失败;如果设为 `True`,当目标节点不可用时,任务可以被调度到其他节点
    result = my_task.options(
    scheduling_strategy=NodeAffinitySchedulingStrategy(node_id=target_node_id, soft=False)
    ).remote()
  • PS:查看集群中有哪些节点:

    1
    2
    3
    4
    5
    6
    7
    8
    import ray
    ray.init()

    # 获取集群中所有节点的信息
    nodes = ray.nodes()
    for node in nodes:
    # NodeID 是节点的唯一标识符
    print(f"Node ID: {node['NodeID']}, Resources: {node['Resources']}")

方法二:使用标签选择器 (Label Selector) 【推荐】

  • 标签选择器是 Ray 2.49 版本后引入的更灵活、更强大的方式

  • 标签选择器通过为节点打上标签,然后在任务中声明需要匹配的标签来选择节点

  • 1)为节点打标签 :在启动 Ray 节点时,通过 --labels 参数为其添加标签

    1
    2
    # 启动一个 Worker 节点,并打上标签
    ray start --address='YOUR_HEAD_NODE_ADDRESS' --labels='{"gpu_type": "A100", "region": "us-west"}'
  • 2)提交任务 :在任务中通过 label_selector 参数指定标签选择条件

    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    12
    13
    14
    15
    16
    17
    18
    19
    20
    @ray.remote
    def my_task():
    return "This task runs on a node with specific labels."

    # 将任务调度到拥有 "gpu_type=A100" 标签的节点上
    result = my_task.options(
    label_selector={"gpu_type": "A100"}
    ).remote()

    # 定义一个标签选择器,要求节点的 "instance-type" 标签不在 ["spot"] 中
    from ray.util.scheduling_strategies import NodeLabelSchedulingStrategy
    strategy = NodeLabelSchedulingStrategy(
    label_selector={
    "instance-type": {
    "operator": "NotIn",
    "values": ["spot"]
    }
    }
    )
    result = my_task.options(scheduling_strategy=strategy).remote()
    • 标签选择器支持 In、NotIn 等复杂操作,可以实现 “运行在非 Head 节点” 或 “运行在region为us-west或us-east的节点” 等精细控制

方法三:使用放置组 (Placement Group)

  • 放置组主要用于跨节点、多资源 的原子性预留,但也可以用来将任务“引导”到特定节点

  • 1)创建放置组 :定义资源包(Bundle),并通过 strategy 策略(如 STRICT_PACK)控制其分布

    1
    2
    3
    4
    5
    6
    7
    from ray.util.placement_group import placement_group, PlacementGroupSchedulingStrategy
    from ray.util.scheduling_strategies import PlacementGroupSchedulingStrategy

    # 创建一个放置组,包含一个需要 1 个 CPU 资源的 Bundle
    # STRICT_PACK 策略会尝试将所有 Bundle 放在同一个节点上
    pg = placement_group([{"CPU": 1}], strategy="STRICT_PACK")
    ray.get(pg.ready()) # 阻塞当前程序,直到指定的放置组在集群中成功创建并完成资源预留,若集群资源不足,该调用会永久阻塞,建议设置超时(如 pg.ready(timeout_seconds=30))
  • 2)获取 Bundle 所在节点 :通过 placement_group 的 API 可以获取其 Bundle 被调度到的节点 ID

  • 3)绑定任务 :将任务绑定到该放置组,从而让它运行在放置组所在的节点上

    1
    2
    3
    4
    5
    6
    7
    8
    @ray.remote
    def my_task():
    return "This task runs where the placement group is."

    # 将任务调度到该放置组所在的节点
    result = my_task.options(
    scheduling_strategy=PlacementGroupSchedulingStrategy(placement_group=pg)
    ).remote()

方法四:使用自定义资源 (Custom Resource) 【旧方案,不推荐】

  • 这是一种旧方案,通过伪造“资源”来实现节点选择。例如,为特定节点设置一个{"node-type": 1}的资源,然后在任务中请求该资源

  • 示例:

    • 启动节点

      1
      2
      # 在目标节点启动时声明自定义资源
      ray start --address='HEAD_ADDRESS' --resources='{"node-type": 1}'
    • 提交任务

      1
      2
      3
      @ray.remote(resources={"node-type": 1})
      def my_task():
      return "This task runs on a node with the custom resource."
  • 不推荐原因这种方式的原因是:

    • 这种方法混淆了“资源(Resource)”和“约束(Constraint)”的概念,且只能表达“等于”的逻辑,缺乏灵活性
    • 官方建议使用标签选择器作为替代方案

Python——Ray-远程函数与本地函数的区别


整体说明

  • 远程函数与本地函数的区别主要在 序列化机制 和 执行位置 两个维度
  • 序列化本质差异:
    • 本地函数可以理解为“传引用”,依赖执行环境已有定义
      • 注:本地函数也不仅仅是 “传引用”
        • Python 的 pickle 序列化函数时,实际上是序列化函数的名称和所在模块的路径
        • 反序列化时,需要在目标环境中导入同名模块、找到同名函数
        • 因此 Python 本地函数调用依赖目标环境与源环境“同构”
        • Ray 跨节点时,Worker 进程的 __main__ 模块通常与 Driver 不同,所以会失败
    • Ray 远程函数是 “传定义+环境” ,集群自动同步,支持跨节点;
  • 执行位置差异:
    • 本地函数固定在调用方进程,无分布式能力;
    • Ray 远程函数由集群调度,可分布式并发执行;
  • 使用场景:
    • 本地函数:适用于单进程/单节点的简单逻辑,无需分布式;
    • Ray 远程函数:适用于分布式计算、并发任务、跨节点执行,是 Ray 分布式能力的核心
  • 核心差异总览
    对比维度 本地函数(未用 @ray.remote 装饰) Ray 远程函数(用 @ray.remote 装饰)
    序列化方式 依赖 Python 原生 pickle,仅序列化「函数引用」 Ray 自定义序列化(结合 pickle+集群元数据),序列化「函数元信息+代码定义」
    序列化限制 无法跨节点传递(远程节点无函数定义,引用失效) 可跨节点传递(集群自动同步函数定义到执行节点)
    执行位置 固定在「调用方所在的本地进程/线程」 分布式调度到「集群任意节点的 Worker 进程」(可指定资源)
    执行特性 同步执行,阻塞调用方;无并发调度能力 异步执行,返回 ObjectRef;支持集群级并发/分布式调度
    依赖传递 需手动确保执行环境有函数依赖(如导入、变量) Ray 自动打包函数依赖(如嵌套函数、闭包变量)并分发

序列化机制:“仅传引用” vs “传定义+元信息”

  • 序列化的核心目的是:让函数能在「非定义环境」中被正确执行
  • 两者的序列化逻辑完全不同:

本地函数:仅序列化“函数引用”,无实际代码

  • Python 原生 pickle 序列化本地函数时,不会打包函数的代码本身 ,只会记录函数的「模块路径+函数名」(比如 __main__.add)
  • 这种“引用式序列化”仅在「同一进程/同一节点且函数已定义」的场景下有效,跨节点会直接失效
  • 错误示例:
    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    12
    13
    14
    15
    16
    17
    18
    import ray

    ray.init(ignore_reinit_error=True)

    # 本地函数
    def add_remote(a, b):
    return a + b

    # 直接传递远程函数的引用(Ray 自动处理序列化)
    @ray.remote
    def execute_remote_func(func, x, y):
    return func(x,y) # 远程工作进程无法识别调用方的 local func,错误

    # 跨节点调度执行(单节点可以成功,但集群有多个节点会失败)
    result_ref = execute_remote_func.remote(add_remote, 2, 3)
    print(ray.get(result_ref)) # 单节点输出:5(成功执行);多节点执行错误

    ray.shutdown()

Ray 远程函数:序列化“函数元信息+代码定义”

  • Ray 对远程函数的序列化做了增强 :
    • 1)序列化时,不仅记录函数引用,还会打包函数的代码定义、依赖模块、闭包变量(若有);
    • 2)远程节点接收后,会自动还原函数的执行环境(无需手动导入);
    • 3)底层用 Ray 自定义的序列化器(兼容 pickle,但更适合分布式场景)
  • 正确示例:远程函数跨节点调用成功
    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    12
    13
    14
    15
    16
    17
    18
    19
    import ray

    ray.init(ignore_reinit_error=True)

    # Ray 远程函数(已注册,自动序列化代码)
    @ray.remote
    def add_remote(a, b):
    return a + b

    # 直接传递远程函数的引用(Ray 自动处理序列化)
    @ray.remote
    def execute_remote_func(func, x, y):
    return ray.get(func.remote(x, y)) # 远程节点能识别并执行

    # 跨节点调度执行(即使集群有多个节点也能成功)
    result_ref = execute_remote_func.remote(add_remote, 2, 3) # 注意:传入的参数 add_remote 本身也需要是 @ray.remote 封装过的 Ray 远程函数
    print(ray.get(result_ref)) # 输出:5(成功执行)

    ray.shutdown()

补充:Ray 还支持 嵌套远程函数 闭包变量传递

  • 比如在远程函数中引用本地变量,Ray 会自动序列化传递:
    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    12
    13
    14
    15
    import ray

    ray.init(ignore_reinit_error=True)

    @ray.remote
    def outer_remote(x):
    # 闭包变量 x 会被 Ray 自动序列化到远程节点
    @ray.remote
    def inner_remote(y):
    return x + y
    return inner_remote.remote(10)

    print(ray.get(ray.get(outer_remote.remote(5)))) # 输出:15

    ray.shutdown()

执行位置:“本地固定” vs “集群分布式调度”

  • 执行位置的差异是两者最直观的区别,直接决定了是否能利用集群资源:

本地函数:执行在 调用方所在进程

  • 本地函数的执行位置完全固定:
    • 无论在哪里调用(即使在远程函数内部调用本地函数),函数都会在 发起调用的进程 中执行【存疑】
      • 问题:这里描述有错(部分书籍会这样写),理论上远程函数内部无法调用本地函数,所以应该加上一句,在可以调用成功的前提下
    • 若在远程函数中调用本地函数,本质是在「远程节点的 Worker 进程」中执行,但该进程没有本地函数的定义(除非手动同步代码),所以必然失败;
    • 无并发能力:多个调用会串行执行在同一个进程/线程(或 Python 多进程的子进程,但需手动管理)

远程函数:执行在「集群 Worker 进程」

  • Ray 远程函数的执行位置由 Ray 集群的调度器统一管理:

    • 1)调用 func.remote() 时,会向 Ray 调度器提交一个任务
    • 2)调度器根据集群节点的资源(CPU、GPU、内存)情况,将任务分配到任意可用节点的 Worker 进程
    • 3)执行完成后,结果会存储在 Ray 的对象存储中,通过 ray.get() 可获取
    • 4)支持并发:多个 remote() 调用会被调度到不同 Worker 进程/节点,并行执行
  • 示例:远程函数分布式执行(多节点/多进程并发)

    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
       import ray
    import os
    import time

    # os.environ["RAY_DEDUP_LOGS"] = "0" # 本意是让每个进程结果都完整输出,但这行代码仅当前进程生效,需要启动前配置环境变量才可以
    # # 如果是通过代码定义,则在 os.environ 设置在 ray.init() 之前进行才能生效,因为 Worker 进程在初始化时读取环境变量
    ray.init(ignore_reinit_error=True)

    # Ray 远程函数:打印执行节点的进程 ID 和节点名
    @ray.remote
    def add_remote(a, b):
    node_name = ray.util.get_node_ip_address() # 获取执行节点 IP
    pid = os.getpid() # 获取执行进程 ID
    print(f"在节点 {node_name} 的进程 {pid} 执行 add({a}, {b})")
    time.sleep(1) # 模拟耗时操作
    return a + b

    # 提交 5 个并发任务(会被调度到不同 Worker 进程)
    start = time.time()
    result_refs = [add_remote.remote(i, i*2) for i in range(5)]
    results = ray.get(result_refs) # 等待所有任务完成
    end = time.time()

    print("结果:", results) # 输出:[0, 3, 6, 9, 12]
    print(f"总耗时: {end - start:.2f}s") # 约 1s(并发执行,而非 5s 串行)

    ray.shutdown()
  • 执行上述脚本:

    1
    2
    export RAY_DEDUP_LOGS=0
    python demo.py
    • 注意:仅在代码里面添加 os.environ["RAY_DEDUP_LOGS"] = "0" 是不够的,因为:
      • Ray 的日志去重功能是在 Worker 进程启动时就决定的,而 Worker 是由 Ray 的主进程(Driver)启动的
      • 上面的代码在 ray.init() 之后才启动 Worker,那么环境变量必须在 Driver 启动 Worker 之前就传递过去,否则 Worker 进程会继承默认的去重配置
      • 所以最安全的打印所有日志的方式就是再启动脚本前配置环境变量
    • 另一种实现方式是在远程函数中返回 PID,然后由 Driver 打印
  • 输出示例:

    1
    2
    3
    4
    5
    6
    7
    8
    2025-11-04 11:42:43,175 INFO worker.py:1918 -- Started a local Ray instance. View the dashboard at 127.0.0.1:8265 
    (add_remote pid=14393) 在节点 127.0.0.1 的进程 14393 执行 add(2, 4)
    (add_remote pid=14399) 在节点 127.0.0.1 的进程 14399 执行 add(4, 8)
    (add_remote pid=14398) 在节点 127.0.0.1 的进程 14398 执行 add(1, 2)
    (add_remote pid=14400) 在节点 127.0.0.1 的进程 14400 执行 add(3, 6)
    (add_remote pid=14396) 在节点 127.0.0.1 的进程 14396 执行 add(0, 0)
    结果: [0, 3, 6, 9, 12]
    总耗时: 1.62s

附录:远程调用时传入的函数指针必须是远程函数

  • 在 Ray 中不支持直接传入 local 函数指针作为远程函数的执行对象,需通过 Ray 装饰器(@ray.remote)将函数注册为远程可执行,再通过 函数名.remote() 调用(本质是基于函数标识而非指针传递)
  • 总结:
    • 不推荐将普通函数作为参数传递给 Ray 远程函数
    • 推荐使用 @ray.remote 装饰器或在远程函数内部定义逻辑
    • 注意:一些代码在单机环境下可能碰巧能运行,但不具有可移植性和可靠性(这一点需要注意 Ray 本地调试通过可能也无法分布式运行)

错误示例(未注册本地函数)

  • 若 add 未被 @ray.remote 注册,它只是一个本地函数 ,无法在 Ray 分布式环境中执行,直接传递给远程函数(如 execute_func)会报错

    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    12
    13
    14
    15
    16
    17
    18
    19
    20
    21
    22
    23
    import ray

    ray.init(ignore_reinit_error=True)

    # 未注册的本地函数
    def add(a, b):
    return a + b

    # 已注册的远程函数
    @ray.remote
    def execute_func(func, x, y):
    # 此处调用本地函数会失败,因为 func 在远程节点无定义
    # # 远程节点的工作进程无法导入本地主模块的 add_local 函数,也无法序列化传递普通函数,可能会直接抛出 SerializationError
    # # 单进程/单节点下调用指针函数可以执行,但是分布式情况下,local_func 无法被序列化,会出错
    return func(x, y) # 报错:NameError,PicklingError 或 SerializationError

    # 调用会抛出异常
    try:
    result = ray.get(execute_func.remote(add, 4, 6))
    except Exception as e:
    print("错误:", e) # 提示无法序列化或找不到函数

    ray.shutdown()
  • 核心原因:Ray 远程函数执行依赖序列化传输和集群节点间代码同步

    • 未注册的本地函数无法被序列化为集群可识别的任务,且远程节点没有该函数的定义,会导致执行失败

正确示例(远程函数调用)

  • Ray 的远程函数依赖集群调度,通过 @ray.remote 显式注册后使用远程调用函数调用

    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
    30
    31
    import ray

    ray.init(ignore_reinit_error=True)

    # 定义远程函数(会注册到 Ray 集群)
    @ray.remote
    def add(a, b):
    return a + b

    # 远程函数,可接收其他远程函数的调用结果
    @ray.remote
    def execute_func(func, x, y):
    # 这里 func 是远程函数标识,通过 .remote() 触发执行
    result = ray.get(func.remote(x, y)) # 使用远程调用的方式调用函数指针,实现调用远程函数,正确!
    # result = func(x, y) # remote 函数无法被直接调用,错误!
    # result = add(x,y) # remote 函数无法被直接调用,错误!
    # result = add_local(x, y) # add_local 当做 local 函数调用(注意:不再是指针传入),正确!
    return result

    # # 不使用 remote 直接调用 远程函数,错误
    # result1 = add(2, 3)

    # 使用remote直接调用远程函数,正确
    result1 = ray.get(add.remote(2, 3))
    print("直接调用结果:", result1) # 输出:5

    # 间接通过另一个远程函数调用(模拟"传递函数逻辑")
    result2 = ray.get(execute_func.remote(add, 2, 3))
    print("间接调用结果:", result2) # 输出:10

    ray.shutdown()
  • Ray 的远程函数依赖集群调度,需通过 @ray.remote 显式注册,无法像本地代码那样传递函数指针(内存地址在分布式环境中无效)

  • 若需在远程函数中复用其他函数逻辑,直接传递已注册的远程函数名(如示例中的 add),再通过 func.remote() 调用即可

1…127128129…352
San Ye

San Ye

Stay Hungry. Stay Foolish.

704 posts
53 tags
© 2026 San Ye
Powered by Hexo
|
Theme — NexT.Gemini v5.1.4