跳到主要内容

数据截至 (上游 commit 36c7a7f6eca6)

评估引擎:evaluate() 内部的双线程池流水线与 Scorer 抽象

30 秒导读: mlflow.genai.evaluate(data, scorers, predict_fn) 把「一批测试数据」跑成「一批打了分的 trace」。 它内部不是一个 for 循环,而是两级线程池组成的流水线:预测池负责产出 trace,打分池负责给 trace 挂 Feedback, 中间用信号量做背压、用令牌桶做限流。这一章讲这条流水线的骨架和 Scorer 抽象; judge 内部怎么问 LLM 见 05,线上定时打分见 06


1. 这是什么(零基础也能懂)

一句话定义: 离线评估引擎——给定一批输入(或一批已有 trace)和一组打分器,批量跑出分数,并把分数写回每条 trace。

解决什么问题: 你改了 prompt、换了模型、加了一个工具,怎么知道「变好了还是变差了」? 手工点开十条对话看看不算数。评估引擎做三件事:

它做的事白话
批量跑把 50 条测试输入喂给你的 agent,每条产出一棵 trace
批量打分每棵 trace 交给若干「评审员」(scorer)打分
落库 + 汇总分数作为 Feedback 挂回 trace,再聚合成 run 级指标(correctness/mean 这种)

三种用法(本章第 3 节详解):

用法data 里有什么要不要 predict_fn
① 评已有 tracetrace 列(通常来自 mlflow.search_traces()不要
② 评静态数据inputs / outputs / expectations不要
③ 现跑现评只有 inputs(+ 可选 expectations

用起来什么样:

import mlflow
from mlflow.genai.scorers import Correctness, Safety, scorer

data = [
{"inputs": {"question": "What is MLflow?"}, "expectations": {"expected_facts": ["ML platform"]}},
]

@scorer # 自定义打分器:函数就是打分器
def fast_enough(trace) -> bool:
return trace.info.execution_duration < 1000

result = mlflow.genai.evaluate(
data=data,
predict_fn=my_agent, # 现跑现评
scorers=[Correctness(), Safety(), fast_enough],
)
print(result.metrics) # {'correctness/mean': 1.0, 'safety/mean': 1.0, ...}

一句话直觉: 把它当成一条流水线车间——左边一台机器负责「造零件」(调用你的 agent,产出 trace), 右边一台机器负责「质检」(跑 scorer 打分)。两台机器并行,中间放一个有容量上限的传送带(背压)。


2. 顶层全景(它大概怎么转)

2.1 一张图看全流程

怎么读这张图:从上到下是一次 evaluate() 的时间顺序,中间那段双线框才是真正并发的部分。

mlflow.genai.evaluate(data, scorers, predict_fn)


┌──────────────────────────────────────┐
│ ① 准备(单线程,base.py) │
│ · 校验 scorers / 校验数据列 │
│ · 各种输入 → 统一的 pandas DataFrame │
│ · 开 run、记 dataset input │
└──────────────────────────────────────┘
│ DataFrame

┌──────────────────────────────────────┐
│ ② DataFrame 每行 → 一个 EvalItem │
└──────────────────────────────────────┘
│ list[EvalItem]

╔══════════════════════════════════════╗
║ ③ 双线程池流水线(harness.py,本章重心)║
║ ║
║ 提交线程 ──背压信号量──▶ 预测池 ║
║ │ trace ║
║ ▼ ║
║ 打分池 ──▶ Feedback 写回 trace
║ ║
║ (全部单轮打完,再跑会话级 scorer) ║
╚══════════════════════════════════════╝
│ list[EvalResult]

┌──────────────────────────────────────┐
│ ④ 收尾(单线程) │
│ · trace 关联到 run、刷新 trace │
│ · 聚合指标 → mlflow.log_metrics │
│ · 拼 result_df → EvaluationResult │
└──────────────────────────────────────┘

2.2 部件一句话职责

部件干什么在哪个文件
evaluate / _run_harness公开入口 + 准备阶段(校验、归一化、开 run)mlflow/genai/evaluation/base.py:56:302
_convert_to_eval_set五花八门的 data → 统一 DataFramemlflow/genai/evaluation/utils.py:151
EvalItem / EvalResult一行数据 / 一行结果的载体mlflow/genai/evaluation/entities.py:78:219
harness.run评估主流程(建进度条、跑流水线、收尾聚合)mlflow/genai/evaluation/harness.py:661
_run_pipeline建两个池、跑生产者—消费者主循环mlflow/genai/evaluation/harness.py:521
_PredictSubmitter / _ScoreSubmitter预测池 / 打分池的封装mlflow/genai/evaluation/harness.py:236:384
RPSRateLimiter / call_with_retry令牌桶限流 + 429 重试mlflow/genai/evaluation/rate_limiter.py:70:168
Scorer / scorer打分器基类 + 装饰器mlflow/genai/scorers/base.py:290:1163
BuiltInScorer 家族内置打分器(多数是 LLM judge 的薄封装)mlflow/genai/scorers/builtin_scorers.py:314
compute_aggregated_metrics逐行分数 → run 级指标mlflow/genai/scorers/aggregation.py:26

2.3 主线走一遍(不进代码)

  1. 你传进来的 data(DataFrame / list[dict] / list[Trace] / Spark DF / 托管数据集)被压平成一个 DataFrame。
  2. 每一行变成一个 EvalItem:它同时是输入容器结果容器——预测阶段往里塞 outputstrace
  3. 预测池对每个 item 调 predict_fn(或复制已有 trace,或造一棵最小 trace)。
  4. 某个 item 一预测完,立刻被丢给打分池;打分池对它跑所有单轮 scorer。
  5. 单轮全部结束后,会话级 scorer 才按 session 分组跑(因为它需要成型的 trace 才能分组)。
  6. 分数以 Feedback 形式 log_assessment 到 trace 上,再聚合成 run 指标。

3. 入口层:三种输入形态如何收敛成一种

这节讲 evaluate() 在进入线程池之前做的所有准备工作。

3.1 evaluate 只是一层壳

evaluate 本身只有一行实质代码,把活全交给 _run_harnessmlflow/genai/evaluation/base.py:298):

result, _ = _run_harness(data, scorers, predict_fn, model_id)
return result

多出来的那个返回值是遥测数据——_run_harness@record_usage_event(GenAIEvaluateEvent) 装饰 (base.py:301),装饰器要拿它上报数据规模和字段(_get_eval_data_size_and_fieldsutils.py:91)。

3.2 准备阶段的六件事

_run_harnessbase.py:302)在一个 with 块里依次做:

