Day 13 / 共 20 天 · 第 3 周 数据与评测

实验主链路:把题库、被测对象、裁判拧到一起

评测集(题库)和评估器(裁判)都齐了,今天看谁把它们组装成一次可执行的"实验"(Experiment)——backend/modules/evaluation/application/experiment_app.go 是应用层入口,真正一条条跑数据靠 backend/cmd/consumer.go 启动的一串 RocketMQ 消费者。这是本仓库"异步任务"设计的最典型样本。

📍 你在整门课的位置 · 第 3 周 数据与评测(共 4 周 · 20 天)
D11 Dataset D12 评测集+评估器 D13 实验主链路 D14 Trace 观测 D15 跨模块组装图
L01

实验:把三件套拧成一次可执行的任务

🤔 痛点:光有题库和裁判还不能自动跑起来 你需要指定"用哪份题库(EvalSet)"、"测哪个对象(Target,可以是一个 Prompt、一个 Agent)"、"用哪几个裁判打分(Evaluators)",然后让系统把题库里每一条数据都喂给被测对象、再喂给每个裁判,最后汇总出一份报告。这整个过程就是一次实验(Experiment)
💡 本质:实验是一张"三元组的执行计划" Experiment 本身不产生新的评测能力,它只是把 Day 11/12 学到的三样东西"引用"起来,再加一层调度、重试、聚合统计的能力。
生活类比:一次"模拟考" 评测集是"试卷",被测对象是"考生",评估器是"改卷老师"。实验就是"组织一场模拟考"这件事本身——发卷子、监考、收卷、送去改、汇总分数出成绩单。组织考试本身也是个需要调度的活儿:几百个考生同时考、改卷老师可能改得慢,得有个监考系统盯着进度。
EvalSet 题库 Target 被测对象 Evaluators 裁判组 Experiment 实验
Experiment 是引用三者、驱动执行的"胶水层",不重新定义它们
L02

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 // 最终聚合分数
}
💡 全部"引用 ID",不复制数据 注意 EvalSetVersionIDTargetVersionIDEvaluatorVersionRef 全是 ID 引用,不是把题库/裁判的内容拷贝一份进实验表——这正是 Day 11/12 强调"版本号不可变"的价值:实验只要存一个版本号,就能保证下次重新查询时数据不会变,不需要靠"复制一份快照"来防止篡改。
状态含义是否终态
Pending已创建,排队等待调度
Processing正在跑
Terminating / Draining正在停止中(expt_run.go:75 判定为"停止中")
Success / Failed / Terminated / SystemTerminated跑完了(expt_run.go:79 判定为"终态")
L03

创建到提交:两个接口两步走

experiment_app.go:125CreateExperimentexperiment_app.go:467SubmitExperiment 分工明确:

CreateExperiment只做"登记":解析评估器版本 ID 列表、去重、组装 Experiment 记录,状态是 Pending,不真正开始跑数据
SubmitExperiment做校验(评估器不能重复、条目数不超过 100)+ 鉴权,内部再调 CreateExperiment 逻辑,然后真正投递调度事件,让实验开始跑
为什么要拆两步? 这和 Day 12 讲的 Run/Debug 拆分是同一种思路——"先攒草稿、再正式提交"在这个仓库里反复出现(数据集有草稿/版本、评估器有草稿/提交、实验有创建/提交)。好处是用户可以先配置好参数存个草稿,检查无误后再点击"开始运行",避免误触发一次昂贵的批量调用。
⚠️ SubmitExperiment 里的硬编码上限 experiment_app.go:472-474if len(req.ItemIds) > 100 { return error }——这是防止用户一次性点"只跑这些条目"选了过多数据导致单次请求过大/调度压力过大的保护性上限。读源码时留意这类"看似随意的数字",它们通常是真实压测/事故后加上的安全阀。
L04

为什么实验要用 RocketMQ,而不是直接同步跑

🤔 痛点:一次实验可能要跑几千条数据 × 好几个裁判 假设题库有 2000 条数据,配了 3 个评估器——那就是 2000 次"调被测对象"+最多 6000 次"调评估器打分"。每次调用可能是一次大模型请求,动辄几百毫秒到几秒。如果同步跑完再返回,一次 HTTP 请求要扛住几十分钟——网关、浏览器都受不了。
💡 本质:把"大批量、慢、可能失败重试"的活拆成消息,交给后台慢慢消化 SubmitExperiment 只是往 RocketMQ 扔一条"实验调度"消息就立刻返回;真正"跑一条数据"、"跑完一条后算聚合分"、"实验该结束了"分别对应不同的消费者,各自独立、可以水平扩容、可以单独重试失败的那一条,不会因为第 1999 条失败就让前 1998 条也白跑。
生活类比:流水线 而不是 一个人从头做到尾 想象工厂加工 2000 个零件:不是一个工人从原料到成品全程盯着(那样任何一步出错都要重新等),而是分成"下料工位""打磨工位""质检工位",零件在工位间流转,每个工位专心做自己的事、出错了只重做那一个零件。RocketMQ 就是这些工位之间传递零件的"传送带"。
为什么是 RocketMQ 而不是 Redis 队列/内存 channel? RocketMQ 提供了"消息持久化"(服务重启不丢消息)、"消费失败自动重试"、"多消费者组水平扩容"这些内存队列做不到(或做起来很脆弱)的能力,这正是"跑几千次可能失败的外部调用"这种场景最需要的。
L05

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 事件收尾。每个消费者只关心"收到什么消息、该干什么活、干完发什么新消息",彼此不直接调用对方的函数,这正是消息驱动架构和"直接函数调用"最大的区别。

L06

一条数据的完整旅程(动态视角)

把自己想象成题库里的第 1 条数据 跟着下面这条线走一遍,比记住六个消费者的名字更重要:
用户点提交 Scheduler 拆任务 RecordEval 调 Target RecordEval 调 Evaluator 写 ExptTurnResult AggrCalculate 更新统计 全部完成→Lifecycle 收尾

落库的核心结果表是 ExptTurnResultdomain/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/无值)时,裸值类型天然只能表示两态,得用指针或额外的状态字段。
L07

今日小结 + 动手 + 预告

🧠 今天你应该能回答

  • Experiment 实体本身存的是数据还是"引用关系"?为什么?
  • CreateExperiment 和 SubmitExperiment 各自做什么,为什么要拆开?
  • 为什么实验要用 RocketMQ 而不是同步循环跑完?
  • 六个消费者分别负责实验生命周期的哪一段?
  • WeightedScore 为什么用指针类型?
🎵 记忆口诀三件套引用不复制、创建提交两步走、大批量慢活走异步、六个消费者接力跑」——这是本仓库"异步任务"设计的最完整样本,Day 14 的 Trace 上报、Day 11 的数据导入都是这个思路的简化版。

✋ 动手 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
明天预告 · Day 14:实验跑的时候,被测对象和评估器内部到底发生了什么?我们需要"看得见"每一次模型调用——今天认识 Trace/Span 可观测性模块,看 SDK 上报的一条 Span 怎么流进 ClickHouse,又怎么被前端查出来画成火焰图。
← 上一天 Day 12 下一天 · Trace 观测 →