跳到主要内容

数据截至 (上游 commit 0261ea4f33d4)

评估层:dataset × task × metric 怎么跑成一次 experiment

30 秒导读: 这一层解决的是"我改了 prompt/模型,到底有没有变好"。你给它一份数据集(dataset)、一个待测函数(task)、一组打分器(metric),它把每条数据喂给 task、把 task 的一次执行包成一条 trace、拿 metric 打分、再把分数写回那条 trace,最后所有这些被归到一个 experiment 下面供对比。本章只讲编排骨架——具体某个指标怎么算分、LLM 裁判怎么提问,看 打分器那一章


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

一句话定义: 离线评估引擎 = 一个"批量跑 + 批量打分 + 结果归档"的编排器。

它解决的问题。 你写了个客服机器人。改了系统提示词之后,你想知道:100 个历史问题里,答对的比例是升了还是降了?手工点 100 次不现实,写个 for 循环又缺三样东西——并发、失败不中断、结果能存下来跟上一版比。这三样就是这一层补的。

三个名词先说清(后面反复出现,一词一义):

名词在 Opik 里指什么
dataset item一条测试数据。一个 dict,字段随便你定({"question": ..., "expected": ...}
task你的待测函数。签名固定:吃一个 dict,吐一个 dict
metric打分器。有个 score(...) 方法,参数名决定它要哪些字段
experiment一次评估跑批的归档单位。它把"哪份数据 × 哪次执行 × 得了几分"三者绑在一起

用起来什么样。 最小可运行的一段(真实 API,参数见 sdks/python/src/opik/evaluation/evaluator.py:137 evaluate):

import opik
from opik.evaluation.metrics import Equals

def my_task(item: dict) -> dict: # 待测函数:吃 dict 吐 dict
answer = call_my_llm(item["question"])
return {"output": answer}

result = opik.evaluate(
dataset=opik.Opik().get_dataset("qa-v1"),
task=my_task,
scoring_metrics=[Equals()], # Equals.score(output, reference)
scoring_key_mapping={"reference": "expected"}, # 数据集字段名 → 指标参数名
task_threads=16,
trial_count=3, # 同一条数据跑 3 遍,看稳定性
)
print(result.aggregate_evaluation_scores().aggregated_scores)

一句话直觉: 把它当成"单元测试框架 + 打分"。dataset item 是测试用例,task 是被测函数,metric 是断言——只不过断言的结果不是 pass/fail 而是一个 0~1 的分,且允许同一个用例跑多遍取统计。


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

2.1 一次 evaluate() 的主干

怎么读这张图:从上到下是时间顺序;虚线框是"每条 dataset item 都要走一遍"的循环体。

用户调用 evaluate(dataset, task, metrics)


① 建 experiment(拿到 experiment_id,把 resume 状态写进 experiment_config)


② 解析数据(流式 or 抽样后物化)→ Iterator[DatasetItem]


③ EvaluationEngine.run_and_score(...)

├─ 指标分流:只看输入输出的 / 要看整棵 span 树的


┌ ─ ─ ─ ─ ─ ─ ─ ─ ─ 每条 item × runs_per_item 次 ─ ─ ─ ─ ─ ─ ─ ─┐
│ ④ 开一条 trace,把 task 调用包进去 │
│ ⑤ task(item_content) → task_output │
│ ⑥ 排干 streamer → 建 agentic 上下文 → metric 打分 │
│ ⑦ 分数写回这条 trace(feedback score) │
│ ⑧ 收尾:落 trace + 插一条 experiment item(串起三者) │
└ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ┘


⑨ 汇总:终端表格 + experiment 级别的聚合分 + 返回 EvaluationResult

2.2 部件一句话职责

部件干什么在哪个文件
公开入口参数校验、建 experiment、解析数据、调引擎、打印汇总sdks/python/src/opik/evaluation/evaluator.py
EvaluationEngine编排:跑 task、包 trace、调打分、并发sdks/python/src/opik/evaluation/engine/engine.py:67
StreamingExecutor线程池 + 进度条,边提交边执行sdks/python/src/opik/evaluation/engine/evaluation_tasks_executor.py:18
MetricsEvaluator指标分流、字段名映射、逐个打分并容错sdks/python/src/opik/evaluation/engine/metrics_evaluator.py:357
evaluate_llm_task_context上下文管理器:开/关 trace,收尾插 experiment itemsdks/python/src/opik/evaluation/engine/helpers.py:28
Dataset / DatasetVersion数据源,含内容哈希去重与版本信息sdks/python/src/opik/api_objects/dataset/dataset.py
Experiment归档单位,负责插 experiment item、写 experiment 级分数sdks/python/src/opik/api_objects/experiment/experiment.py:52
resume/断点续跑:状态、checkpoint、待办计算、结果合并sdks/python/src/opik/evaluation/resume/

2.3 三个数据对象的关系

评估层从头到尾在搬运三样东西,理解它们的关系就理解了一大半:

DatasetItem ──(get_content)──► dict ──► task ──► task_output(dict)
│ │
│ ▼
│ TestCase{trace_id, dataset_item_id,
│ task_output, dataset_item_content}
│ │
│ ▼ metric.score(**mapped)
│ TestResult{test_case, score_results, trial_id}
│ │
▼ ▼
ExperimentItem ◄──── trace_id ────────────── feedback scores 挂到 trace 上
(dataset_item_id + trace_id)

一句话:TestCase 是"输入 + 输出 + 那条 trace 的 id",TestResultTestCase 加上分数evaluation/test_case.py:9evaluation/test_result.py:10)。