顺序做什么关键调用
1校验 scorers 是实例不是类validate_scorersscorers/validation.py:26
2开 run + 设 active model + 打开自动埋点_start_run_or_reuse_active_runconfigure_autologging_for_evaluationbase.py:341-346
3会话级 scorer 与 predict_fn 互斥校验validate_session_level_evaluation_inputssession_utils.py:190
4data → DataFrame_convert_to_eval_setutils.py:151
5内置 scorer 要的列是否齐valid_data_for_builtin_scorersvalidation.py:113)——缺列只 log info,不报错
6记录 dataset 与 model 到 run_log_dataset_inputbase.py:442

注意第 2 步的注释:自动埋点必须在线程池外打开,否则多线程同时改 patch 会 race (base.py:344)。这条约束直接决定了「准备阶段单线程、执行阶段多线程」的分层。

第 6 步用 MlflowClient().log_inputs 把数据集(以及可选的 model_id)挂到 run 上, 这样 UI 里能看到「这次评估用的是哪个数据集、哪个模型版本」。

3.3 归一化:怎么把五种输入压成一种

_convert_to_eval_setutils.py:151)是个四段管道:

data ──▶ _convert_eval_set_to_df # list/DataFrame/Spark/托管数据集 → pandas
──▶ _deserialize_trace_column_if_needed # trace 列的 JSON 串/dict → Trace 对象
──▶ _extract_request_response_from_trace # 没有 inputs/outputs 就从 trace 根 span 取
──▶ _extract_expectations_from_trace # 没有 expectations 就从 trace 的 Expectation 断言取

后两段是「trace 形态」和「静态数据形态」合流的地方:只要有 trace,就能反推出 inputs/outputs/expectationsutils.py:226utils.py:268)。所以下游的 scorer 不必关心数据是哪来的。

一个真实的坑:如果 trace 是用 search_traces(..., include_spans=False) 拿的,就没有根 span, inputs/outputs 会是 None——代码专门为此打了一条带解决方案的 warning(utils.py:257-263)。

3.4 predict_fn 的包装

用法 ③ 传进来的 predict_fn 会先过 convert_predict_fnmlflow/genai/utils/trace_utils.py:557), 它做三件事:

  • 异步函数(async def)自动包成同步调用;
  • NoOpTracerPatcher 试跑一次样例输入,数它有没有产生 span;一个 span 都没有就自动套 mlflow.trace
  • 最后返回 lambda request: predict_fn(**request)——把 inputs 字典展开成关键字参数

最后一条解释了文档里那句要求:inputs 必须是字典,键名要和 predict_fn 的参数名对得上。

3.5 拿端点当 predict_fn:to_predict_fn

to_predict_fn(endpoint_uri)base.py:529)把一个 Databricks 端点变成可评估的函数, 按 URI schema 分流(base.py:611-620):apps:/<name>_create_app_predict_fnendpoints:/<name>_create_endpoint_predict_fn

_create_endpoint_predict_fnbase.py:623)里有一段值得学的 trace 缝合逻辑:

  1. 请求里注入 {"databricks_options": {"return_trace": True}},让服务端把 trace 一起返回;
  2. 如果返回的 trace 已经在当前 experiment 里(服务端双写模式),直接复用,不重复拷贝base.py:672-678);
  3. 否则 copy_trace_to_experiment 拷过来;
  4. 端点根本不返回 trace 时,用 start_span_no_context 手工造一棵只有根 span 的 trace, 时间戳取调用前后的真实毫秒(base.py:692-700)。

4. 数据实体:EvalItem 与 EvalResult

这节讲流水线上流动的两个对象。

4.1 EvalItem:既是输入也是暂存区

EvalItementities.py:78)是个 dataclass,字段分两类:

字段谁写的说明
request_id / inputs / expectations / tags / source数据集输入侧
outputs数据集 预测阶段_run_predict 会覆盖它
trace数据集 预测阶段三种来源,见下
error_message预测阶段predict_fn 抛异常时填这里,不中断整体评估

