数据截至 (上游 commit d3b797b4c5cc)
评测子系统 — 实验执行引擎
本章讲什么: 一次 experiment(实验)从被提交到出分,中间到底发生了什么。重点是"执行引擎"那条主线:调度器怎么把一个实验拆成成千上万个小执行单元,每个单元怎么跑 target、跑 evaluators,分数怎么加权聚合。
1. 先建直觉:实验 = 一张大表逐格填充
把一次实验想成一张二维表:
- 行 = 评测集里的每个 item(一条数据),item 里又可能有多个 turn(多轮对话的每一轮)。
- 列 = 实验绑定的每个 evaluator(评分员)。
- 每格 = 这个 evaluator 对"这一行的(输入, target 输出)"打出的分。
实验要做的事,就是把这张表逐格填满,然后把每行的格子加权成一个"行分",再把所有行分聚合成"实验分"。
这张表很大(数据集动辄几千行),所以执行是异步、并发、可中断、可重试的——这正是本章所有复杂度的来源。
2. 核心数据模型:Experiment 聚合根
Experiment 是评测的聚合根,一个实例 = 一次评测的全部配置 + 运行态:
// 示意,非源码:只保留最能说明问题的字段
type Experiment struct {
ID int64
EvalSetID int64 // 用哪个评测集
TargetID int64 // 被测对象(可为空:仅记录型)
Evaluators []*Evaluator // 一组评分员
EvalConf *EvaluationConfiguration // 字段映射、并发数、重试数…
Status ExptStatus // 状态机:Pending→Processing→Success/Failed/Terminated
ExptType ExptType // Offline(离线批量) 还是 Online(在线/仅记录)
Stats *ExptStats // 跑了多少、成功多少、失败多少
}
真实定义在 modules/evaluation/domain/entity/expt.go:150(type Experiment struct)。几个关键点:
- 状态机枚举在
expt.go:56-74:ExptStatus_Pending(2) → Processing(3) → Success(11)/Failed(12)/Terminated(13)/SystemTerminated(14),外加流式的Draining(21)与Terminating(15)。判定函数IsExptFinished/IsExptFinishing在expt_run.go:74-80。 - EvalConf 里最值得注意的是
Connector(expt.go:268):它用 FieldAdapter / FieldConf(expt.go:355-363)描述"评测集的哪个字段喂给 target 的哪个入参"、"target 的哪个输出字段喂给 evaluator"。这套字段映射是把异构数据对齐的关键。 - AsyncExec()(
expt.go:217):若 target 或某个 evaluator 是异步的(如 WebAgent、Agent 评估器),整条执行走异步回报路径。
3. 执行的最小单元:一个 turn 怎么被评
抛开调度,先看最里层——对单独一行(一个 turn)打分。这是 DefaultExptTurnEvaluationImpl.Eval,逻辑清爽得像教科书:
// 示意,非源码:抓主干
func Eval(etec *ExptTurnEvalCtx) *ExptTurnRunResult {
targetResult := CallTarget(etec) // ① 跑被测对象,拿输出
if 失败或需中止 { return }
evalResults := CallEvaluators(etec, targetResult) // ② 每个评估器并发打分
return 装好 target+evaluators 结果
}
真实实现 modules/evaluation/domain/service/expt_run_item_turn_impl.go:66(Eval)。重点看两步:
① CallTarget —— 拿被测输出
CallTarget(expt_run_item_turn_impl.go:105)先判断要不要跳过:skipTargetNode(同文件 :148)在三种情况下跳过——没绑 target、实验是 ExptType_Online、target 是"仅记录型"。这对应实体里 EvalTargetType.IsRecordOnlyType()(target.go:85,那些 *Online 类型)。"仅记录型"是个巧妙设计:在线评测时被测对象的输出早已存在于 trace 里,不需要重新执行,只记录类型即可。
真正执行时(callTarget,:193),它按 target 类型用 FieldConf 把评测集字段映射成 target 入参(buildEvalSetFields,:700),同步对象走 ExecuteTarget, 异步对象(WebAgent 等)走 AsyncExecuteTarget 并把上下文存进 evalAsyncRepo 等回报(:288-309)。
② CallEvaluators —— 并发打分
CallEvaluators(:312)的精华在并发与正确性:
// 示意,非源码:每个评估器一个 goroutine,受并发池限流
pool := goroutine.NewPool(evaluatorsConf.GetEvaluatorConcurNum()) // 默认 3
for _, ev := range expt.Evaluators {
inputData := buildEvaluatorInputData(...) // 按字段映射组装输入
inputCopy := deepCopyEvaluatorInputData(inputData) // 关键:深拷贝
pool.Add(func() error {
record := evaluatorService.RunEvaluator(baseRunReq) // 真正打分
recordMap.Store(ev.versionID, record)
return err
})
}
pool.ExecAll()
真实在 callEvaluators(:381)。两个容易踩的坑、它都处理了:
- 闭包捕获(
:438-443):循环变量必须先复制到局部,否则 goroutine 读到最后一次迭代的值。 - 深拷贝输入(
deepCopyEvaluatorInputData,:655):多个 evaluator 并发时,保存记录会就地裁剪大字段(content_omitted),若共享同一个*Content指针,先跑完的 evaluator 会污染没跑的。所以每个 evaluator 拿一份深拷贝。这条注释(:440-441)是真实踩坑的结晶。
此外还支持评估器劫持(ShouldInterceptEvaluator,:446):某些情况下根据输入前置判断直接给结果、跳过真实执行;以及异步评估器(asyncCallEvaluator,:524):调用后把上下文存 evalAsyncRepo,等结果回报。
4. 调度器:把一个实验拆成无数个 turn
上一节是"一行怎么评"。但谁来决定"评哪些行、按什么顺序、失败了怎么办"?答案是调度器,它是 MQ 事件驱动的。
6 种运行模式
实验不是只有"从头跑"一种姿势。ExptRunMode(expt_run.go:21-37)定义了 6 种:
| 模式 | 含义 |
|---|---|
EvaluationModeSubmit | 创建后首次提交全量跑 |
EvaluationModeTrialRun | 试运行(跑少量条目预览) |
EvaluationModeFailRetry | 失败后重试所有失败项 |
EvaluationModeAppend | 追加模式(评测集新增了数据) |
EvaluationModeRetryAll | 重跑全部 |
EvaluationModeRetryItems | 重跑指定 item |
工厂 DefaultSchedulerModeFactory.NewSchedulerMode(expt_run_scheduler_mode_impl.go:120)按 mode 分发到不同 exec 实现。TrialRun 直接内嵌复用 Submit(:147 ExptTrialRunExec 组合 *ExptSubmitExec),是典型的"小变体靠组合复用"。