Skip to main content

mini-ray:手写一个单机版 Ray 运行时

项目信息

GitHubNEDONION/mini-ray

C++17 CMake 3.15+ pybind11 POSIX 共享内存 Python 3.8+ PyTorch (示例)

定位:面向学习的 alpha 项目,只支持单机多进程,不是 Ray 的替代品,也不与 Ray 兼容。

单机多进程推理演示

1 为什么要造这个轮子

直接读 Ray 源码的问题是:gRPC、GCS、Raylet、Plasma、容错、资源调度全都缠在一起,看不出哪些是核心抽象、哪些是工程细节。

mini-ray 的做法是把分布式那一维砍掉,只保留单机多进程,然后完整保留 Ray 的核心分层:

这个分层不是随便定的,而是照着真实 Ray 的设计哲学

层次语言职责为什么是这个语言
API 层Python用户接口、装饰器、序列化易用性、灵活性
绑定层pybind11Python ↔ C++ 类型转换跨语言桥梁
核心层C++调度、存储、通信、资源管理性能,且绕开 GIL

2 最小 API

整个用户界面只有四个函数,和 Ray 一模一样:

import miniray

miniray.init(num_workers=2)

@miniray.remote
def fibonacci(n):
if n <= 1:
return n
a, b = 0, 1
for _ in range(n - 1):
a, b = b, a + b
return b

result = miniray.get(fibonacci.remote(10))
print(result) # 55

miniray.shutdown()

fibonacci.remote(10) 不返回结果,返回一个 ObjectRef——这是理解 Ray 的第一个关键点:提交任务和获取结果是两件事,中间这段时间可以继续提交别的任务,并行度就是这么来的。

3 一次任务提交的完整链路

三个设计细节值得注意:

1. Worker 是 Pull 而不是 Push。 调度器不主动派发,空闲 Worker 主动来拿。这样天然做到负载均衡——快的 Worker 自然拿得多,不需要额外的负载探测。

2. C++ 只搬运不解析 Python 对象。 函数用 cloudpickle 在 Python 侧序列化成字节数组,C++ 侧只当二进制存储和传输:

struct Task {
// 序列化的函数(不反序列化,只传输)
std::vector<uint8_t> serialized_function;
std::vector<uint8_t> serialized_args;
};

这条边界划得很清楚:跨语言的地方只传字节,不传语义,避免了在 C++ 里实现 Python 对象模型的噩梦。

3. ObjectRef 是 128-bit UUID,本地生成。 不需要中心化的 ID 分配服务,提交任务时立刻就能返回引用。

4 核心组件

4.1 ObjectStore:为什么必须是 C++

class ObjectStore {
public:
ObjectRef Put(const std::shared_ptr<Buffer>& data);
std::shared_ptr<Buffer> Get(const ObjectRef& object_ref);
void Delete(const ObjectRef& object_ref);
bool Contains(const ObjectRef& object_ref);

private:
std::unordered_map<ObjectID, void*> objects_; // ObjectID → 共享内存地址
boost::interprocess::managed_shared_memory segment_;
};

为什么不用 Python 的 multiprocessing.Manager().dict() 这是这个项目最值得记住的一条对比:

Manager().dict()共享内存
访问方式每次访问都是一次 IPC 往返直接读内存
拷贝必然拷贝可零拷贝
内存布局不可控精确控制
GIL受限不受限

对于要传递大张量的 ML 负载,这个差距是数量级的。Ray 用 Plasma(Apache Arrow)也是同一个理由。

4.2 Scheduler:条件变量而不是轮询

class Scheduler {
private:
std::queue<Task> task_queue_;
std::mutex queue_mutex_;
std::condition_variable queue_cv_;

public:
void SubmitTask(const Task& task) {
std::lock_guard<std::mutex> lock(queue_mutex_);
task_queue_.push(task);
queue_cv_.notify_one(); // 唤醒一个等待的 Worker
}

std::optional<Task> GetTask() {
std::unique_lock<std::mutex> lock(queue_mutex_);
queue_cv_.wait(lock, [this] { return !task_queue_.empty(); });
Task task = task_queue_.front();
task_queue_.pop();
return task;
}
};