from_dataset_rowentities.py:125)负责逐列容错解析:inputs 允许是 JSON 串, trace 允许是 Trace 对象或 JSON 串,expectations 缺失补 {}。 没有 request_id 时用第一个非空字段的 sha256 当 ID(entities.py:154-164)。

预测阶段 _run_predictharness.py:806)给 trace 定了三条来路:

情况怎么拿 trace
predict_fn调用时挂一个 eval_request_id,跑完 mlflow.get_trace(eval_request_id) 取回
数据里已有 trace跨 experiment 或后端不支持 run 关联时克隆一份,否则直接 link_traces_to_run
只有静态 inputs/outputscreate_minimal_trace 现造一棵只有根 span 的 trace,好让分数有地方挂

「一行数据 = 一个可变对象」是这套设计的核心简化:预测线程往 EvalItem 里写 outputs/trace, 打分线程直接读同一个对象。因为每个 item 只被一个预测任务和一个打分任务串行触碰,不需要加锁。

4.2 EvalResult:一行的打分结果

EvalResultentities.py:219)= eval_item + assessments(Feedback 列表)+ scorer_stats(每个 scorer 调了几次、失败几次)。

get_assessments_dictentities.py:239)把每条 Feedback 摊平成四列: <name>/value<name>/rationale<name>/error_message<name>/error_code。 这个「斜杠列名」约定贯穿全链路——最终的 result_df 和断言判定都靠它。

4.3 EvaluationResult:返回给用户的东西

EvaluationResultentities.py:255)有 run_id / metrics / result_df,外加一对测试友好的属性:

  • passedentities.py:281):所有行、所有 scorer 都通过才 True;
  • reasonentities.py:295):失败原因的人话汇总。

判定规则在 _assertion_outcomeentities.py:20):优先用 scorer 自带的 pass_if 谓词, 否则只认 yes/no 字符串和 bool;其他值(比如 0.87 这种分数)一律判失败, 并在提示里教你写 @scorer(pass_if=lambda v: v >= 0.8)


5. 流水线调度(本章技术重心)

这节讲 harness.py 里那台并发机器。

5.1 为什么不能只用一个线程池

朴素做法是「一个池,每个任务里先 predict 再 score」。这有两个问题:

  1. 两段的并发度需求不一样。 predict 打的是你的 agent(可能一次一个 LLM 调用), score 打的是 judge 模型(N 个 scorer 就是 N 倍调用量)。用一个池,只能取一个折中值。
  2. 两段的限流对象不一样。 你的 agent 端点和 judge 端点是两套配额,得分开管。

所以 _run_pipelineharness.py:521)建了两个池、两个限流器。

5.2 生产者—消费者的四个角色

背压信号量(容量 = 2 × 打分线程数)
│ acquire
[提交线程]─────────────┴──────▶ [预测池 N₁ 线程]
_submit_all │ future
│ 把 future 放进 queue │
▼ ▼
┌─────────────────────────────────────────┐
│ [主循环] wait(pending, FIRST_COMPLETED) │
│ · 是预测 future → 提交打分任务 │
│ · 是打分 future → 记结果、release 信号量 │
└─────────────────────────────────────────┘
│ submit

[打分池 N₂ 线程]

▼ (每个任务内部还有第三层小池:一个 scorer 一根线程)

四个角色的分工:

角色代码职责
提交线程_PredictSubmitter._submit_allharness.py:318顺序遍历 items,取背压槽位后提交预测
预测池_PredictSubmitter._poolharness.py:274_run_predict
主循环_run_pipeline 的 while(harness.py:587-617单线程调度:分发、收结果、更新进度条
打分池_ScoreSubmitter._poolharness.py:422_run_score

主循环最巧的一点:pending 集合里混着两种 future,靠 predictor.owns(future)harness.py:370) 区分是预测完成还是打分完成,于是一个 wait(..., FIRST_COMPLETED) 就同时驱动了两段流水线(harness.py:601-617)。

还有一处刻意为之的「不优化」:即便一个单轮 scorer 都没有,item 也照样提交给打分池。 因为 _run_score 顺带做了写 expectations 和写 tags 的事,短路掉会丢东西——注释里挂了 issue 号 #23746harness.py:605-608)。

5.3 背压:一个信号量

如果提交线程无脑把 1 万条都提交进预测池,1 万棵 trace 会同时堆在内存里等打分。 _PredictSubmitter 用一个信号量卡住这件事:

self._in_flight = threading.Semaphore(backpressure_buffer(score_workers)) # harness.py:277

backpressure_bufferharness.py:151)就是 2 × 打分线程数。释放时机是关键—— 不是预测完就释放,而是该 item 打分完成后release_slot()harness.py:615:379)。 所以「已预测但还没打完分」的数量恒定不超过这个上限。

5.4 限流:令牌桶 + AIMD 自适应

RPSRateLimiterrate_limiter.py:70)是标准令牌桶:acquire() 消耗一个令牌, 不够就按 (1 - tokens) / rps 睡到够(rate_limiter.py:109-129)。桶容量等于 1 秒的量, 即允许一秒的突发。

自适应部分是 AIMD(加性增、乘性减,TCP 拥塞控制的老套路):

