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 为什么需要引用计数?