3. 公开入口全家福

evaluator.py 对外导出 7 个入口,加一个给测试套件/优化器共用的 __internal_api__run_test_suite__(导出清单见 sdks/python/src/opik/evaluation/__init__.py;同处导出的 evaluate_threadsevaluation/threads/ 独立模块,评的是会话线程而非 dataset,不在本章范围)。7 个入口的差别就三点:task 谁给、数据从哪来、结果挂到哪个 experiment

入口task 由谁提供数据从哪来experiment典型场景
evaluate用户传Dataset / DatasetVersion新建通用离线评估
evaluate_prompt引擎用 messages + model 现造Dataset / DatasetVersion新建只测提示词,不想写函数
evaluate_optimization_trial用户传Dataset / DatasetVersion新建,type="trial" 且挂 optimization_id优化器的一次试验
evaluate_experiment不跑 task从已有 experiment 的 item 反查复用已有换一批指标重新打分
evaluate_resume用户重新传按 experiment 里存的配置 + 本地 checkpoint 重建复用已有上次跑一半崩了
evaluate_on_dict_items用户传内存里的 List[dict]不建优化器内循环,要快、不想污染 experiment 列表
run_tests用户传TestSuite / TestSuiteVersion新建,evaluation_method="evaluation_suite"断言式测试,出 pass/fail 报告
__internal_api__run_test_suite__用户传(外面包了结果校验)同上同上,可挂 optimization_idrun_tests 与优化套件共用的实现体

几条容易踩的差别,逐条说:

evaluate_experiment 是唯一不跑 task 的。 它调的是 EvaluationEngine.score_test_casesengine/engine.py:637),TestCaserest_operations.get_experiment_test_casesevaluation/rest_operations.py:43)从后端已有的 experiment item 里重建。注意它会跳过没有 evaluation_task_output 的 itemevaluation/rest_operations.py:63-73)——那代表上次那一趟没跑完,没有输出可重新打分。

evaluate_on_dict_items 最轻。 它把裸 dict 现场包成 DatasetItem(id=f"temp_item_{i}", ...)experiment_=None 传进引擎(evaluation/evaluator.py:1770-1796)。没有 experiment 就不插 experiment item,只留 trace。默认 verbose=0,因为它是给优化器内循环用的。

run_tests 的指标不是用户传的,是从 dataset 版本上读的:dataset.get_evaluators(evaluator_model) 把后端存的 llm_judge 配置实例化成 LLMJudgeapi_objects/dataset/dataset.py:468);执行策略也来自 dataset.get_execution_policy()。它还会在跑之前临时打开本地 emulator,跑完在 finally 里恢复原状态(evaluation/evaluator.py:752-800)——理由见 §4.4。

evaluate_prompt 现造 task。 _build_prompt_evaluation_taskevaluation/evaluator.py:981)返回一个闭包:把 dataset item 的字段填进 chat 模板、调模型、返回 {"input": 渲染后的消息, "output": 模型回复}。填模板前会检查模型是否支持 vision/video,不支持就降级成文本占位并告警。


4. 核心机制

4.1 EvaluationEngine:为什么"无状态"是刻意的

要解决的小问题: 一个引擎实例要同时服务好几种流程(跑 task 的、只打分的、带 span 指标的),如果把 metric 列表、字段映射这些流程数据存成实例属性,多线程复用时就会互相污染。

思路: 实例上只放配置——client、project_name、workers、verbose、source、flush_timeout;所有流程数据(metric 列表、scoring_key_mapping、experiment、执行策略)一律走方法参数。类的 docstring 把这条写死了:

class EvaluationEngine:
"""
Stateless evaluation executor.

Only stores configuration (client, workers, verbosity).
All flow-specific data (metrics, key mappings) is passed as method parameters.
"""

—— sdks/python/src/opik/evaluation/engine/engine.py:67-73,构造函数只赋这 6 个字段(engine.py:75-93)。

两条公开路径,对应"跑不跑 task":

方法跑 task 吗入参核心用它的入口
run_and_scoreengine.py:595dataset_items 迭代器 + task + 执行策略evaluate / evaluate_prompt / run_tests / on_dict_items / resume
score_test_casesengine.py:637不跑已有 List[TestCase]evaluate_experiment

两条路都先做同一件事:split_into_regular_and_task_span_metrics(scoring_metrics) 把指标分成两类(§4.6)。区别是 score_test_cases 直接丢掉 task-span 那一类——没跑 task 就没有 span 树可看。

4.2 一次执行怎么被包进一条 trace

要解决的小问题: 用户的 task 是个普通函数,怎么让它这一次调用连同内部所有 LLM 调用,变成后台可查的一棵 span 树,并且和某条 dataset item 对上号?

思路: 先用 @opik.track 把 task 装饰成会自动出 span 的函数(如果用户没装饰过),再手工建一条 TraceData 塞进 context storage,让装饰器产生的 span 自动挂到这条 trace 下。追踪装饰器本身的原理见 追踪层那一章

原理演示:

# 示意,非源码
if not hasattr(task, "opik_tracked"): # 用户没标注过就替他标注
task = opik.track(name=task.__name__)(task)

trace_data = TraceData(input=item_content, name="evaluation_task")
with evaluate_llm_task_context(experiment, item.id, trace_data, client) as state:
output = task(item_content) # 内部所有 span 自动挂到 trace_data 下
update_current_trace(output=output)
result = score(...) # 打分也在 with 里面(关键)
state.evaluation_completed = True # 只有全程无异常才会走到这行
# __exit__ 里才真正把 trace 发出去,并插 experiment item