事件动作参数
撞到 429rps *= 0.5,但不低于 1.0_beta=0.5_min_rps=1.0rate_limiter.py:102-105
一次成功rps += 1.0 / rps,不超过初始值的 2 倍_alpha=1.0_max_rpsrate_limiter.py:153
连续降速5 秒冷却期内只降一次_throttle_cooldown=5.0rate_limiter.py:136-138

降速要冷却,是因为并发几十个线程可能同时撞 429——没有冷却的话一瞬间会被连着砍很多次。 增速用 1.0 / rps 而不是固定值,则让速率越高涨得越慢。

call_with_retryrate_limiter.py:168)把限流和重试缝在一起:每次尝试前 acquire(), 非限流错误直接抛出(不浪费重试次数),限流错误则 report_throttle() + 指数退避(上限 60 秒)。

配套的 eval_retry_contextrate_limiter.py:16)是个反直觉但必要的设计: 它关掉底层 HTTP 层和 litellm 的 429 重试,好让错误冒泡到这一层—— 否则底层默默重试完,上层的 AIMD 永远看不到拥塞信号。

5.5 池子多大:从限流速率反推线程数

_pool_sizeharness.py:136)的公式一句话:线程数 = 峰值 rps × 平均延迟, 再夹到 [10, 500],其中平均 LLM 延迟写死为 2 秒。

默认配置下的实际数字(假设 3 个 scorer):

来源
predict 限流"auto" → 初始 10 rps,adaptiveMLFLOW_GENAI_EVAL_PREDICT_RATE_LIMITenvironment_variables.py:841
scorer 限流未设时 = predict_rps × scorer 数 = 30 rps_get_scorer_rate_configharness.py:157
预测池线程10 × 2(AIMD 上限) × 2(延迟) = 40_get_pool_sizesharness.py:173
打分池线程30 × 2 × 2 = 120同上
背压容量2 × 120 = 240backpressure_buffer

MLFLOW_GENAI_EVAL_MAX_WORKERSenvironment_variables.py:825)一旦显式设置, 就同时覆盖两个池为同一个值(harness.py:183-185)——这是给「我就想固定并发」的用户的逃生舱。

5.6 第三层并发:每个 item 内部的 scorer 池

_compute_eval_scoresharness.py:916)在打分任务内部又开了一个线程池, 让这一行的多个 scorer 并行跑:

max_scorer_workers = min(len(scorers), MLFLOW_GENAI_EVAL_MAX_SCORER_WORKERS.get()) # harness.py:946

默认上限 10(environment_variables.py:833)。所以理论峰值并发 = 打分池线程数 × 每行 scorer 数, 真正的闸门是上一层的令牌桶而不是线程数。

5.7 进度反馈:进度条 + 心跳

  • 进度条tqdm 是可选依赖,没装就是 None,全链路用 if progress_bar 保护(harness.py:688-697)。 总任务数 = item 数 + session 数(harness.py:686)。结束时把「predict 占多少 %、scorer 占多少 %」 拼进 bar_format(harness.py:723-728),这是排查「到底谁慢」最直接的线索。
  • 心跳_Heartbeatharness.py:192)每 15 秒打一条 debug 日志, 内容是已预测/已打分/两边挂起数/两边当前 rps。它同时被用作 wait() 的 timeout, 保证主循环即使没有任务完成也会周期性醒来打日志(harness.py:596-600)。

5.8 会话级 scorer 为什么必须排在后面

多轮 scorer 在单轮流水线全部结束后才跑(harness.py:625)。 原因写在注释里:用法 ③ 下,trace 是单轮阶段才创建出来的,而 session 分组要读 trace 元数据里的 session id,两个阶段没法重叠(harness.py:621-624)。

分组逻辑 group_traces_by_sessionsession_utils.py:52)先看 trace 元数据的 mlflow.trace.session,再退回数据集记录的 source_data["session_id"]evaluate_session_level_scorerssession_utils.py:98)把整个 session 的 trace 列表传给 scorer, 结果挂在时间上最早的那条 trace 上(get_first_trace_in_sessionsession_utils.py:85)。


6. 打分与写回

这节讲一个 item 进了打分池之后发生什么。

6.1 一次打分的五步

_run_scoreharness.py:866)的顺序:

① 把 run_id 写进线程上下文 ← 打分线程不在开 run 的那个线程里
② _compute_eval_scores ← 并行跑所有 scorer,产出 Feedback 列表
③ 追加 expectations ← _get_new_expectations,只加 trace 上还没有的
④ 写 tags 到 trace ← 跳过 IMMUTABLE_TAGS
⑤ _log_assessments ← 把 Feedback 挂回 trace

第 ① 步不是可有可无的:MLflow 的 active run 是线程局部的,工作线程拿不到主线程开的 run, 所以 context.get_context().set_mlflow_run_id(run_id) 手工传递(harness.py:873-875, 配合 context.py:91set_mlflow_run_id)。

6.2 参数注入与「错误也是一种分数」

_invoke_scorerharness.py:907)永远按四个关键字调用:

return scorer_func(inputs=..., outputs=..., expectations=..., trace=...)

真正的「按签名取用」发生在 Scorer.run 内部(见第 7 节)。

_compute_eval_scores 里的 run_scorerharness.py:928)有个关键设计: scorer 抛异常不会让整行失败,而是被转成一条带 error 的 Feedbackharness.py:942-953), error_code 固定为 "SCORER_ERROR",还带完整 stack trace。

