实验主链路:把题库、被测对象、裁判拧到一起
评测集(题库)和评估器(裁判)都齐了,今天看谁把它们组装成一次可执行的"实验"(Experiment)——backend/modules/evaluation/application/experiment_app.go 是应用层入口,真正一条条跑数据靠 backend/cmd/consumer.go 启动的一串 RocketMQ 消费者。这是本仓库"异步任务"设计的最典型样本。
实验:把三件套拧成一次可执行的任务
Experiment 本身不产生新的评测能力,它只是把 Day 11/12 学到的三样东西"引用"起来,再加一层调度、重试、聚合统计的能力。Experiment 实体拆解
backend/modules/evaluation/domain/entity/expt.go:150-198 的真实字段:
type Experiment struct {
ID int64
SpaceID int64
Name string
EvalSetVersionID int64 // 引用哪个版本的题库
EvalSetID int64
TargetType EvalTargetType // 被测对象类型:Prompt / Agent / …
TargetVersionID int64
EvaluatorVersionRef []*ExptEvaluatorVersionRef // 引用哪几个评估器版本(可以是多个裁判)
EvalConf *EvaluationConfiguration // 打分权重等配置
Status ExptStatus // Pending / Processing / Success / Failed / Terminated …
LatestRunID int64 // 支持"重跑",每次跑生成一个 RunID
TriggerType string // manual / openapi / schedule —— 谁触发的
Stats *ExptStats // 实时统计:跑了多少、成功多少
AggregateResult *ExptAggregateResult // 最终聚合分数
}
EvalSetVersionID、TargetVersionID、EvaluatorVersionRef 全是 ID 引用,不是把题库/裁判的内容拷贝一份进实验表——这正是 Day 11/12 强调"版本号不可变"的价值:实验只要存一个版本号,就能保证下次重新查询时数据不会变,不需要靠"复制一份快照"来防止篡改。| 状态 | 含义 | 是否终态 |
|---|---|---|
Pending | 已创建,排队等待调度 | 否 |
Processing | 正在跑 | 否 |
Terminating / Draining | 正在停止中(expt_run.go:75 判定为"停止中") | 否 |
Success / Failed / Terminated / SystemTerminated | 跑完了(expt_run.go:79 判定为"终态") | 是 |
创建到提交:两个接口两步走
experiment_app.go:125 的 CreateExperiment 和 experiment_app.go:467 的 SubmitExperiment 分工明确:
CreateExperiment只做"登记":解析评估器版本 ID 列表、去重、组装 Experiment 记录,状态是 Pending,不真正开始跑数据SubmitExperiment做校验(评估器不能重复、条目数不超过 100)+ 鉴权,内部再调 CreateExperiment 逻辑,然后真正投递调度事件,让实验开始跑SubmitExperiment 里的硬编码上限
experiment_app.go:472-474:if len(req.ItemIds) > 100 { return error }——这是防止用户一次性点"只跑这些条目"选了过多数据导致单次请求过大/调度压力过大的保护性上限。读源码时留意这类"看似随意的数字",它们通常是真实压测/事故后加上的安全阀。为什么实验要用 RocketMQ,而不是直接同步跑
SubmitExperiment 只是往 RocketMQ 扔一条"实验调度"消息就立刻返回;真正"跑一条数据"、"跑完一条后算聚合分"、"实验该结束了"分别对应不同的消费者,各自独立、可以水平扩容、可以单独重试失败的那一条,不会因为第 1999 条失败就让前 1998 条也白跑。consumer.go 与六个消费者:谁负责哪一步
服务入口 backend/cmd/consumer.go 启动一个独立进程,专门跑消息消费者(和 HTTP 服务是两个不同的可执行程序):
func MustInitConsumerWorkers(
cfactory conf.IConfigLoaderFactory,
mqFactory mq.IFactory,
experimentApplication exptapp.IExperimentApplication,
datasetApplication dataapp.IJobRunMsgHandler, // 记得 Day 11 的 Job?消费者就是它的下游
...
) []mq.IConsumerWorker {
workers, err := evalconsumer.NewConsumerWorkers(loader, experimentApplication, publisher)
res = append(res, workers...)
workers, err = dataconsumer.NewConsumerWorkers(cfactory, datasetApplication) // Day 11 的数据导入导出也在这里消费
...
}
真正的六个实验消费者在 backend/modules/evaluation/infra/mq/rocket/consumer/consumer.go:25-32:
ExptSchedulerEventConsumer
接到"开始跑实验"的调度事件,把题库条目拆分成一条条子任务再发出去。
ExptRecordEvalEventConsumer
真正"跑一条数据":调被测对象拿结果,再调评估器打分,写入结果表。
ExptAggrCalculateEventConsumer
某条/若干条跑完后,重新计算实验的聚合统计(平均分、通过率等)。
ExptTurnResultFilterEventConsumer
维护结果的筛选加速索引,方便前端"按分数筛选"能快速查询。
ExptExportEventConsumer
把实验结果导出成 CSV,异步生成下载文件(大文件同理不能同步等)。
ExptLifecycleEventConsumer
处理实验生命周期事件(开始/结束/终止),还负责失败时的 webhook 重试通知。
👶 小白问:这六个是不是一个套着一个链式调用?
👨🏫 老师:对,是一条事件驱动的流水线:Scheduler 拆任务 → 发消息 → RecordEval 消费一条条任务 → 每条跑完发"该更新统计了"的事件 → AggrCalculate 消费更新统计 → 全部跑完触发 Lifecycle 事件收尾。每个消费者只关心"收到什么消息、该干什么活、干完发什么新消息",彼此不直接调用对方的函数,这正是消息驱动架构和"直接函数调用"最大的区别。
一条数据的完整旅程(动态视角)
落库的核心结果表是 ExptTurnResult(domain/entity/expt_result.go:292-306),一条记录对应"一条数据在一次实验里的评测结果":
type ExptTurnResult struct {
ExptID int64
ExptRunID int64 // 支持重跑,每次跑对应不同 RunID
ItemID int64
TurnID int64 // 多轮对话场景下,第几轮
Status int32
TargetResultID int64 // 被测对象这一条的回答
EvaluatorResults *EvaluatorResults // 每个评估器给的分数
WeightedScore *float64 // 加权后的最终分(指针:nil=未算,非 nil=已算)
}
*float64 而不是 float64 的深意
WeightedScore 用指针类型是为了区分"分数是 0"和"还没算出分数"两种情况——如果用普通 float64,初始值默认就是 0,你没法区分"真的打了 0 分"还是"压根还没跑完"。这是 Go 里一个常见但容易被忽视的坑:需要"三态"语义(有值A/有值B/无值)时,裸值类型天然只能表示两态,得用指针或额外的状态字段。今日小结 + 动手 + 预告
🧠 今天你应该能回答
- Experiment 实体本身存的是数据还是"引用关系"?为什么?
- CreateExperiment 和 SubmitExperiment 各自做什么,为什么要拆开?
- 为什么实验要用 RocketMQ 而不是同步循环跑完?
- 六个消费者分别负责实验生命周期的哪一段?
WeightedScore为什么用指针类型?
✋ 动手 5 分钟(可选)
# 1. 看实验状态机的完整枚举
grep -n "ExptStatus_" backend/modules/evaluation/domain/entity/expt.go
# 2. 看六个消费者是怎么注册的
sed -n '19,33p' backend/modules/evaluation/infra/mq/rocket/consumer/consumer.go
# 3. 找到独立的消费者进程入口
cat backend/cmd/consumer.go
# 4. 看结果表的真实字段
grep -n "type ExptTurnResult struct" -A20 backend/modules/evaluation/domain/entity/expt_result.go