真实实现: _compute_test_result_for_llm_taskengine/engine.py:269)。自动装饰在 engine.py:279-281TraceData 构造在 engine.py:284-290,名字用常量 EVALUATION_TASK_NAME = "evaluation_task"engine.py:28)。

这里有个不显然的设计——"完成标志"。 EvaluationContextState.evaluation_completed 默认 False,引擎在打完分之后才置 True(engine.py:376)。上下文管理器的 finally 里:

if not state.evaluation_completed:
trace_data.output = None

—— engine/helpers.py:51-56

为什么要把 output 抹掉?因为**"persisted trace 有没有 output" 就是断点续跑判断这一趟跑完没有的唯一信号**。is_trial_fully_completed 直接读这个(resume/context.py:175-186)。任务抛异常、打分崩了、写分数失败——任何一处没走到那行,output 就是 None,下次 resume 会重跑这一趟。这是一个把"事务性"塞进日志字段的做法,不用额外状态表。

__exit__ 的收尾顺序是:写 error_info(若有)→ 按标志决定是否抹 output → init_end_time() → 落 trace → 若有 experiment 则插一条 ExperimentItemReferencesengine/helpers.py:41-70)。experiment item 就是这里把 dataset item 和 trace 绑在一起的,它只存两个 id 加一份执行策略(api_objects/experiment/experiment_item.py:11)。

4.3 重跑几遍、几遍算过:执行策略的合并规则

要解决的小问题: LLM 有随机性,跑一遍的分数不可信。想让同一条数据跑 N 遍;而且不同数据的"难度"不同,有的需要 5 遍过 4 次才算过,有的跑 1 遍就够。

两个参数:

参数含义谁在用
runs_per_item同一条 dataset item 跑几遍引擎,决定提交几个任务
pass_threshold至少几遍通过才算这条数据通过test suite 的结果构造器

合并规则:item 级覆盖 suite 级,逐字段生效。 get_item_execution_policyengine/engine.py:33-64):item 没自己的策略就整份用默认;有的话逐字段挑——item 的字段非 None 就用 item 的,否则回落默认。注意回落用的是 default_policy.get("runs_per_item", 1),所以两级都缺时兜底是 1。

evaluate() 这类入口构造默认策略时把两个值都设成 trial_countevaluation/evaluator.py:634-637)——即"跑 N 遍就得 N 遍全过"。test suite 则从 dataset 版本上读真实配置。

重跑怎么落地。 _compute_test_results_with_execution_policyengine.py:382)对每条 item:先解析策略、把解析结果写回 item.execution_policy(后面构造结果时还要读)、先声明组大小再提交,然后提交 item_runs 个任务,每个带同一个 group_id=item.id、不同的 trial_id=run_id

executor.set_group_size(item.id, item_runs) # 先声明,避免和回调抢
for run_id in range(item_runs):
executor.submit(functools.partial(..., trial_id=run_id), group_id=item.id)

—— engine.py:420-436。先声明组大小是为了避免竞态:回调在工作线程里跑,如果它比 set_group_size 先到,进度条就会拿不到组大小。

阈值在哪里判。 引擎不判——它只负责把 trial_id 不同的多个 TestResult 都产出来。真正判定在 build_suite_resultapi_objects/dataset/test_suite/suite_result_constructor.py:18):一趟里所有断言都过 → 这趟 passed;runs_passed >= pass_threshold → 这条 item passed;所有 item 都过 → suite passed(判定逻辑 suite_result_constructor.py:62-73)。有个细节:没有任何断言的一趟算通过not r.score_results or all(...))。

4.4 排干 streamer:为什么打分前必须等一下

要解决的小问题: agentic 类型的 LLM 裁判要读"刚才这次执行产生的 span 树"来判断(比如"agent 有没有调用搜索工具")。这些 span 是 @opik.track 在函数进出时提交的,但提交只是丢进 streamer 队列,真正应用到本地 emulator 是消费者线程异步做的。打分如果立刻开始,可能读到不完整的视图。

时序问题长这样:

task 执行中 ──► @opik.track 提交 span 消息 ──► [streamer 队列] ──► 消费者线程 ──► emulator
▲ │
│ ▼
打分开始 ────────────────────────────────────────────┘ 不等 = 可能读到旧状态

解法:一次有超时上限的排干。 _build_trace_tool_contextengine/engine.py:193)在建上下文前调:

drained = self._client.__internal_api__drain_to_processors__(timeout=self._flush_timeout)

—— engine.py:245-247self._flush_timeout 默认取常量 DEFAULT_STREAMER_DRAIN_TIMEOUT_SECONDS = 5.0engine.py:30,构造函数在 engine.py:91-93 兜底赋值)。

三个设计点,逐条:

  1. 是 drain 不是 flush。 __internal_api__drain_to_processors__ 只等本地处理器把队列吃完,跳过文件上传和 replay 的刷新(api_objects/opik_client.py:2064-2077)——打分只关心本地状态,不关心是否已送达后端。
  2. 超时是有意的、且超时不报错。 超时只打一条 debug 日志然后继续,拿当前状态尽力而为(engine.py:248-253)。一个卡住的消费者线程不该把整轮评估拖死。
  3. 排干之后还要摸一下 emulator.trace_trees _ = emulator.trace_treesengine.py:258)强制重建树,让本进程记录的 span 挂到 trace 上。

还有个顺序悖论要处理:打分发生在 evaluate_llm_task_contextwith 体内,而这条 trace 的 CreateTraceMessage 要等 __exit__ 才发。也就是说打分时 span 在 emulator 里,trace 不在。解法是两条构造路径:

场景trace_data 传了吗走哪条为什么
现跑现打分传了(engine.py:352build_trace_tool_context_from_trace_data用内存里的 TraceData 合成 TraceModel,span 仍从 emulator 拿
重新打分没传(engine.py:663build_trace_tool_contexttrace 是上一轮落的,emulator 里查得到

两条路径的取舍写在 suite_evaluators/agentic/context.py:265-291 的 docstring 里。查不到 trace 时返回 None,裁判退回一次性(非 agentic)路径。

另一半问题:别让评估自己产生的 span 污染被评估对象。 引擎给自己的打分环节加了 @opik.track(为了可观测),这些 span 的 input 里含整个 TestCase 信封——包括 LLMJudgeConfig 里的断言文本。如果裁判读到它们,就会看见"断言原文"出现在它要评判的 trace 里,分不清那是 agent 真干了这事、还是断言自己在回声。

解法是一个命名空间化的 sentinel 标签:

INTERNAL_SPAN_TAG = "__opik_eval_internal__"

—— suite_evaluators/agentic/context.py:177。引擎的两个打分方法都带上它:metrics_calculationengine.py:97-99)和 task_span_metrics_calculationengine.py:163-165)。过滤在 _filter_internal_spanscontext.py:185-229):先按标签选出种子集合,再顺着 parent 映射反复扫描到闭包,整棵子树一起丢——被标记 span 底下的都是评估插件(其它 scorer、模型包装层),一样不该被看见。

为什么用标签而不是按名字匹配?因为 span 名字是用户可设的(@opik.track(name=...)),匹配 name == "metrics_calculation" 会误杀恰好同名的用户 span。这段权衡原文写在 context.py:154-181 的注释里。

4.5 并发、进度与失败归类

StreamingExecutor 的三个特点engine/evaluation_tasks_executor.py:18):

  1. 边提交边执行。 不用先攒完再 mapsubmit() 立刻丢进线程池(evaluation_tasks_executor.py:138),所以数据是流式从后端来的也能马上开跑。
  2. workers == 1 不开线程池。 直接在调用线程里同步跑完,用一个手工 Future 填结果,等价于普通 for 循环(evaluation_tasks_executor.py:118-129)。调试时省掉一层线程。
  3. 进度按"条"不按"趟"。group_id 时,只有一组全完成才 update(1)evaluation_tasks_executor.py:186-200),所以 trial_count=3 时进度条显示的是 item 数而不是 3 倍的运行数。

进度条上还挂着实时均分。 每完成一个任务,回调在锁内累加 score_totals/score_counts,用 set_postfix 刷出各指标的跑动平均(evaluation_tasks_executor.py:161-183)。scoring_failed 的分不计入。test suite 关掉这个后缀(show_scores_in_progress_bar=Falseevaluation/evaluator.py:794),因为断言型分数的均值没有意义。

失败处理是两段式的:

任务抛异常 ──► _on_future_done: 记 WARNING,不入 results(进度也不推进)


get_results(): futures.wait 等全部结束


逐个 future.result() ──► 重抛第一个异常

—— 回调在 evaluation_tasks_executor.py:148-157,重抛在 evaluation_tasks_executor.py:232-233注意:一个 item 失败会让整轮 evaluate() 抛异常,不是静默跳过。

exception_analyzer.py 只做一件事:认出限流。 整个文件 14 行,is_llm_provider_rate_limit_error 判三种情况——openai.RateLimitErrorlitellm.exceptions.RateLimitError、或者任何带 status_code == 429 的异常(engine/exception_analyzer.py)。命中就额外打一条针对性日志,提示是被限流而不是代码坏了。两处在用:task 执行失败时(engine.py:313-316)和指标计算失败时(metrics_evaluator.py:347-350)。

4.6 指标分流与字段名映射

分流:两类指标,靠函数签名区分。 split_into_regular_and_task_span_metricsengine/metrics_evaluator.py:63)用 inspect.signaturemetric.score 有没有 task_span 参数:

类别判据能看到什么什么时候算
regular metric签名里没有 task_spandataset item 字段 + task 输出在 task 的 with 体内,紧接着执行
task span metric签名里 task_span整棵执行 span 树的根 span(token、延迟、子调用)全部跑完之后,另开一轮

task span 那一类要走完全不同的路:_execute_evaluationengine.py:534)发现有这类指标,就把整轮执行包进 local_recording.record_traces_locallyengine.py:583),跑完再从录到的 trace 树里,取第一个 span 当作 evaluation spanengine.py:462-467),并行地补算这批分数。没有 span 会直接抛 ValueError。没有 task span 指标时这层完全不启用(engine.py:566-567),省掉本地录制的开销。

字段名映射:scoring_key_mapping 解决"数据集字段名和指标参数名对不上"。

打分输入是 {**dataset_item, **task_output} 合并出来的一个平铺 dict(task 输出覆盖同名字段),然后逐条应用映射(evaluation/metrics/arguments_helpers.py:121-155):

# 示意,非源码
mapped = {**dataset_item, **task_output}
for target_key, source in scoring_key_mapping.items():
if callable(source):
mapped[target_key] = source(mapped) # 值可以是函数,现算
else:
mapped[target_key] = mapped[source] # 值是字符串,就是改名

方向容易搞反:key 是指标要的参数名,value 是数据里现有的字段名。 数据里叫 expected、指标参数叫 reference,就写 {"reference": "expected"}。源字段不存在时不报错,只打 debug 日志——错误留到调用指标时由参数校验报,那里能把"缺哪些参数、有哪些参数、映射里哪些没用上"一起说清(arguments_helpers.py:36-49)。