这意味着:失败在结果里是可见的一等公民,不是日志里的一行噪音。 ScorerStatentities.py:51)同时记录调用数和失败数,收尾时 _log_scorer_failure_summaryharness.py:1098)打一条 'correctness': 2/50 failed 这样的汇总。

打开 MLFLOW_GENAI_EVAL_ENABLE_SCORER_TRACING(默认 false,environment_variables.py:907)后, scorer 自身的执行也会被 trace,span 类型是 EVALUATOR,并把 scorer 的 trace id 写进 Feedback 元数据 (harness.py:955-967)——「给评审员本身也录像」,调试 judge 时很有用。

6.3 写回:Feedback 挂到 trace 上

_log_assessmentsharness.py:1015)对每条 assessment 做三件事再落库:

做什么为什么
trace_id预测阶段可能新建/克隆了 trace,ID 已经变了
元数据里塞 SOURCE_RUN_ID同一条 trace 可能被多次评估,靠它区分是哪个 run 的分
span_id 指向根 spanDatabricks 评估 UI 需要它才能展示

最后调 mlflow.log_assessment这一步是「评估」与「追踪」两个子系统的接缝: 分数不是存在某张评估表里,而是存成 trace 上的断言(见 01 的数据模型)。

6.4 收尾阶段的四个动作

harness.run 在流水线结束后按顺序做(harness.py:731-803):

  1. 合并多轮结果:会话级 assessment 已在工作线程写回 trace,这里只是补进 eval_results 好进 DataFrame(harness.py:733-738)。
  2. 关联与打标batch_link_traces_to_run 把 trace 挂到 run; _tag_mlflow_test_tracesharness.py:633)在 @mlflow.test 里跑时,给每条 trace 打上测试名和用例 ID——回归测试 UI 靠它分组。
  3. 刷新 trace_refresh_eval_result_tracesharness.py:1042)用线程池并行重新拉一遍 trace, 让内存里的对象包含刚写进去的全部断言(单轮 + 多轮一次性刷)。
  4. 聚合 + 清理:算指标、mlflow.log_metrics、搜出 run 下所有 trace、 用 clean_up_extra_traces 删掉评估过程中产生的噪音 trace(比如 judge 自己的调用), 最后 construct_eval_result_dftrace_utils.py:942)拼出结果表。

7. Scorer 抽象

这节讲「评审员」这个概念在代码里长什么样。

7.1 三种写法,一个基类

Scorerscorers/base.py:290)是个 pydantic BaseModel,只有 name / aggregations / description 三个公开字段。 子类只需实现 __call__。三条产出 Scorer 的路径:

写法长什么样kind
装饰器@scorer def my_check(outputs): ...ScorerKind.DECORATOR
内置类Correctness()Safety()ScorerKind.BUILTIN
自然语言 judgemake_judge(...)(见 05INSTRUCTIONS / GUIDELINES / MEMORY_AUGMENTED

ScorerKindscorers/base.py:54)这个枚举不是装饰品——注册、复制、序列化都按它分流。

7.2 参数注入:run__call__ 的分工

用户写的 scorer 通常只要一两个参数(def my_check(outputs)),但引擎总是传四个。 桥梁是 Scorer.runscorers/base.py:717):

sig = inspect.signature(self.__call__)
filtered = {k: v for k, v in merged.items() if k in sig.parameters}
result = self(**filtered)

mergedinputs / outputs / expectations / trace / session 五件套, 按 __call__ 的签名筛掉不要的。可为什么 CustomScorer.__call__ 只是 def __call__(self, *args, **kwargs), 签名里根本没有 outputs?答案在装饰器结尾——它手工改写了 __signature__scorers/base.py:1514-1518):

new_params = [inspect.Parameter("self", ...)] + list(signature.parameters.values())
CustomScorer.__call__.__signature__ = signature.replace(parameters=new_params)

于是 inspect.signature 看到的是原函数的参数表,筛选就对了。这是整个抽象最巧的一处: 用户函数的签名被当成一份「我要哪些字段」的声明。

run 之后还做两件收尾:校验返回类型(只允许 int/float/bool/str/Feedback/list[Feedback]/None), 以及把匿名 Feedback 的名字改成 scorer 名——注释解释了原因:失败路径下只知道 scorer 名, 成功路径若保留默认名 "feedback",同一个 scorer 会在结果里出现两种列名(scorers/base.py:750-767)。

7.3 scorer 装饰器做了什么

scorerscorers/base.py:1290)在函数上动态造一个 CustomScorer 类(scorers/base.py:1486):

  • 看签名里有没有 session,有就是会话级 scorer,并禁止同时出现 inputs/outputs/tracescorers/base.py:1473-1484);
  • _original_func 存原函数(序列化要用)、_pass_if 存断言谓词;
  • __call__ 直接转调原函数;
  • 名字默认取函数名(scorers/base.py:1521)。

顺带一提基类的 __init_subclass__scorers/base.py:305):每个子类的 __call__ 会被自动包一层遥测, 用 ContextVar 判断是否嵌套调用,只上报最外层——并且显式拒绝 async 的 __call__scorers/base.py:264-268)。

7.4 序列化:把源码存下来

Scorer 要能存进服务端(注册后用于线上打分),但它可能是个闭包函数。MLflow 的答案是存源码文本