condition_variable 而不是 while(true) { sleep(0.01); check(); }:空闲时零 CPU 占用,有任务时微秒级唤醒。

4.3 CoreWorker:每个进程一个

class CoreWorker {
public:
ObjectRef SubmitTask(const TaskSpec& task_spec);
std::vector<std::shared_ptr<Buffer>> GetObjects(
const std::vector<ObjectRef>& object_refs);
ObjectRef Put(const std::shared_ptr<Buffer>& data);
ActorHandle CreateActor(const ActorCreationSpec& actor_spec);

private:
std::unique_ptr<ObjectStore> object_store_;
std::unique_ptr<TaskSubmitter> task_submitter_;
};

CoreWorker 是 Ray 架构里最容易被忽略但最重要的抽象:每个进程(包括驱动进程和 Worker 进程)都有一个,它让「提交任务」和「执行任务」用的是同一套接口。驱动进程和 Worker 进程在架构上是对等的,所以任务里可以再提交任务(嵌套并行)。

5 上层库:从运行时到能跑训练

有了 Task + ObjectRef + Actor,就可以往上叠库了。仓库里实现了 Ray 生态的几个简化对应物:

模块对应 Ray作用
miniray.actorRay Actor有状态对象,跨调用保持状态
miniray.psRay Parameter Server 范式参数服务器与同步策略
miniray.trainRay TrainDataParallelTrainer 数据并行训练
miniray.tuneRay Tune超参搜索、Trial 管理、结果分析
miniray.dashboardRay DashboardFlask + 前端的任务/资源监控

多进程训练

用它跑了一个 GAN 训练实验,验证整条链路在真实 ML 负载下能跑通:

GAN 实验

GAN 实验期间的 GPU 监控

GPU 不是必需条件,GAN 示例在 CUDA 不可用时会回退到 CPU。

6 和真实 Ray 的差距

诚实标注边界,比假装是生产级更有用:

特性mini-ray真实 Ray
架构Python API + C++ 核心 ✅Python API + Cython + C++ 核心
对象存储C++ 共享内存(简化版 Plasma)✅Plasma(Apache Arrow)
调度器C++ FIFO 调度器 ✅Raylet(分布式调度)
网络通信单机(共享内存 + Pipe)gRPC(分布式)
GCS全局控制存储
容错自动重试、容错
资源管理简单 Worker 池CPU/GPU/内存 调度
语言支持Python onlyPython, Java, C++

核心思想是一致的:Task 抽象、ObjectRef、CoreWorker 架构、Python/C++ 分层、Actor 模型。这五点学会了,再读 Ray 源码就有地图了。

7 沉淀下来的经验

1. 想读懂一个复杂系统,先砍掉一个维度重写它。 把「分布式」砍掉,Ray 的核心抽象立刻清晰了。分布式带来的复杂度(GCS、故障检测、网络分区)和核心抽象(Task/ObjectRef/Actor)是正交的。

2. 跨语言边界只传字节,不传语义。 C++ 侧完全不理解 Python 对象,只当二进制搬运。这条边界一旦模糊,就要在 C++ 里实现半个 CPython。

3. 「提交」与「取回」分离是并行的前提。 remote() 立即返回 ObjectRef 是整个模型的核心。如果它同步返回结果,就退化成了普通函数调用。

4. Pull 模式的调度天然负载均衡。 Worker 主动拉取,比调度器主动派发少了一整套负载探测逻辑。

5. 共享内存 vs IPC 是数量级差异。 在 ML 负载(大张量传递)下,Manager().dict() 的每次访问一次 IPC 是不可接受的。这就是 Ray 要自己做 Plasma 的原因。

6. C++ 核心的真正理由是 GIL,不只是性能。 调度器要在后台线程跑、要用条件变量精确唤醒,这些在 GIL 下都做不干净。

参考