逐指标打分与容错。 _compute_metric_scoresmetrics_evaluator.py:269)对每个指标:

┌─ ScorerWrapperMetric? ──是──► 用原始 dataset_item / task_outputs 调,不套映射

└─ 否 ──► 校验参数齐不齐 ──► 能吃 trace_tool_context 就注入 ──► metric.score(**kwargs)

▼ 抛 ScoreMethodMissingArguments
默认容错级下直接向上抛(配置错,不该被吞)

▼ 抛其它异常
记 ERROR,产出 ScoreResult(value=0.0, scoring_failed=True, reason=str(e))

—— 分支在 metrics_evaluator.py:297-317,容错在 metrics_evaluator.py:331-356

三个细节值得单独点出:

  • "参数缺失"默认不容错、别的都容错。 容错分级由 ErrorToleranceevaluation/types.py:17)控制:默认 METRIC_ERRORSScoreMethodMissingArguments 原样上抛(metrics_evaluator.py:331-338),因为那是你映射没配对,跑下去也全是 0 分,早失败更好;显式选 ALL_SCORING_ERRORS 才连它也降级成失败分。
  • trace_tool_context 只注给"接得住"的指标。 _accepts_trace_tool_contextmetrics_evaluator.py:46)检查签名里显式写了这个参数、或者有 **kwargsLLMJudge 属于后者);否则不注入,避免窄签名的指标炸在"unexpected keyword argument"。
  • 失败的分是 0 不是缺失。 这样汇总时能看到"N failed",而不是分母悄悄变小。终端汇总里 _compute_average_scores 对失败项 score_counts += 0evaluation/report.py:26-33),所以均值不被 0 分拉低,但失败数会单独红字显示。

item 级指标。 build_metrics_evaluatormetrics_evaluator.py:181)除了 suite 级指标,还会读 item.evaluators 现场实例化 LLMJudge(目前只支持 type == "llm_judge",其它类型告警并记为一条失败分而非静默跳过,metrics_evaluator.py:146-163),然后把所有 LLMJudge 合并成一个(LLMJudge.merged(judges)metrics_evaluator.py:197-199)——多个断言合成一次裁判调用,省钱省时。


5. 数据与结果对象

5.1 Dataset:内容哈希去重与懒同步

要解决的问题: 反复 insert() 同一批数据(脚本重跑、增量补数据)不该产生重复行。

做法:按内容算 SHA-256,本地维护两份缓存。

字段类型作用
_hashesSet[str]已知内容哈希集合,插入时查它去重
_id_to_hashDict[str, str]item id → 哈希,删除时用它反查并从集合里摘掉
_hashes_syncedbool本地缓存是否和后端一致

—— 三个字段定义在 api_objects/dataset/dataset.py:346-355

哈希算的是"内容"而非整个对象:get_content()(额外字段)再加上 description / evaluators / execution_policy(若有),json.dumps(..., sort_keys=True) 后取 sha256(api_objects/dataset/dataset_item.py:91-106)。id 不参与,所以同样内容不同 id 也算重复。

懒同步是这里的精髓。 直接 create_dataset() 造出来的 Dataset,_hashes_synced = True——本地就是全部,没什么可同步的。而从后端拉回来的(from_publicdataset.py:387)设成 False:后端可能已经有本地没见过的 item。同步不在拉取时做,而是推迟到第一次 insert()

if not self._hashes_synced:
self.__internal_api__sync_hashes__()

—— dataset.py:620-621,同步实现在 dataset.py:687-700(流式拉全量、重算哈希)。

为什么这么设计?注释里说得很直白:避免在 list 数据集时付出 N+1 次同步(dataset.py:348-353)。你 list_datasets() 出 50 个数据集只是想看看名字,不该为此拉 50 份全量数据。

删除时要双向清理:从 _id_to_hash 查出哈希、从 _hashes 摘掉、再删映射(dataset.py:757-761),漏一步就会让删掉的内容再也插不进来。

5.2 DatasetVersion:版本与差异

DatasetVersion 是数据集在某个时间点的只读快照(dataset.py:148)。评估层关心的属性:

属性含义位置
version_hash该版本的内容指纹,流式取数时作为 dataset_version 参数下传dataset.py:212,用在 dataset.py:282
version_name人可读版本名(v1v2),断点续跑就是靠它锁定版本dataset.py:217
items_added / items_modified / items_deleted相对上一版的增改删条数dataset.py:237 / :242 / :247
is_latest是不是最新版dataset.py:227

Dataset.get_version_info()dataset.py:439)取最新一版;后端未开版本功能时会返回 403,这里吞掉并返回 Nonedataset.py:456-462)——这个 None 会一路影响到断点续跑能不能用(见 §6)。

5.3 Experiment:把三样东西串起来

一条 experiment item 只存两个引用 + 一份策略:

@dataclasses.dataclass
class ExperimentItemReferences:
dataset_item_id: str
trace_id: str
project_name: Optional[str] = None
execution_policy: Optional[Dict[str, Any]] = None

—— api_objects/experiment/experiment_item.py:11。分数不在这里,分数是挂在 trace 上的 feedback score。读回来时后端做 join,ExperimentItemContentexperiment_item.py:19)就带上了 dataset_item_dataevaluation_task_outputfeedback_scoresassertion_results 四样。