SerializedScorerscorers/base.py:155)是扁平的 dataclass,按 scorer 类型分成五组互斥字段, __post_init__ 强制「有且只有一组」(scorers/base.py:194-227):

字段组存什么对应 kind
builtin_scorer_class + builtin_scorer_pydantic_data类名 + 构造参数BUILTIN
call_source + call_signature + original_func_name函数体源码 + 签名串 + 函数名DECORATOR
instructions_judge_pydantic_data指令 + 模型INSTRUCTIONS
memory_augmented_judge_data对齐后的 judge 数据MEMORY_AUGMENTED
third_party_scorer_data模块 + 类 + metric 名 + kwargsTHIRD_PARTY

写: model_dumpscorers/base.py:382)对装饰器 scorer 调 _extract_source_code_info, 后者用 AST 抽出函数体(extract_function_bodyscorers/scorer_utils.py:99)——会跳过 docstring、 做 dedent,只留纯函数体。结果缓存在 _cached_dump,避免动态函数被二次序列化时出错。

读: model_validatescorers/base.py:455)按哪组字段非空来分流。 装饰器 scorer 走 _reconstruct_decorator_scorerscorers/base.py:667): recreate_functionscorer_utils.py:115)用 def 名字(签名):\n<缩进后的函数体> 拼出源码, 在一个预置了 Feedback/Trace/CategoricalRating 等符号的命名空间里 exec,再重新套一遍 @scorer

安全边界很明确:exec 有代码执行风险,所以非 Databricks 的 tracking URI 直接拒绝反序列化, 并把源码原样打印出来让你手工粘贴(scorers/base.py:675-686)。 第三方 scorer 同理,只允许从四个白名单模块导入(THIRD_PARTY_SCORER_ALLOWED_MODULESscorer_utils.py:46)。

另有一处向前兼容的细节:SerializedScorer.from_dictscorers/base.py:229)会丢弃未知字段并打一条 「你的 mlflow 太旧,升级到 >= X」的错误日志,而不是抛 TypeError

7.5 注册与生命周期:register / start / stop

这三个方法把 scorer 接到线上定时打分上(细节见 06):

方法干什么位置
register存进 scorer store,返回带后端标记的新副本scorers/base.py:878
start设采样率(必须 > 0)+ 可选 filter,开始自动评估scorers/base.py:944
update改采样率 / filterscorers/base.py:1030
stop本质就是 update(sample_rate=0.0)scorers/base.py:1118

三个设计点:

  • 全部返回新实例,不原地修改(_create_copyscorers/base.py:1165)——所以 Scorer 用起来像值对象。
  • 状态是推导出来的,不是存的:ScorerStatusscorers/base.py:76)由 _registered_backendsample_rate 反推, 文档说明白了原因——后端压根没有「启动/停止」这个概念(scorers/base.py:79-81)。
  • 准入检查集中在 _check_can_be_registeredscorers/base.py:1219):kind 白名单、 装饰器 scorer 只许在 Databricks 注册、第三方 scorer 反过来不许在 Databricks 注册、 judge 模型必须是 databricks:/ 开头。

8. 内置 scorer 家族

这节给个分类地图,具体 judge prompt 与调用见 05

8.1 BuiltInScorer 的共同点

BuiltInScorer(Judge)builtin_scorers.py:314)在 Scorer 之上加了三样:

  • required_columns:声明「我需要哪些列」,validate_columns 在准备阶段被调用做预检;
  • instructions 抽象属性:这个 scorer 到底在评什么(也是 judge 的提示词来源);
  • 自己的 model_dump / model_validate:只存类名 + pydantic 字段,反序列化时按类名从模块里取类再实例化(builtin_scorers.py:350-382)。

8.2 按「要不要 LLM」分三类

类别代表要什么列说明
确定性规则RegexMatch:3229outputs正则 search/fullmatch,返回 yes/no
PIIDetection:3338outputs正则扫邮箱/电话/SSN/卡号/IP,发现 PII 返回 no
ResponseLength:3447outputs长度上下界
LLM judgeCorrectness:1699inputs+outputs+expected_response/expected_facts答案是否被期望事实支持
Safety:1593inputs+outputs有害内容
Guidelines:1156inputs+outputs是否遵守自然语言规则(kind 是 GUIDELINES
RelevanceToQuery:1475)/ Fluency:1895)/ Completeness:2975)/ Summarization:3084inputs+outputs单点质量维度
要读 trace 结构RetrievalRelevance:393inputs+trace逐个 chunk 判相关性
RetrievalGroundedness:687inputs+trace回答有没有脱离检索到的上下文
RetrievalSufficiency:551inputs+trace检索到的内容够不够回答
ToolCallCorrectness:875)/ ToolCallEfficiency:794trace工具选得对不对、调得冗不冗余

RegexMatch 是理解「非 LLM scorer 也是 BuiltInScorer」的最好样本: 它照样实现 instructions(返回一句人话描述)、照样返回 Feedback(value=CategoricalRating.YES/NO)builtin_scorers.py:3333-3358),只是中间没有模型调用。

RetrievalRelevance 则是「一个 scorer 产出多条 Feedback」的样本:每个 chunk 一条, 再加一条名为 <scorer 名>/precision 的 span 级平均分(builtin_scorers.py:542-547)。

8.3 字段抽取:resolve_scorer_fields

