Skip to content

Object Store 实现

Object Store 是 Mini Ray 的第一块地基。没有它,.remote() 就只能返回线程池 Future,无法表达 Ray 风格的对象依赖。

数据结构

每个对象有一个 _ObjectEntry

python
class _ObjectEntry:
    def __init__(self) -> None:
        self.ready = threading.Event()
        self.value = None
        self.error = None
        self.created_at = time.time()

关键是 ready。任务提交时先创建 entry,但值还没有计算出来。任务完成后设置值并唤醒等待者。

reserve 与 put

reserve 创建一个未完成对象:

python
def reserve(self) -> ObjectRef:
    object_id = uuid.uuid4().hex
    self._entries[object_id] = _ObjectEntry()
    return ObjectRef(object_id)

put 创建一个已完成对象:

python
def put(self, value):
    ref = self.reserve()
    self.set_result(ref, value)
    return ref

这个区别很重要:

  • Task 返回值:先 reserve,执行完成后 set_result
  • 用户 ray.put:对象已经在 Driver 手里,可以马上 ready。

get

get 的核心逻辑是等待:

python
def _get_one(self, ref, timeout=None):
    entry = self._entry(ref)
    if not entry.ready.wait(timeout=timeout):
        raise TimeoutError(...)
    if entry.error is not None:
        raise MiniRayTaskError(...) from entry.error
    return entry.value

这解释了为什么远程异常会在 get 时出现。任务执行线程只是把异常记录进对象 entry,Driver 取结果时才看到。

wait

wait 不需要取值,只需要知道哪些对象 ready:

python
ready = [ref for ref in refs if self._entry(ref).ready.is_set()]

真实 Ray 的 wait 要处理分布式对象位置和事件通知。Mini Ray 用轮询实现,足够说明 API 语义。

依赖解析

当任务参数里包含 ObjectRef,Runtime 会递归解析:

python
def resolve(self, value):
    if isinstance(value, ObjectRef):
        return self.get(value)
    if isinstance(value, tuple):
        return tuple(self.resolve(item) for item in value)
    if isinstance(value, list):
        return [self.resolve(item) for item in value]
    if isinstance(value, dict):
        return {key: self.resolve(item) for key, item in value.items()}
    return value

这让下面代码成立:

python
values_ref = ray.put([1, 2, 3])
total_ref = total.remote(values_ref)

total 函数拿到的是列表,而不是 ObjectRef

和真实 Ray 的差异

Mini Ray 递归解析容器里的 ObjectRef 是为了教学直观。真实 Ray 对嵌套引用、所有权、生命周期和序列化有更细的规则。

小练习

给 Object Store 增加一个 delete(ref) 方法,然后思考:

  • 如果有人正在 get(ref),删除应该怎么处理?
  • 如果某个 pending task 依赖这个 ref,删除是否安全?
  • 真实 Ray 为什么需要引用计数?

面向学习目的的 Ray Core 中文导读与 Mini Ray 机制预览。