dataset item ─┐
├─► experiment item(只存 id 对) ──► experiment
trace ────────┘ │
│ └─ 读回时 join 出输出与分数
└─ feedback scores(分数真正存在这里)

插入走 streamer 异步批量(experiment.py:100-122,批大小 FEEDBACK_SCORES_MAX_BATCH_SIZE),上报管道细节见 上报层那一章

写分数分两类。 log_test_result_feedback_scoresevaluation/rest_operations.py:87)遍历分数:category_name == "suite_assertion" 的走 log_assertion_results,转成 passed/failed;其余走 log_traces_feedback_scoresscoring_failed 的一律不写(rest_operations.py:97-98)——失败的 0 分只活在本地汇总里,不污染后端统计。

experiment 级分数是另一条路:experiment_scoring_functions 吃整个 List[TestResult]ScoreResult,通过 update_experiment 直接写在 experiment 上(experiment.py:150-171)。它的执行是吞异常的——某个聚合函数崩了只打 warning,不影响返回结果(evaluation_result.py:33-38)。

5.4 结果对象与终端汇总

链路:TestCaseTestResultEvaluationResult → 两种视图。

对象关键字段位置
TestCasetrace_iddataset_item_idtask_outputdataset_item_contentmapped_scoring_inputsevaluation/test_case.py:9
TestResulttest_casescore_resultstrial_idtask_execution_timescoring_timeevaluation/test_result.py:10
EvaluationResultexperiment_iddataset_idtest_resultsexperiment_urltrial_countexperiment_scoresevaluation/evaluation_result.py:131
EvaluationResultOnDictItems只有 test_results(无 experiment 场景)evaluation/evaluation_result.py:219

EvaluationResult 给两种视图:

  • aggregate_evaluation_scores()evaluation_result.py:142)→ 按指标名聚合,每个指标一份 ScoreStatistics{mean, max, min, values, std}
  • group_by_dataset_item_view()evaluation_result.py:170)→ 按 dataset item 分组,组内按 trial_id 排序,每组再算一份统计。看某条数据在多次 trial 间稳不稳定用这个。

calculate_aggregated_statisticsevaluation/score_statistics.py:21)只纳入 scoring_failed=False 且值有限(math.isfinite)的分;样本数 < 2 时 std 为 None 而不是 0(score_statistics.py:51)——1 个样本的标准差没定义,给 0 是撒谎。

终端汇总用 rich 画evaluation/report.py):

函数触发条件输出
display_experiment_resultsreport.py:44verbose >= 1面板:总耗时、样本数(trial>1 时显示 "N items (M runs)")、各指标均分 + 失败数
display_evaluation_scores_statisticsreport.py:135verbose >= 2表格:每个指标的 Mean / Min / Max / Std
display_experiment_linkreport.py:119verbose >= 1一行可点击的 experiment 链接

6. 断点续跑与抽样

6.1 resume 要解决的问题

跑 1000 条、跑到 700 条网断了。重跑全部既费钱又费时。目标:只补跑没跑完的,且返回的结果看起来像完整跑了一遍。

6.2 状态存在哪:两份,一份在后端一份在本地

存储存什么为什么这么分
experiment 的 experiment_config["_opik_resume"]一个 JSON 字符串:default_runs_per_itemdataset_filter_stringdataset_version_namenb_samplesrequires_local_checkpoint都是小而可复现的配置,不含解析后的数据列表
~/.opik/resume/<experiment_id>.json解析后的 dataset item id 列表只有配置重建不了迭代顺序时才写(用了 sampler 或显式 ids)

—— schema 定义在 evaluation/resume/state.py:40-64,编码成单个 JSON 字符串是为了让 experiment 配置界面只显示一行而不是每个字段一行(state.py:69-75)。本地 checkpoint 在 evaluation/resume/checkpoint.py,写入走临时文件 + os.replace 保证原子(evaluation/resume/checkpoint.py:49-51)。

这里有个明确的拒绝:没有钉住数据集版本就不许 resume。 resume_state_for_evaluateresume/integration.py:42)拿不到 version_name 时写入的是 NonResumableState,理由字符串写死在 integration.py:35-39。解码端还要再拦一次:即使标了 resumable,缺 dataset_version_name 也降级为不可续(state.py:143-151)。原因很实在——对着会变的 dataset HEAD 迭代,会悄悄多算或漏算原来那次跑过的 item。

6.3 续跑的四步

prepare_resume_context(client, experiment_id) resume/context.py:63
│ 读后端 resume blob → 必须是 ResumableState,否则抛 ExperimentNotResumable
│ 需要 checkpoint 就读本地文件,读不到抛 LocalCheckpointMissing
│ 按 version_name 取 DatasetVersion(钉版本)
│ 统计每条 item 已完成几趟(output 非 None 才算数)

_resolve_resume_items(context) evaluator.py:1558
│ 有 checkpoint 用 checkpoint 的 id 列表;否则按原 filter + nb_samples 重新迭代

build_pending_items_iterator(items, context) resume/iteration.py:56
│ remaining = max(0, expected - completed)
│ remaining == 0 就跳过;否则把 item 的 runs_per_item 改写成 remaining

_evaluate_task(...) 照常跑 → merge_resume_results 合并 evaluator.py:1519

为什么改写 runs_per_item 就够了? 因为引擎本来就支持 item 级策略覆盖(§4.3),resume 只要把 item 预先标注好,引擎不需要任何改动。这一点原文写在 resume/iteration.py 的模块 docstring 里。