内置 scorer 的 __call__ 允许「只给 trace,别的自己想办法」。这件事由 resolve_scorer_fieldsbuiltin_scorers.py:233)统一处理,逻辑分三级降级:

没有 trace ──▶ 原样返回调用方给的 inputs/outputs/expectations

有 trace
├─▶ ① 从 trace 根 span 直接取 inputs / outputs
│ (resolve_inputs_from_trace / resolve_outputs_from_trace)
├─▶ ② 需要 expectations 时,从 trace 的 Expectation 断言里取
└─▶ ③ 还缺 judge 声明需要的字段 → 拿 LLM 读整棵 trace 抽出来
(_construct_field_extraction_config + 结构化输出)

第 ③ 级是个有意思的兜底:现造一个只含缺失字段的 pydantic schema, 让模型带着 trace 去「找出用户的原始问题 / 系统的最终回答」(builtin_scorers.py:269-290)。 抽取失败只 warning,不炸——随后由 _validate_required_fieldsbuiltin_scorers.py:199)统一报缺什么。

判断「需要哪些字段」靠 judge.get_input_fields()——同一套 JudgeField 声明既驱动 prompt 构造,也驱动字段抽取

8.4 会话级分支

SessionLevelScorerbuiltin_scorers.py:2185)是另一条继承线:

  • is_session_level_scorer 恒为 True,get_input_fields() 只声明一个 session
  • __call__ 只接受 sessionexpectations,传别的直接 TypeErrorbuiltin_scorers.py:2226-2253);
  • 内部 judge 惰性创建并缓存(_get_judgebuiltin_scorers.py:2208)。

公开的会话级 scorer 继承 BuiltInSessionLevelScorerbuiltin_scorers.py:2256,同时吃到内置序列化): UserFrustration:2260,返回 none/resolved/unresolved 三态)、ConversationCompleteness:2339)、 ConversationalSafety:2418)、ConversationalToolCallEfficiency:2499)、 ConversationalRoleAdherence:2577)、ConversationalGuidelines:2654)、KnowledgeRetention:2789)。

get_all_scorers()builtin_scorers.py:3630)不维护硬编码列表,而是递归遍历 BuiltInScorer.__subclasses__() 挑出本模块里的非抽象类(builtin_scorers.py:3602),再逐个尝试无参构造,构造不出来的(比如 RegexMatch 必须给 pattern)跳过。


9. 结果聚合:逐行分数怎么变成 run 指标

compute_aggregated_metricsaggregation.py:26)分三步:

  1. 收集:遍历所有 EvalResult 的 Feedback,按 assessment 名字归到一个列表里。
  2. 转 float_cast_assessment_value_to_floataggregation.py:77)—— 数值/bool 直接转;字符串里只有合法的 CategoricalRating 会被转成 yes → 1.0、其他 → 0.0; 转不出来的返回 None,不参与聚合
  3. 聚合:按 scorer 声明的 aggregations 算,默认只算 mean。 可选项是 min/max/mean/median/variance/p90aggregation.py:16),也允许传自定义函数 (算错了只 log error 并跳过,aggregation.py:109-113)。

指标名的拼法是 {assessment 名}/{聚合名},比如 correctness/mean

一个容易踩的细节:查 scorer 聚合配置时用的 key 是 name.split("/", 1)[-1]aggregation.py:64), 即斜杠后面那截。所以 RetrievalRelevance 产出的 retrieval_relevance/precision 会用 "precision" 去查配置、 查不到就退回 ["mean"],最终指标名是 retrieval_relevance/precision/mean

算完直接 mlflow.log_metrics(aggregated_metrics)harness.py:766),指标就进了 run。


10. 巧妙之处(可以带走的技术)

  1. 两种 future 混在一个 pending 集合里调度。 一个 wait(FIRST_COMPLETED) 同时驱动两段流水线, 靠 owns() 分辨归属,省掉了一整套跨池协调代码(harness.py:601-617:370)。

  2. 背压槽位在「打分完成」时才释放,不是「预测完成」。 直接把「内存里堆积的 trace 数」 变成一个可证明的上限(harness.py:277:602)。

  3. 底层重试被主动关掉。 eval_retry_contextrate_limiter.py:16)关掉 HTTP 层和 litellm 的 429 重试, 让拥塞信号能到达 AIMD。「为了让上层看得见,先把下层的自动补救关掉」这个思路很值得记。

  4. 线程数从限流速率推导,而不是让用户猜。 线程数 = 峰值 rps × 平均延迟harness.py:136), 限流器负责排队,抢不到令牌的线程阻塞在 acquire()。用户只需要说「我的配额是多少 rps」。

  5. __signature__ 改写把用户函数签名变成依赖声明。 @scorer def f(outputs) 里的参数表, 就是「我需要 outputs」的声明(scorers/base.py:1514-1518:627)。

  6. 失败被建模成数据,而不是异常。 scorer 抛错 → 带 SCORER_ERROR 的 Feedback, 照样写回 trace、照样进结果表(harness.py:942-953)。评估流程从不因为单个评审员崩溃而中断。

  7. 可执行代码的反序列化有明确信任边界。 装饰器 scorer 存源码、读时 exec, 所以非 Databricks 环境直接拒绝并把源码打给你看(scorers/base.py:667-686); 第三方 scorer 走模块白名单(scorer_utils.py:46)。

  8. 主线程外的 run 上下文靠显式传递。 每个工作线程开头 ctx.set_mlflow_run_id(run_id)harness.py:814-816:847-849),绕开「active run 是线程局部」这个坑。

  9. 自动埋点在池外配置。 注释直说是为了避开 race(base.py:344); 同理,litellmdatabricks.sdk 在主线程预先 import,避免多线程首次 import 撞上模块导入锁死锁 (harness.py:25-36)。


