Ray分布式计算框架

Ray 分布式计算框架:从 Actor 模型到 Unified AI Runtime 的深度工程实践

Ray 最初由加州大学伯克利分校 RISElab 孵化,2017 年开源,现在已是 AI/ML 领域分布式计算的事实标准框架。不同于 Spark 批量数据处理模型,Ray 原生支持有状态计算(Actor)、细粒度任务图调度、共享内存对象存储和容错重试机制,这些特性让它天生适配 ML 训练批处理并行、模型服务弹性扩缩容、强化学习环境交互等复杂场景。

一、Ray Core 调度模型

1.1 GCS 调度器核心架构

Ray 的调度是一个层级体系:每个 Worker 节点运行本地调度器负责节点内资源分配(CPU/GPU 亲和、本地对象优先),全局 GCS 调度器负责跨节点资源放置。GCS 基于 Redis 或内置实现,维护所有节点的资源视图和任务等待队列。

核心调度流程:

import ray
ray.init()

@ray.remote(num_cpus=2, num_gpus=1)
def gpu_task(data):
    import torch
    return torch.tensor(data).cuda().mean().item()

# 任务提交时不立即执行,返回 ObjectRef 句柄
refs = [gpu_task.remote([i]*1000) for i in range(100)]
results = ray.get(refs)  # 阻塞等待结果

1.2 全局调度 vs 两级调度

当任务需要多节点资源时:

  1. GCS 选择目标节点(资源足够且数据局部性高的节点)
  2. 目标节点 LocalScheduler 再决定具体 Worker
  3. 任务序列化通过 Plasma 共享内存零拷贝传递

这种两级调度在大规模 ML 训练场景中表现出色:数据密集任务优先调度到数据节点,GPU 任务避免跨 NUMA 迁移。

二、Actor 模型深入

Ray Actor 解决了无状态 remote function 的表达局限。Actor 是有状态的,每个 Actor 被调度到固定节点,所有方法调用排队串行执行:

@ray.remote
class Counter:
    def __init__(self):
        self.value = 0
    def increment(self):
        self.value += 1
        return self.value
    def get(self):
        return self.value

counter = Counter.remote()
# 并发调用自动排队串行执行
refs = [counter.increment.remote() for _ in range(1000)]
print(ray.get(counter.get.remote()))  # 1000

Actor 通过命名引用实现松耦合服务发现:

# 全局可访问,任何节点可获取引用
counter = ray.get_actor("global_counter")
ray.get(counter.get.remote())

实际工程中,Actor 常用于实现:

  • 参数服务器:Actor 持有模型参数副本,Worker 轮询梯度
  • 环境模拟器:强化学习环境中,Actor 维护 Gym 环境状态
  • 在线路由层:请求调度器按哈希/负载策略转发

三、Plasma 对象存储

Ray 的对象存储基于 Apache Arrow Plasma,核心特点是共享内存零拷贝:

  • 写入:put(value) → ObjectID,数据拷贝进共享内存
  • 读取:get(ObjectID) → 直接映射同一段物理内存,无序列化开销
  • 溢出:超阈值自动 spill 到本地 SSD

该特性对 GPU 场景至关重要:同一节点不同 Worker 可通过共享内存交换中间张量,避免 PCIe 带宽瓶颈。

实测数据:1GB 张量跨节点传输,Ray Plasma 比基于 gRPC 的序列化快 5-8 倍;节点内同 NUMA 下达到内存带宽上限。

四、Data / Train / Tune / Serve 生态

Ray 生态覆盖 AI 工作负载全生命周期:

层 模块 功能
编排层 Ray Core 任务调度、Actor、对象存储
AI 工作负载 Ray Train 分布式 PyTorch / Horovod 训练
Ray LLM / vLLM 模型可滚动升级部署
Ray Tune 超参调优(PBT、ASHA 调度)
Ray Data 批量数据流式预处理
原生 Serve Ray Serve 请求级副本管理
经典 ML RLlib 强化学习

Ray Serve 示例——零停机多模型灰度部署

from ray import serve

@serve.deployment(num_replicas=3, ray_actor_options={"num_gpus": 0.5})
class ModelV1:
    def __call__(self, request):
        return {"pred": model_v1(request)}

