01 · 控制面 LightningStore 与数据模型
本章讲什么: 整个 Agent Lightning 围着一本中央账本转,这本账本就是
LightningStore。搞懂它记了哪四种数据、每种数据的状态怎么变,你就抓住了这个项目的骨架——后面所有章节都是在往这副骨架上挂肉。
1. 为什么需要一本中央账本
先回到问题。训练侧(算法)和执行侧(Runner 跑 agent)通常在不同进程、甚至不同机器上:算法可能在带 GPU 的训练节点,Runner 可能是十几个并行 worker。它们要交换两样东西:
- 执行侧要往训练侧送遥测(这趟 agent 干了啥、得了多少分)。
- 训练侧要往执行侧送新资源(这批任务该用哪个模型/哪份提示词)。
如果让两侧直接互相调用,就会耦合死:谁先启动、谁认识谁、并发怎么办,全是麻烦。Agent Lightning 的选择是插一个中间层——LightningStore。两侧都只跟 store 打交道,谁也不认识谁。
它同时扮演三个角色:
| 角色 | 含义 |
|---|---|
| 数据库 | 持久化任务、尝试、span、资源 |
| 消息队列 | 任务排队(enqueue)+ 领取(dequeue),FIFO |
| 状态机引擎 | 驱动 rollout/attempt 的状态流转、超时、重试 |
契约定义在 agentlightning/store/base.py:104 的 LightningStore 基类(一堆 async def ... raise NotImplementedError()),具体实现有内存版、SQLite 版、Mongo 版(agentlightning/store/memory.py、sqlite.py、mongo.py)。默认用内存版(依据:trainer/trainer.py:300 InMemoryLightningStore)。
2. 四种核心数据
store 里就四类东西。一句话先记住:
| 数据 | 白话 | 类型定义 |
|---|---|---|
| Rollout | 一道「题」——让 agent 在一个输入上跑一趟的请求 | types/core.py:174 Rollout |
| Attempt | 一次「答题尝试」——同一道题失败可重试,每次是一个 attempt | types/core.py:136 Attempt |
| Span | 一个「遥测事件」——一次 LLM 调用、一次工具调用、一次打分 | types/tracer.py:251 Span |
| Resources | 一份「可训练的东西」——LLM 端点、提示词模板,带版本号 | types/resources.py:192 ResourcesUpdate |
下面逐个拆。
2.1 Rollout:一道题
Rollout 是一次 agent 运行的请求单元。核心字段(依据:types/core.py:174-199):
rollout_id:唯一 id。input:任务输入,类型是TaskInput = Any——任意载荷,一道数学题、一段 SQL 需求都行(依据:types/core.py:259)。mode:"train"/"val"/"test",给下游分析用(依据:types/core.py:132RolloutMode)。resources_id:这趟该用哪个版本的资源。status:当前状态(见下节状态机)。config:重试/超时策略(RolloutConfig,types/core.py:161,含max_attempts、retry_condition)。
2.2 Attempt:一次尝试
为什么 rollout 之外还要 attempt?因为一道题可能要重试。Rollout 是「逻辑上的题」,Attempt 是「物理上的某一次跑」。一道题重试三次,就是一个 rollout 挂三个 attempt,sequence_id 从 1 递增(依据:types/core.py:142 sequence_id)。
Attempt 记的是执行细节(依据:types/core.py:136-158):attempt_id、sequence_id、start_time/end_time、status、worker_id(哪个 worker 在跑)、last_heartbeat_time(最后一次报活时间——用来判断卡死)。
有个方便的组合类型 AttemptedRollout(types/core.py:202)——就是「Rollout + 当前 attempt」打包,Runner 领到的就是 它,一手拿到题面和尝试 id。
2.3 Span:一个遥测事件
Span 借用了 OpenTelemetry(一套业界标准的分布式追踪规范) 的概念:agent 每做一件可观测的事,就产生一个 span。一次 openai.chat.completion 调用是一个 span,一次工具调用是一个 span,emit_reward 打一次分也是一个 span。
Agent Lightning 在标准 OTel span 上加了三个「定位」字段,把 span 钉到具体的题和尝试上(依据:types/tracer.py:262-266):
rollout_id:属于哪道题。attempt_id:属于哪次尝试。sequence_id:严格递增的序号,用来给同一尝试内的 span 排一个全序。
为什么要单独一个 sequence_id 而不靠时间戳?因为 span 可能来自不同机器、时钟不同步,时间戳不可靠;序号由 store 统一发放(get_next_span_sequence_id,store/base.py:555),才能保证顺序确定。这在第 03 章的 triplet 提取里是命根子。
2.4 Resources:可训练的东西
这是 Agent Lightning 一个关键抽象:「要被优化的东西」被统一建模成「资源」。目前有三种(依据:types/resources.py):
| 资源 类型 | 是什么 | 定义 |
|---|---|---|
LLM | 一个 OpenAI 兼容的模型端点(endpoint + model + 采样参数) | types/resources.py:43 |
ProxyLLM | 会经过 LLMProxy 改写端点的 LLM,把 rollout/attempt 信息塞进 URL | types/resources.py:65 |
PromptTemplate | 一份提示词模板(template + 引擎,如 f-string) | types/resources.py:146 |
多个资源用名字组织成 NamedResources(就是 Dict[str, 资源],types/resources.py:172),比如 {"main_llm": LLM(...), "system_prompt": PromptTemplate(...)}。
资源是不可变、带版本的。每次算法优化出新东西,不是原地改,而是发布一个新快照 ResourcesUpdate(带 resources_id、version、create_time,types/resources.py:192)。store 永远记着「最新版是哪个」,Runner 领任务后去取当前该用的版本。这样 RL 训练里「第 5 轮用的是哪版模型」永远可追溯。
一句话把四者串起来: 算法发布一版 resources → 排一批 rollout → Runner 领走建 attempt → 跑出一串 span → 算法读 span 学习 → 发布下一版 resources。
3. 两套状态机(这是 store 的灵魂)
store 不只是存数 据,它驱动状态流转。有两套状态机,一套管 rollout,一套管 attempt。
3.1 Rollout 状态流转
七个状态(依据:types/core.py:109-117 RolloutStatus)。怎么读下图:从左到右是正常生命周期,向下的分支是异常/重试。
enqueue dequeue 收到首个 span 正常结束
queuing ───────────▶ preparing ─────────▶ running ───────────▶ succeeded
▲ │ │ failed
│ │ │ cancelled
│ requeuing ◀────────┴────────────────────┘ (满足 retry_condition 时重排队)
└──────────────┘
三个辅助判断函数把七状态归成三类(依据:store/base.py:27-36):
is_queuing:queuing或requeuing——在队列里等着。is_running:preparing或running——正在跑。is_finished:succeeded/failed/cancelled——到终点了。
3.2 Attempt 状态流转
Attempt 的状态更「物理」,没有排队/取消这类调度态,只有一个进程的生死(依据:types/core.py:120-129 AttemptStatus):preparing → running → succeeded/failed,外加两种「疑似出事」:unresponsive(worker 一段时间没报活)和 timeout(还在报活但跑太久)。
3.3 状态怎么被推动
关键在于:很多状态转移不是显式调用,而是「副作用」。最典型的——span 到达即心跳。往 store add_span 一个 span,会顺手做三件事(依据:store/base.py:336-341 add_span 文档契约):
- 校验对应的 rollout/attempt 存在。
- 更新该 attempt 的
last_heartbeat_time(证明 worker 还活着)。 - 如果 rollout/attempt 还停在
preparing/requeuing,把它们推进到running。
也就是说,「第一个 span 到了」= 「这题真的开跑了」。这个设计让「进度」和「数据」共用一条通道,不需要额外的心跳协议来判断卡死。
另一处联动:update_attempt 会根据 attempt 的新状态反向更新 worker 状态(成功/失败→worker 空闲 idle;卡死/超时→worker 未知 unknown;否则忙 busy,依据:store/base.py:776-780)。
4. store 的接口长什么样(按用途分组)
把 LightningStore 上几十个方法按「谁会用」归一下类,就不吓人了:
| 用途 | 代表方法 | 谁调 |
|---|---|---|
| 排队 / 领取任务 | enqueue_rollout / dequeue_rollout | 算法排,Runner 领 |
| 直接开跑(跳过队列) | start_rollout / start_attempt | Runner 单步调试、在线 RL |
| 落 span | add_span / add_otel_span / add_many_spans | Runner |
| 发号(span 序号) | get_next_span_sequence_id | Runner |
| 发布 / 读资源 | add_resources / update_resources / get_latest_resources | 算法发,Runner 读 |
| 查询(做看板用) | query_rollouts / query_spans / query_attempts | 算法、仪表盘 |
| 等待完成 | wait_for_rollouts | 算法(阻塞到一批题跑完) |
| 心跳 | update_worker | Runner 定期报活 |
依据:以上方法均定义在 store/base.py,如 enqueue_rollout(:200)、dequeue_rollout(:250)、add_span(:330)、get_latest_resources(:543)、add_resources(:681)、wait_for_rollouts(:593)。
一个易漏的细节:start_rollout vs enqueue_rollout
两个都是「新建一道题」,但语义不同(依据:store/base.py:158 与 :200 的 docstring):
enqueue_rollout:只排队,状态queuing,不建 attempt,等着谁来领。批量训练用这个。start_rollout:立即开跑,直接建好第一个 attempt(sequence_id=1),状态preparing,调用方不必走公共队列。单步执行/在线场景用这个(Runner 的step()就走它,见第 02 章)。
5. 能力与统计(实现可选声明)
store 实现可以声明自己的能力(LightningStoreCapabilities,store/base.py:59):是否线程安全、是否异步安全、是否零拷贝(跨线程只有一份)、是否支持 OTLP/HTTP 收 span。比如客户端-服务端执行策略下,内存 store 会被要求 thread_safe=True(依据:trainer/trainer.py:300)。这让上层能按部署形态挑合适的 store,而不必写死。
小结: LightningStore = 数据库 + 队列 + 状态机三合一,管 rollout/attempt/span/resources 四种数据。记住两句话——「span 到达即心跳并推动状态」、「资源不可变、按版本发布」——就抓住了它的精髓。下一章看执行侧怎么真正把一道题跑出来、把 span 生出来。