合并要小心快照时机。 reconstruct_previous_test_resultsresume/merge.py:27)从后端已有的 experiment item 重建 TestResult——直接拿存着的 evaluation_task_outputfeedback_scores不重新算指标、不写新分数。关键是这个快照必须在 _evaluate_task 开跑之前取(evaluator.py:1595-1601),否则这次新写的 experiment item 会被一起重建进来,重复计数。

聚合分要重算并覆盖。 _evaluate_task 内部已经就"本次这一小片"算过一遍 experiment 级分数写上去了;evaluate_resume 拿合并后的全集再算一遍,覆盖写(evaluator.py:1635-1641)。注释坦承两次写之间有一小段窗口后端上是"只有片段"的分数,权衡后接受。

6.4 抽样

samplers/ 只有一个抽象基类和一个实现:

作用位置
BaseDatasetSampler抽象方法 sample(items) -> List[DatasetItem]evaluation/samplers/base_dataset_sampler.py
RandomDatasetSampler随机取 max_samples 条,可选 shuffle 与 seedevaluation/samplers/random_dataset_sampler.py

抽样有个明确代价:关掉流式。 采样器的接口吃的是完整 list,所以 resolve_dataset_itemsevaluation/helpers.py:80)在有 sampler 时会先把整个流物化成 list,日志明说 "Dataset streaming disabled due to sampler"evaluation/helpers.py:111)。返回值再包回迭代器,让上层形状一致。返回值不是 list 直接抛 TypeError,明确不支持流式采样器(evaluation/helpers.py:121-125)。

RandomDatasetSampler.sample 里有个小优化:rng.sample 再 shuffle,而不是先 shuffle 整个数据集——大数据集上省一大截(random_dataset_sampler.py:42-46)。

三种取数模式与 checkpoint 的关系_materialize_for_checkpointevaluator.py:99)统一裁决:

模式迭代方式写 checkpoint 吗
有 sampler已物化,抽干一次取 id 再重新包迭代器写(抽样后的 id)
只有显式 dataset_item_ids保持惰性,id 本来就知道
都没有全程流式,batch 200(evaluation/helpers.py:10不写,配置足以重建

sampler 优先于显式 ids 是刻意的:checkpoint 必须记录引擎实际迭代的那批,不是原始输入(evaluator.py:114-119)。


7. 巧妙之处(可借鉴)

  • 用日志字段当事务标志。 不额外建状态表,"trace 有没有 output"就是"这一趟跑完没有"。写点在 engine.py:376,擦除点在 engine/helpers.py:51-56,读点在 resume/context.py:186。三处一致,语义就闭合了。
  • 无状态引擎 + 参数传流程数据。 同一个 EvaluationEngine 实例被多线程复用不会串味(engine.py:67-73)。
  • 有超时的 drain 而不是 flush。 只等本地处理器、不等网络;超时降级不报错(engine.py:245-253)。这是"一致性 vs 可用性"在客户端 SDK 里一个很小但很实的取舍。
  • 命名空间 sentinel 标签胜过按名字过滤。 __opik_eval_internal__ 只有引擎会打,用户 span 再怎么起名也不会误伤;配合"父子闭包扫描"整棵子树一起丢(context.py:182context.py:185-229)。
  • 懒同步哈希缓存。 把 N+1 的同步开销从"列举时"推到"第一次写入时",读多写少的场景直接省掉(dataset.py:348-353dataset.py:620-621)。
  • 先声明组大小再提交任务。 一行 set_group_size 消掉了回调与提交之间的竞态(engine.py:420evaluation_tasks_executor.py:94-100)。
  • 配置错早失败、运行错晚失败。 参数缺失原样上抛,其它异常降级成 scoring_failed 的 0 分(metrics_evaluator.py:331-356)。区分标准是"这错重跑还会不会一样"。
  • 样本不足时 std 给 None 而不是 0score_statistics.py:51)。小事,但避免了下游把 0 当成"很稳定"。

8. 边界与局限

明说的不支持:

  • test suite 不能 resume。 __internal_api__run_test_suite__ 刻意不写 resume 状态,原文注释说持久化原语已经预留、以后可加(evaluator.py:442-445)。
  • 流式采样器不支持。 sample() 必须返回 list,返回别的直接 TypeErrorevaluation/helpers.py:121-125)。
  • item 级 evaluator 只认 llm_judge 其它 type 打 warning 并记为失败分(metrics_evaluator.py:146-163dataset.py:483-498)。
  • 没开版本功能的数据集不能 resume。 见 §6.2。

会崩或会让人意外的地方:

  • 一个 item 失败 = 整轮抛异常。 get_results 逐个 future.result() 重抛第一个异常(evaluation_tasks_executor.py:232-233)。想"跳过坏数据继续跑"得自己在 task 里 try。
  • task span 指标要求至少有一个 span。 没有就 ValueError,并且默认拿第一个 span 当 evaluation spanengine.py:462-467)——依赖 span 顺序,注释也只写了 "the first span is the evaluation span"。
  • evaluate_resume 不校验你传的 task 和上次是不是同一个。 docstring 明说框架不做一致性检查(evaluator.py:1539-1550),scoring_key_mapping 同理。
  • 本地 checkpoint 绑机器。 换台机器 resume 就会抛 LocalCheckpointMissingresume/context.py:94-99)。
  • _hashes 只在本进程内准。 两个进程同时往一个 dataset 插同样内容,各自的本地集合看不见对方,去重会漏。
  • project_name 参数已废弃。 数据集自带 project_name 时,用户传的会被忽略并告警(evaluation/helpers.py:13-41)。

9. 横向对比与延伸阅读

本章只讲编排。同组其它章:

章节与本章的接口
01 追踪层引擎用 @opik.track 把用户 task 包成 span 树;本章的 trace/span 语义由那一章定义
02 上报层experiment item、feedback score 都走 streamer 异步批量;本章的 drain 就是在跟这条管道打交道
03 服务端存储experiment item 读回时的 join(dataset item + trace + feedback score)在服务端完成
05 打分器本章只管调用 metric.score(...);分数怎么算、LLM 裁判怎么提问在那一章
06 生产闭环evaluate_optimization_trialevaluate_on_dict_items 是优化器的两个调用入口

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

路径相对克隆根 opik/

主题文件路径符号名
通用评估入口sdks/python/src/opik/evaluation/evaluator.pyevaluate
提示词评估sdks/python/src/opik/evaluation/evaluator.pyevaluate_prompt_build_prompt_evaluation_task
重新打分(不跑 task)sdks/python/src/opik/evaluation/evaluator.pyevaluate_experiment
优化器试验sdks/python/src/opik/evaluation/evaluator.pyevaluate_optimization_trial
轻量内存评估sdks/python/src/opik/evaluation/evaluator.pyevaluate_on_dict_items
断点续跑入口sdks/python/src/opik/evaluation/evaluator.pyevaluate_resume_resolve_resume_items
测试套件入口sdks/python/src/opik/evaluation/evaluator.pyrun_tests__internal_api__run_test_suite__
checkpoint 物化裁决sdks/python/src/opik/evaluation/evaluator.py_materialize_for_checkpoint
引擎主体sdks/python/src/opik/evaluation/engine/engine.pyEvaluationEnginerun_and_scorescore_test_cases
单条执行 + trace 包装sdks/python/src/opik/evaluation/engine/engine.py_compute_test_result_for_llm_task
streamer 排干 + agentic 上下文sdks/python/src/opik/evaluation/engine/engine.py_build_trace_tool_contextDEFAULT_STREAMER_DRAIN_TIMEOUT_SECONDS
执行策略合并sdks/python/src/opik/evaluation/engine/engine.pyget_item_execution_policy
task span 指标补算sdks/python/src/opik/evaluation/engine/engine.py_execute_evaluation_update_test_result_with_task_span_metrics
trace 生命周期 + experiment itemsdks/python/src/opik/evaluation/engine/helpers.pyevaluate_llm_task_contextEvaluationContextState
并发执行与进度sdks/python/src/opik/evaluation/engine/evaluation_tasks_executor.pyStreamingExecutorexecuteset_group_size
限流识别sdks/python/src/opik/evaluation/engine/exception_analyzer.pyis_llm_provider_rate_limit_error
指标分流与打分sdks/python/src/opik/evaluation/engine/metrics_evaluator.pysplit_into_regular_and_task_span_metrics_compute_metric_scoresMetricsEvaluator
item 级 evaluator 合并sdks/python/src/opik/evaluation/engine/metrics_evaluator.pybuild_metrics_evaluator_extract_item_evaluators
字段名映射sdks/python/src/opik/evaluation/metrics/arguments_helpers.pycreate_scoring_inputsraise_if_score_arguments_are_missing
内部 span 过滤sdks/python/src/opik/evaluation/suite_evaluators/agentic/context.pyINTERNAL_SPAN_TAG_filter_internal_spansbuild_trace_tool_context_from_trace_data
数据源解析与抽样sdks/python/src/opik/evaluation/helpers.pyresolve_dataset_itemsEVALUATION_STREAM_DATASET_BATCH_SIZE
抽样器sdks/python/src/opik/evaluation/samplers/random_dataset_sampler.pyRandomDatasetSampler
数据集去重sdks/python/src/opik/api_objects/dataset/dataset.py__internal_api__insert_items_as_dataclasses____internal_api__sync_hashes__
item 内容哈希sdks/python/src/opik/api_objects/dataset/dataset_item.pyDatasetItem.content_hashExecutionPolicyItem
版本信息sdks/python/src/opik/api_objects/dataset/dataset.pyDatasetVersion.version_hashitems_addedget_version_info
experiment 归档sdks/python/src/opik/api_objects/experiment/experiment.pyExperiment.insertlog_experiment_scores
experiment item 结构sdks/python/src/opik/api_objects/experiment/experiment_item.pyExperimentItemReferencesExperimentItemContent
分数写回sdks/python/src/opik/evaluation/rest_operations.pylog_test_result_feedback_scoresget_experiment_test_cases
结果对象sdks/python/src/opik/evaluation/evaluation_result.pyEvaluationResultmerge_resume_resultscompute_experiment_scores
统计聚合sdks/python/src/opik/evaluation/score_statistics.pycalculate_aggregated_statisticsScoreStatistics
终端汇总sdks/python/src/opik/evaluation/report.pydisplay_experiment_resultsdisplay_evaluation_scores_statistics
续跑状态 schemasdks/python/src/opik/evaluation/resume/state.pyResumableStateNonResumableStateread_resume_state
本地 checkpointsdks/python/src/opik/evaluation/resume/checkpoint.pywrite_checkpointread_checkpoint
续跑上下文sdks/python/src/opik/evaluation/resume/context.pyprepare_resume_contextis_trial_fully_completed
待办计算sdks/python/src/opik/evaluation/resume/iteration.pybuild_pending_items_iteratorremaining_runs_for_item
历史结果重建sdks/python/src/opik/evaluation/resume/merge.pyreconstruct_previous_test_results
通过阈值判定sdks/python/src/opik/api_objects/dataset/test_suite/suite_result_constructor.pybuild_suite_result