11. 边界与局限

  • evaluate() 明确不是线程安全的,docstring 里写了警告(base.py:292-295)。 根因之一是 NoOpTracerPatcher 会全局改 NoOpTracer.start_spantrace_utils.py:586)。
  • predict_fn 必须一次调用产出一棵 trace,否则引擎按 eval_request_id 取不到 trace(harness.py:821-837)。
  • 会话级 scorer 不能和 predict_fn 直接组合,必须走 ConversationSimulator 或已有 trace(session_utils.py:190)。
  • 缺列不阻断评估:内置 scorer 要的列不全时只打一条 info 日志,对应分数留空(validation.py:172-186)。
  • 平均 LLM 延迟写死 2 秒harness.py:144)。如果你的 agent 一次跑 30 秒,默认线程数会偏小, 得手工设 MLFLOW_GENAI_EVAL_MAX_WORKERS
  • 装饰器 scorer 在 OSS 后端不能注册/加载——不是功能没做,是安全上的刻意选择(scorers/base.py:1237-1244)。
  • 非 yes/no、非 bool 的分数在 result.passed 里一律算失败,除非显式声明 pass_ifentities.py:42-47)。

12. 代码地图(导航索引)

主题文件路径符号名
公开入口mlflow/genai/evaluation/base.pyevaluate
准备阶段全流程mlflow/genai/evaluation/base.py_run_harness
数据集记录到 runmlflow/genai/evaluation/base.py_log_dataset_input
端点/App 转 predict_fnmlflow/genai/evaluation/base.pyto_predict_fn_create_endpoint_predict_fn
输入归一化管道mlflow/genai/evaluation/utils.py_convert_to_eval_set_extract_request_response_from_trace
返回值 → Feedbackmlflow/genai/evaluation/utils.pystandardize_scorer_value
行实体mlflow/genai/evaluation/entities.pyEvalItem.from_dataset_rowEvalResultEvaluationResult
断言判定规则mlflow/genai/evaluation/entities.py_assertion_outcomeScorerStat
评估主流程mlflow/genai/evaluation/harness.pyrun
流水线主循环mlflow/genai/evaluation/harness.py_run_pipeline
预测池 / 打分池mlflow/genai/evaluation/harness.py_PredictSubmitter_ScoreSubmitter
背压与池大小mlflow/genai/evaluation/harness.pybackpressure_buffer_pool_size_get_pool_sizes
心跳日志mlflow/genai/evaluation/harness.py_Heartbeat
单行预测 / 单行打分mlflow/genai/evaluation/harness.py_run_predict_run_score
多 scorer 并行mlflow/genai/evaluation/harness.py_compute_eval_scores_invoke_scorer
分数写回 tracemlflow/genai/evaluation/harness.py_log_assessments
测试标记与刷新mlflow/genai/evaluation/harness.py_tag_mlflow_test_traces_refresh_eval_result_traces
令牌桶与 AIMDmlflow/genai/evaluation/rate_limiter.pyRPSRateLimitercall_with_retryeval_retry_context
会话分组与多轮打分mlflow/genai/evaluation/session_utils.pyclassify_scorersgroup_traces_by_sessionevaluate_session_level_scorers
线程内 run 上下文mlflow/genai/evaluation/context.pyeval_contextRealContext.set_mlflow_run_id
Scorer 基类与注入mlflow/genai/scorers/base.pyScorerScorer.run
装饰器mlflow/genai/scorers/base.pyscorerCustomScorer
序列化 schemamlflow/genai/scorers/base.pySerializedScorerScorer.model_dumpScorer.model_validate
源码级重建mlflow/genai/scorers/base.py_reconstruct_decorator_scorer
注册生命周期mlflow/genai/scorers/base.pyregisterstartstop_check_can_be_registered
源码抽取与 execmlflow/genai/scorers/scorer_utils.pyextract_function_bodyrecreate_function
内置 scorer 基类mlflow/genai/scorers/builtin_scorers.pyBuiltInScorerget_all_scorers
字段抽取降级mlflow/genai/scorers/builtin_scorers.pyresolve_scorer_fields_validate_required_fields
会话级 scorermlflow/genai/scorers/builtin_scorers.pySessionLevelScorerUserFrustration
指标聚合mlflow/genai/scorers/aggregation.pycompute_aggregated_metrics_cast_assessment_value_to_float
列预检mlflow/genai/scorers/validation.pyvalidate_scorersvalid_data_for_builtin_scorers
predict_fn 包装 / 结果表mlflow/genai/utils/trace_utils.pyconvert_predict_fnconstruct_eval_result_dfcreate_minimal_trace
并发相关环境变量mlflow/environment_variables.pyMLFLOW_GENAI_EVAL_MAX_WORKERSMLFLOW_GENAI_EVAL_PREDICT_RATE_LIMIT