@serve.deployment(num_replicas=1)
class ModelV2:
    def __call__(self, request):
        return {"pred": model_v2(request)}

# 灰度 10% 流量到 V2
app = ModelV1.bind().route_prefix("/v1") | ModelV2.bind().route_prefix("/v2")

Ray Serve 的请求级路由特性让其区别于传统 serving 框架:每个 HTTP 请求独立路由到指定 Deployment,天然支持蓝绿发布、金丝雀测试和流量切分。

五、生产部署要点

5.1 集群启动

# Head 节点启动
ray start --head \
  --port=6379 \
  --dashboard-host=0.0.0.0 \
  --num-gpus=4 \
  --max-worker-count=512

# Worker 节点加入
ray start --address="head:6379" --num-gpus=4

生产环境推荐使用 KubeRay Operator,将 Ray 集群声明式部署到 K8s:

apiVersion: ray.io/v1alpha1
kind: RayCluster
metadata:
  name: ml-cluster
spec:
  headGroupSpec:
    rayStartParams:
      dashboard-host: "0.0.0.0"
    template:
      spec:
        containers:
          - name: ray-head
            image: rayproject/ray:latest
            resources:
              limits:
                memory: "32Gi"
  workerGroupSpecs:
    - groupName: gpu-workers
      replicas: 4
      template:
        spec:
          containers:
            - name: ray-worker
              resources:
                limits:
                  nvidia.com/gpu: "1"

5.2 故障容错策略

场景 Ray 行为 应用层应对
Worker 崩溃 任务自动重试(max_retries 配置) 幂等设计,checkpoint 幂等
Actor 崩溃 Actor 不可恢复(默认行为) 持久化状态到外置存储
节点下线 GCS 标记 drain,等待任务安全迁移 K8s 优雅停机
Head 节点故障 内置 GCS 高可用(v2.4+ Raft 协议) 前置 VIP / ALB

5.3 生产调优清单

资源分配:避免 CPU 0.1 这种细粒度 request,容易引起 OOM;DL 场景推荐 integer GPU 分配。

对象存储内存:object_store_memory 默认物理内存 30%,Spill 路径配置高速 NVMe 以避免 GPU 空闲等待。

Dashboard 安全:生产环境暴露 --dashboard-host=0.0.0.0 需配合内网或 Sidecar 代理认证(Ray 2.6+ 引入 token 认证)。

GCS RPC 长连接:大规模集群(>500 节点)跳过外部 Redis,使用内置 GCS(v2.0+),内置 Raft 协议一致性保证。

六、与其他框架的对比选型

Ray vs Spark

Ray 适合有状态在线 serving、Actor 模型时间驱动型任务;Spark 适合 ETL 批量数据处理。两者可互补:Spark 预处理 → Ray 在线推理。

Ray vs K8s + 框架

K8s 是容器编排基础设施,Ray 是计算框架。Ray 运行于 K8s 之上(KubeRay Operator),负责任务级秒级调度,K8s 负责集群级节点管理和故障替换。

Ray vs Dask

Dask 强在 DataFrame/ML 生态;Ray 生态更成熟,生产级 Serve、Train、RLlib 实战验证更多。新项目建议优先考虑 Ray。

七、实战观点

Ray 正在成为 AI 工程化主流选择:AnyScale 公司提供商业化支持,Spark-RSS 引擎也在向 Ray 迁移。如果你在做 LLM 推理、多模型在线服务、或者大规模在线服务,Ray 是值得重点投入的基础设施。

当前最值得关注的方向:

  1. Ray + vLLM 推理服务:request 级别调度 + continuous batching 组合,推理吞吐提升 2-4x
  2. Ray Data 流式预处理:实时数据流水线替代 Lambda 架构
  3. Ray Train FSDP 集成:DeepSpeed ZeRO-3 平替方案,开箱即用
  4. Ray Dashboard 可观测:替代 Prometheus + Grafana 部分场景,任务级细粒度成本归因

Ray 的核心价值在于:一个引擎解决 AI 从训练到服务的全链路分布式需求,避免 Spark + TorchServe + Tecton 多栈技术带来的运维复杂度。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部