跳到主要内容

数据截至 (上游 commit 36c7a7f6eca6)

落库、查询与生产监控:服务端怎么存 trace、怎么被搜、怎么持续打分

30 秒导读: 前五章都站在应用进程里看——span 怎么产生、怎么被导出、怎么被评估。 这一章翻到另一面:trace 落到服务端之后的一生。它进了哪张表、怎么被一句 filter 字符串搜出来、 谁在半夜把它搬去对象存储、谁在生产环境里不停地给它打分。

本章不重复 SDK 侧的导出逻辑(那在 02 章),只讲服务端与长期运行


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

一句话定义: MLflow 的服务端存储层,是一套「把 trace 拆成可查询的摘要 + 可归档的明细」的双层存储, 外加两个每分钟醒一次的后台调度器。

为什么需要它

前面几章的产物是内存里的 Trace 对象。要让它变成生产可用的观测系统,还差三件事:

缺的东西具体问题本章对应节
落盘一条 trace 可能有几百个 span、几 MB JSON,全塞进关系库会撑爆§3、§4
检索「找出昨天延迟 > 5s 且 safety 判定为 no 的 trace」要能秒回§5
持续打分生产流量不会停,得有人自动采样、自动跑 judge、自动写回§7

用起来什么样

对使用者,整条链路只是三个命令:

# 1. 起一个用 Postgres 做后端的 tracking server
mlflow server --backend-store-uri postgresql://... --artifacts-destination s3://my-bucket

# 2. 应用把 trace 打进去(MLflow SDK,或任何标准 OTLP 客户端)
export OTEL_EXPORTER_OTLP_TRACES_ENDPOINT=http://localhost:5000/v1/traces

# 3. 查(注意选项名是 --filter-string,没有 --filter)
mlflow traces search --experiment-id 1 --filter-string "trace.status = 'ERROR'"

第三条的选项名定义在 mlflow/cli/traces.py:172--experiment-id 是必填项,也可以用 MLFLOW_EXPERIMENT_ID 环境变量给(mlflow/cli/traces.py:93)。

一句话直觉

把 trace 当成一封邮件: 信封(发件人、时间、主题、大小)进数据库,方便按条件筛; 正文附件(几 MB 的 span 明细)另存一处,只在你真的点开某封信时才去取。 本章大半篇幅在讲这个「信封 / 附件」的分家怎么做、以及附件后来还搬了几次家。


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

怎么读这张图: 从上往下是一次写入的流向;右侧两个方框是独立于请求的后台循环, 它们不接请求,靠 Huey 定时任务被唤醒。

MLflow SDK 客户端 任意 OTLP 客户端
(StartTraceV3 / REST) (Claude Code / Codex / …)
| |
v v
+---------------------+ +----------------------+
| Flask 路由 | | FastAPI /v1/traces |
| handlers.py | | otel_api.py |
+---------------------+ +----------------------+
\ /
\ /
v v
+--------------------------------+ +-------------------+
| AbstractStore(存储抽象) |<---------| ① 归档调度器 |
| SqlAlchemy / File / Rest 三选一 | | 每分钟一次 |
+--------------------------------+ +-------------------+
| | +-------------------+
v v <-----| ② 在线打分调度器 |
+------------------+ +-----------------+ | 每分钟一次 |
| 关系表 | | span 明细 | +-------------------+
| trace_info + 从表 | | spans 表 / S3 |
+------------------+ +-----------------+

部件一句话职责

部件干什么在哪个文件
AbstractStore定义 trace 相关的 30 来个方法签名,绝大多数默认 raise NotImplementedErrormlflow/store/tracking/abstract_store.py
SqlAlchemyStore唯一功能完整的实现,近 1 万行,本章主角mlflow/store/tracking/sqlalchemy_store.py
FileStore每条 trace 一个目录,搜索靠全量扫描后内存过滤mlflow/store/tracking/file_store.py
RestStore客户端侧实现,把方法调用翻成 HTTP 打给服务端mlflow/store/tracking/rest_store.py
Flask handlersMLflow 自家 REST API 的入口,protobuf 收发mlflow/server/handlers.py
FastAPI OTLP 路由标准 OTLP/HTTP 摄取口,非 MLflow 客户端也能进mlflow/server/otel_api.py
归档调度器把过了保留期的 trace 明细搬去对象存储mlflow/tracing/trace_archival_service.py
在线打分调度器按采样率把新 trace 喂给注册好的 scorermlflow/genai/scorers/job.py

主线走一遍

  1. 客户端调 POST /mlflow/traces,落到 _start_trace_v3()mlflow/server/handlers.py:4085)。
  2. handler 把 protobuf 转成 TraceInfo 实体,交给 _get_tracking_store().start_trace(trace_info)
  3. SqlAlchemyStore.start_tracetrace_info 主行 + tags / metadata / metrics / assessments 从行 (mlflow/store/tracking/sqlalchemy_store.py:3693)。
  4. span 明细走另一条路:要么 SDK 直接传 artifact 仓,要么经 log_spans()spans 表。
  5. 后来查询时,search_traces 只碰关系表;真要看 span 才按 SPANS_LOCATION 标签去对应位置取。

3. 存储抽象与三种实现

3.1 抽象层长什么样

AbstractStore 用「默认抛 NotImplementedError」而不是 @abstractmethod 来声明 trace 相关方法。 这是刻意的——三个实现的能力差得很远,强制全实现会让 FileStore 塞满假方法。

trace 相关的接口大致分五组:

方法抽象层行号
写入start_traceset_trace_tagdelete_trace_taglink_traces_to_runlink_prompts_to_traceabstract_store.py:285:583:594:1520:1547
读取get_trace_infoget_tracebatch_get_tracesbatch_get_trace_infos:350:362:378:394
查询search_tracesquery_trace_metricscalculate_trace_filter_correlation:481:551:1562
评价get_assessmentcreate_assessmentupdate_assessmentdelete_assessment:604:618:633:663
生命周期delete_traces / _delete_tracesarchive_tracesresolve_trace_archival_config:297:341:412:444

一个值得学的小设计:模板方法拆参数校验。 delete_traces 是 public 入口,只做参数互斥校验 (max_timestamp_millistrace_ids 不能同时给、max_traces 必须为正),校验完才转给子类实现的 _delete_tracesabstract_store.py:339)。三个后端因此不必各写一遍同样的校验。

另有一个能力开关:supports_trace_archival 属性,基类返回 FalseSqlAlchemyStore 覆盖为 Trueabstract_store.py:81 vs sqlalchemy_store.py:395)。归档调度器靠它决定要不要空转。

3.2 三个实现怎么选

选哪个不是配置项,而是由 tracking URI 的 scheme 决定,注册表在 mlflow/tracking/_tracking_service/utils.py:243_register_tracking_stores()

URI 形态实现适用场景
空 / file://FileStore本地跑着玩,无 server
postgresql://mysql://sqlite://SqlAlchemyStore生产,服务端进程内使用
http://https://RestStore客户端侧,转发给远端 server
databricksdatabricks-ucDatabricks 专用 REST store托管环境

三者的取舍差别,最集中体现在 search_traces 上:

SqlAlchemyStore.search_traces → 拼 SQL,让数据库做过滤/排序/分页
(sqlalchemy_store.py:3713)

FileStore.search_traces → 遍历 experiment 下每个 trace 目录读 meta,
(file_store.py:2187) 全读进内存后 SearchTraceUtils.filter/sort/paginate

RestStore.search_traces → 组 SearchTracesV3 protobuf,POST /mlflow/traces/search
(rest_store.py:575) 自己不存任何东西

FileStore 的实现只有二十来行,因为它把过滤完全外包给了纯 Python 的 SearchTraceUtils.filter()——正确性一致,但是 O(全量 trace)。生产别用。

RestStore 还有一层版本降级逻辑值得看:start_trace 先打 V3 端点,收到 ENDPOINT_NOT_FOUND 就退回 V2 的 deprecated_start_trace_v2 + deprecated_end_trace_v2 两步式 (rest_store.py:441)。同一个客户端能同时对着新旧服务端工作。


4. 关系表结构:信封与附件的分家

4.1 先看一眼表

trace_info 是主表,其余都挂在它下面。所有从表都用 ondelete="CASCADE" 的外键, 删主行时数据库负责清干净。

表(类名)表名存什么主键行号
SqlTraceInfotrace_info信封:时间、耗时、状态、请求/响应预览request_idmodels.py:768
SqlTraceTagtrace_tags可变标签,值上限 8000 字符(request_id, key):806
SqlTraceMetadatatrace_request_metadata不可变元数据(source run、session、token usage 原文)(request_id, key):836
SqlTraceMetricstrace_metricstrace 级数值指标:token 数、成本(request_id, key):867
SqlSpanMetricsspan_metricsspan 级数值指标(trace_id, span_id, key):903
SqlAssessmentsassessments评价:feedback / expectation / issue 三合一assessment_id:944
SqlIssueissues归纳出来的问题条目(有 status / severity / root_causes)issue_id:1154
SqlEvaluationDatasetevaluation_datasets评估数据集dataset_id:1554
SqlSpanspans附件的其中一种落法:整个 span 的 JSON(trace_id, span_id):1957
SqlEntityAssociationentity_associations通用关联边:trace↔run、trace↔prompt versionassociation_id:2052

trace_info 只建了一条复合索引 (experiment_id, timestamp_ms),注释直说原因:UI 默认视图就是 「某 experiment 下按时间倒序」,且每个搜索查询的 where 里必带 experiment_idmodels.py:828)。

4.2 「TraceInfo 进表、span 明细进 artifact」的分层

这句老话现在要加个补丁:span 明细有三个可能位置,由一个标签 mlflow.trace.spansLocationTraceTagKey.SPANS_LOCATION)指路,取值定义在 mlflow/tracing/constant.py:211SpansLocation 枚举。

trace_info 行(信封,永远在关系库)
|
读 tags['mlflow.trace.spansLocation']
|
+-----------------------+------------------------+
| | |
TRACKING_STORE ARTIFACT_REPO ARCHIVE_REPO
spans 表逐行 traces.json 单文件 OTLP protobuf 单文件
(log_spans 写的) (SDK 直传 S3 的) (归档器搬过去的)
| | |
v v v
SELECT content download_trace_data() download_archived_trace_data()

真实分派就在一个函数里,_get_spans_with_trace_infosqlalchemy_store.py:6568):

spans_location = trace_info.tags.get(TraceTagKey.SPANS_LOCATION)
if spans_location == SpansLocation.ARCHIVE_REPO.value:
artifact_repo = get_artifact_repository(get_archive_uri_for_trace(trace_info))
return artifact_repo.download_archived_trace_data().spans
if spans_location != SpansLocation.TRACKING_STORE.value:
raise MlflowTracingException("Trace data not stored in tracking store")

注意第三种情况是抛异常而不是自己去取:注释解释,ARTIFACT_REPO 的 URI 可能是 mlflow-artifacts:// 这种只有 handler 层才解析得了的代理协议,所以让上层去 fallback (sqlalchemy_store.py:6576-6578)。接住这个异常的就是 _fetch_trace_data_from_store,它把 MlflowTracingException 转成 None 表示「你去 artifact 仓拿」 (mlflow/server/handlers.py:4390)。

artifact 路径怎么算出来的?start_trace 一开始就往 tags 里塞了一条 mlflow.artifactLocation,值是 <experiment 的 artifact_location>/traces/<trace_id>/artifacts_get_trace_artifact_location_tagsqlalchemy_store.py:3450)。注释写明用 /traces 子目录 隔离于 run 的 artifact。文件名固定 traces.jsonmlflow/tracing/utils/artifact_utils.py:6)。

4.3 写入路径上的两处硬功夫

其一:start_trace 与 log_spans 的竞态。 两条路径都可能先创建 trace_info 行。 start_trace 的写法是先乐观 INSERT,撞 IntegrityError 再补救sqlalchemy_store.py:3527):回滚、expunge_all() 清掉失败的 ORM 树、加行锁重读、 然后把 tags / assessments / metadata / metrics 逐条 session.merge()(也就是 upsert)合并进去。

补救分支里有一条注释点破了为什么必须重建子行对象:不重建的话,后面的逐行 merge 会把失败父对象里的 陈旧 trace_info 状态一起拖回来(:3543-3545)。

其二:谁说了算。 start_trace 写完会额外塞一条元数据 TraceMetadataKey.TRACE_INFO_FINALIZED = "true",用途写在注释里:告诉并发的 log_spans() 不要覆盖 request_time / execution_duration / session_id / token usage / cost 这些 trace 级权威值 (sqlalchemy_store.py:3512-3515)。

SqlTraceMetrics 的填法也在这里:从元数据里的 mlflow.trace.tokenUsage JSON 串里, 按 TokenUsageKey.all_keys() 抽出 input/output/total/cache 各项,转成 float 存成独立行 (sqlalchemy_store.py:3501-3509)。同一份数据存两遍——JSON 原文进 metadata 表给人看, 拆开的数值进 metrics 表给 SQL 聚合用。

4.4 assessments 表:一张表装三种东西

assessment_type 列取 "feedback" / "expectation" / "issue" 三值之一,value 列是 JSON 文本。 to_mlflow_entity() 按类型分派回 Feedback / Expectation / IssueReference 三种实体 (models.py:1081-1144)。

两个字段撑起了「评价可被修正」这件事:

  • overrides:指向被它覆盖的旧 assessment_id。
  • valid:布尔,被覆盖的那条会被置 False

create_assessment 建新记录时,顺手把 overrides 指向的旧记录 update({"valid": False})sqlalchemy_store.py:4419);而 delete_assessment 会把这件事撤销回来—— 删掉一条覆盖记录时,被它压住的原记录恢复 valid=Truesqlalchemy_store.py:4633)。 删不存在的 assessment 直接 return,是幂等的。

create_assessment 还有一段专门处理外键报错的代码:flush 撞 IntegrityError 时先回滚, 再查一次 trace 在不在,在就报「约束冲突」,不在就报「trace 已被删除」——为的是不把裸 SQL 错误 泄给用户(sqlalchemy_store.py:4436-4462)。

4.5 关联边:trace ↔ run ↔ prompt

link_traces_to_runsqlalchemy_store.py:4670)和 link_prompts_to_trace:4726) 都往同一张 entity_associations 表写,只是 source_type / destination_type 不同。 两者都先查一遍已有关联再只插差集,所以重复调用是安全的

link_traces_to_run 有个硬上限 MAX_TRACE_LINKS_PER_REQUEST,超了直接报错(:4684-4686)。 prompt 的 destination_id 用 "name/version" 拼串当主键(:4738)。


5. 搜索:从一句 filter 字符串到 SQL

5.1 要解决的小问题

用户想写这样一句话,然后立刻拿到结果:

trace.status = 'ERROR'
AND tags.env = 'prod'
AND span.type = 'LLM'
AND feedback.safety = 'no'

难点是这四个条件落在四张不同的表上,而且不能简单 join——比如 span.name = "x" AND span.status = "OK" 必须匹配同一个 span,不能一个条件命中 span A、 另一个命中 span B。

5.2 三段式流水线

filter 字符串
|
| ① 词法/语法解析:切成 {type, key, comparator, value} 列表
v SearchTraceUtils.parse_search_filter_for_search_traces
(mlflow/utils/search_utils.py:1688 起的类)
|
| ② 按 type 分桶,各自造 SQLAlchemy 片段
v _get_filter_clauses_for_search_traces
(sqlalchemy_store.py:9345)
|
+--> attribute_filters 直接 WHERE 在 trace_info 列上
+--> non_attribute_filters tags / metadata 的子查询,等待 join
+--> span_filters span / assessment / issue / prompt 的子查询
+--> run_id_filter 特殊,单独处理
|
| ③ 挂到主 statement 上
v _apply_trace_filter_clauses (sqlalchemy_store.py:3645)

第二步返回的是个四元组,这个四元组签名被三处复用,是这一层设计的关键: search_traces:3757)、_build_first_trace_filter_subquery:3944)、 _build_trace_filter_subquery:4837)。

5.3 每种标识符落到哪

SearchTraceUtils 里声明了七个标识符(_IDENTIFIERSsearch_utils.py:1766-1774), 外加四个别名(_ALTERNATE_IDENTIFIERS:1760-1765tagstagattributesattributetraceattributemetadatarequest_metadata)。

写法内部 type翻成什么 SQL代码位置
trace.status = 'OK'attributeWHERE trace_info.status = 'OK'sqlalchemy_store.py:9390-9401
trace.end_time_ms > Xattributetimestamp_ms + coalesce(execution_time_ms, 0) 表达式:9391-9395
tags.env = 'prod'tag子查询 trace_tags WHERE key=… AND value=… 再 inner join:9572-9580
tags.env IS NULLtagNOT EXISTS(...) 关联子查询:9449-9457
metadata.mlflow.sourceRun = 'r1'request_metadata特判,抽成 run_id_filter:9404-9411
span.name = 'foo'span条件累积到 span_filter_conditions,最后合成一个子查询:9513
span.attributes.model = 'gpt'spanspans.content JSON 文本上做 LIKE / RLIKE:9478-9507
feedback.safety = 'no'feedbackassessments 子查询 + 同 session 兄弟 trace 的并集:9549-9564
issue.id = 'iss-1'issueassessments WHERE assessment_type='issue' AND name=值:9374-9387
prompt = 'name/1'tag(特判)entity_associations 子查询:9416-9447

span 条件为什么要攒起来一次成型? 源码注释画了个反例(:9582-9588):

span 1. name: foo status: OK
span 2. name: search_web status: ERROR

span.name = "search_web" AND span.status = "OK",如果分成两个子查询各 join 一次, 这条 trace 会被错误命中。所以所有 span 条件先攒进 span_filter_conditions, 最后用一个 SqlSpan 查询把它们 filter(*conditions) 一起放进去(:9589-9597)。

assessment 条件里的 session 扩散是另一个不显然的点。_get_session_scoped_trace_ids (定义在 :9321,两处调用在 :9528:9560)做两步:先找出「带 session 作用域标记的 assessment」所属的 session id,再把这些 session 下的全部 trace id 拉出来。 搜 feedback 时,直接命中的 trace 和这些「同 session 兄弟」取并集(:9545:9563)。 也就是说,给整个会话打的分,会让该会话每条 trace 都被搜到

5.4 join 的正确性靠什么保证

search_traces 主体只有六十来行(sqlalchemy_store.py:3713),核心是三步: selectinload 预加载 tags/metadata/assessments → 挂过滤 → offset/limit。

没有做 DISTINCT,注释给了理由(:3775-3781):右侧子查询都被 (key, value) 条件 限死,而 (request_id, key) 在 tags / metadata 表里是主键,所以 join 的右表在 trace_id 上唯一, 不会放大行数。注释最后一句是警告:改查询构造逻辑时小心别破坏这个唯一性,否则就得加去重,而去重很贵。

排序由 _get_orderby_clauses_for_search_traces:9267)生成。两个细节:

  • 对每个 order by 键额外生成一个 CASE WHEN value IS NULL THEN 1 ELSE 0 列, NULL 值统一排到后面:9298)。
  • 无论用户排什么,末尾都追加 timestamp_ms DESC, request_id ASC 作为兜底和 tie-breaker (:9308-9318)——这个稳定序是后面在线打分 checkpoint 能工作的前提。

分页是偏移量式的:token 里就是个 offset 数字 (SearchTraceUtils.parse_start_offset_from_page_token / create_page_tokenmlflow/utils/search_utils.py:942:986)。

5.5 会话搜索:找「已经聊完」的会话

find_completed_sessionssqlalchemy_store.py:3799)回答的问题是: 哪些 session 的最后一条 trace 落在 [min, max] 窗口内,且窗口之后再没有新 trace(说明会话结束了)。

源码注释里给的例子最清楚(:3855-3858):给定 min=200、max=400,

会话trace 时间戳判定
A100, 200, 300收,最后一条 300 在窗口内
B100, 200, 500不收,500 > 400,还在进行中
C50, 100不收,太旧,没有 ≥200 的 trace

实现是四层子查询叠出来的:候选(_build_candidate_sessions_subquery:3901)→ 首条 trace 过滤(_build_first_trace_filter_subquery:3927)→ 统计首末时间 (_build_session_stats_subquery:4007)→ 取完成的(_build_completed_sessions_query:4043)。

其中 _build_first_trace_filter_subquery 是个巧妙的复用:它把 filter 字符串解析成同样的四元组, 再用 MIN(timestamp_ms) GROUP BY session_id 找出每个会话的首条 trace, 最后 join 到「timestamp 等于该会话最小值」上——于是「filter 只作用于会话第一条 trace」这个语义 被表达成了纯 SQL(:3971-3992)。

5.6 相关性分析:两个 filter 之间的 NPMI

calculate_trace_filter_correlation:4772)回答的是「用了 retrieval 的 trace,是不是更容易出错」 这类问题。做法是把两个 filter 各自变成一个 trace_id 子查询 (_build_trace_filter_subquery:4831——注意它比 _apply_trace_filter_clauses 更简单, 只 join 不处理 run_id),然后一次 SQL 数出四个计数:total / filter1 / filter2 / joint, 再交给 trace_correlation.calculate_npmi_from_counts 算标准化点互信息。

_get_trace_correlation_counts 用 LEFT JOIN 而不是 EXISTS,注释说明是为了 MSSQL 兼容 ——MSSQL 不支持子 join 里的 exists(:4857-4858)。

5.7 指标聚合:另一条查询路径

query_trace_metrics:4097)不返回 trace,返回时间序列数据点。它先自己按 experiment_id 和时间范围收窄 _trace_query,再把这个 query 交给 mlflow/store/tracking/utils/sql_trace_metrics_utils.py:730query_metrics 拼聚合。

三个视图类型决定从哪张表聚(mlflow/entities/trace_metrics.py:8):

view_type数据来源
TRACEStrace_info + trace_metrics
SPANSspans + span_metrics
ASSESSMENTSassessments

percentile 聚合在 MSSQL / MySQL 上需要窗口函数套子查询,走单独分支 (sql_trace_metrics_utils.py:778)。分页尚未实现——源码里就是一行 # TODO: Implement pagination with page_tokensqlalchemy_store.py:4142)。

5.8 客户端侧的分页包装

服务端一次最多给 500 条(_search_traces_v3 里硬校验,handlers.py:4170)。 用户想要更多,靠 mlflow.search_traces() 在客户端循环。

mlflow/tracing/fluent.py:992search_traces 做四件事:

  1. 决定返回类型——装了 pandas 默认返回 DataFrame,否则 list(:1088)。
  2. flush=True 时先把异步导出队列排空,避免刚打的 trace 搜不到(:1116)。
  3. 定义一个 pagination_wrapper_func 闭包,交给通用的 get_results_from_paginated_fnmlflow/utils/__init__.py:255)反复调用直到够数(:1136-1155)。
  4. 对 UC 表位置(location 里带 .)且没写时间过滤时,发一条 warning 说会很慢很贵(:1125-1133)。

search_sessionsfluent.py:1164)则是两阶段:先翻页扫 trace 收集去重后的 session id, 再用线程池按 session id 并发 search_traces 拉全量(:1268-1276)。 线程数取 min(session 数, MLFLOW_SEARCH_TRACES_MAX_THREADS)


6. 服务端入口:两个 HTTP 面

MLflow server 实际上是 FastAPI 套着 Flask:FastAPI 先注册自己的原生路由,最后把整个 Flask app 用 WSGI 中间件挂在 / 上(mlflow/server/fastapi_app.py:254)。 顺序很关键,源码注释专门标了「必须在 include_router 之后,否则 Flask 会吃掉所有请求」(:189)。

FastAPI app (mlflow/server/fastapi_app.py:148 create_fastapi_app)
├── otel_router POST /v1/traces ← 标准 OTLP 入口
├── job_api_router
├── gateway_router
├── assistant_router
└── mount("/", WSGI(flask_app)) ← 其余全部老 API

6.1 Flask 面:MLflow 自家 API

路由不是手写的,而是从 protobuf service 定义里反射生成的:get_service_endpoints() 遍历每个 rpc 方法的 databricks_pb2.rpc 扩展,取出 path 和 method (mlflow/server/handlers.py:7146)。所以 service.proto 里写的路径就是真实路径。

Handler路由落到 store 的哪个方法行号
_start_trace_v3POST /mlflow/tracesstart_tracehandlers.py:4085
_search_traces_v3POST /mlflow/traces/searchsearch_traces:3947
_get_traceGET /mlflow/traces/getget_trace:3927
_batch_get_tracesGET /mlflow/traces/batchGetbatch_get_traces:3897
_get_trace_info_v3GET /mlflow/traces/{trace_id}get_trace_info:3885
_delete_tracesPOST /mlflow/traces/delete-tracesdelete_traces:3986
_link_traces_to_runPOST /mlflow/traces/link-to-runlink_traces_to_run:4126
_calculate_trace_filter_correlationPOST /mlflow/traces/calculate-filter-correlation同名:4023
_query_trace_metricsPOST /mlflow/traces/metricsquery_trace_metrics:4280
_create_assessmentPOST /mlflow/traces/{assessment.trace_id}/assessmentscreate_assessment:4335
_update_assessmentPATCH /mlflow/traces/{trace_id}/assessments/{assessment_id}update_assessment:4370
get_trace_artifact_handlerGET(UI 专用,非 proto 生成)见下:4209

大部分 handler 都是薄的:校验 schema → protobuf 转实体 → 调 store → 实体转 protobuf。 有三处例外值得看。

_delete_traces 的 nullable 陷阱。 protobuf 的可选字段访问器在未设置时返回默认值 0, 而 max_traces=0(删全部)和 max_traces=None(不删)语义完全相反。handler 里专门定义了 _get_nullable_field,用 HasField() 判断后返回 Nonehandlers.py:4224-4228)。

_update_assessment 的 update_mask。 它不整体覆盖,而是遍历 update_mask.paths, 只把被点名的字段拼进 kwargs(:4386-4398)。支持 assessment_name / expectation / feedback / rationale / metadata / valid 六个路径。

get_trace_artifact_handler 是三级 fallback。 它是唯一直接碰 artifact 仓的 handler:

带 path 参数? ── 是 ──> 校验路径安全 → download_trace_attachment(path) → 返回单个附件
│ 否
v
_fetch_trace_data_from_store() ← 先试 store.get_trace(allow_partial=True)
│ 再试 store.batch_get_traces
│ 返回 None(表示不在 store 里)
v
看 SPANS_LOCATION 标签
├── ARCHIVE_REPO → download_archived_trace_data()
└── 其他 → download_trace_data() (traces.json)

对应 handlers.py:4477-4485allow_partial=True 的理由写在注释里:让前端能渲染还在进行中的 trace。 路径参数一定要过 validate_path_is_safe:4222),这是防目录穿越。

6.2 FastAPI 面:标准 OTLP 直接入库

这是 MLflow 观测能力向外扩张的关键一步——不装 MLflow SDK 也能把 trace 打进来。 Claude Code、Codex CLI、Gemini CLI 这些工具本身就是 OTel 客户端,把 OTEL_EXPORTER_OTLP_TRACES_ENDPOINT 指过来就行。

入口是 export_tracesmlflow/server/otel_api.py:97),路径 /v1/tracesOTLP_TRACES_PATHmlflow/tracing/utils/otlp.py:20)。流程:

  1. 必须带 header x-mlflow-experiment-id,FastAPI 用 Header(...) 声明为必填(:98)。 experiment 是唯一的落点标识,OTLP 协议本身没有这个概念,所以只能靠 header 补。
  2. Content-Type 只收 application/x-protobufapplication/json,去掉 charset 后比对(:130)。
  3. Content-Encoding 解压(:141)。
  4. JSON 编码要打一个补丁:OTLP 规范用小写 hex 表示 trace_id/span_id,而 protobuf 的 JSON 映射 对 bytes 字段要 base64。_convert_otlp_json_ids_to_base64:69)在 Parse 之前把三个字段 逐个 hex→base64 转掉,注释说不转的话下游整数转换会溢出(:150-152)。
  5. 逐 span 转成 MLflow Span,转失败直接 422。
  6. 一次 store.log_spans(experiment_id, all_spans) 全批落库。没有部分成功—— docstring 明说是 all-or-nothing,要按 trace 隔离错误请客户端自己分批(:110-112)。
  7. 带了 x-mlflow-run-id 就顺手 link_traces_to_run,失败只记日志不影响主流程(:229-233)。

一个安全设计:service.name 白名单。 resource 属性里的 service.name 会被提到根 span 上 以便 UI 显示,但只有出现在 _KNOWN_SERVICE_NAMES 冻结集合里的值才会被采纳 (otel_api.py:51,当前收了 claude-code、codex_cli_rs、codex_vscode、gemini-cli、qwen-code 五个)。 注释直说理由:防止存下不受信客户端的任意自由文本(:48-49)。

同样的防御在写库侧也有一遍——log_spans 把 resource 属性存成 trace tag 时, 跳过 telemetry.sdk.mlflow. 两个前缀,免得客户端用 resource 属性覆写 SPANS_LOCATION 这类记账标签(sqlalchemy_store.py:5377-5382)。

store 不支持 log_spans 时(比如 FileStore)返回 501 并带上 store 类名(otel_api.py:217-221)。


7. 归档与长期留存

7.1 要解决的问题

trace 明细放在 spans 表里查得快,但一个月后没人再看,还占着数据库。归档就是把「冷」trace 的 span 明细搬去便宜的对象存储,同时把 trace_info 这封信封留在关系库里——搜索结果不变,只是点开变慢

7.2 两条落法

MLflow 里有两套完全独立的归档,别混:

开源自管归档Databricks UC Delta 归档
入口服务端 YAML 配置 + 内建调度器enable_databricks_trace_archival()
代码trace_archival_config.py / trace_archival_service.py / sqlalchemy_store.archive_tracesmlflow/tracing/archival.py:7
落到哪任意 artifact repository(S3/GCS/…)Unity Catalog Delta 表
格式OTLP TracesData protobuf 单文件Delta 表 / inference table
实现在哪MLflow 本体委托给 databricks-agents

mlflow/tracing/archival.py 全文只有 70 行,两个函数都是薄转发try: from databricks.agents.archive import enable_trace_archival, 导入失败就报「请 pip install databricks-agents」(:33-36)。也就是说 UC Delta 这条路的真正逻辑 不在这个仓库里。下面讲的都是开源自管那条。

7.3 配置从哪来

配置是一个 YAML 文件,路径由环境变量 MLFLOW_TRACE_ARCHIVAL_CONFIG 指定 (trace_archival_config.py:52)。解析成一个冻结 dataclass:

trace_archival:
enabled: true
location: s3://my-bucket/mlflow-archive # 必填,且要过 repository 支持性校验
retention: 30d # 必填
long_retention_allowlist: ["1", "2"] # 这些 experiment 不受 retention 限制
interval_seconds: 300 # 默认 300,上限 86400
max_traces_per_pass: 1000 # 可选

对应 TraceArchivalServerConfig:34)和 load_trace_archival_server_config:105)。

配置带 5 秒 TTL 缓存_TRACE_ARCHIVAL_SERVER_CONFIG_CACHE_TTL_SECONDS:25), 所以改 YAML 不用重启。刷新失败时的行为很克制:如果之前有过一份有效配置, 就打 warning 并继续用旧的,不让一次手滑的编辑打断归档(:76-89)。

7.4 调度与执行

调度器是 Huey 的每分钟定时任务,注册在 mlflow/server/jobs/utils.py:807register_periodic_tasks() 里,带一把 lock_task 分布式锁防止上一轮没跑完就重入(:785)。

真正的节流在 run_trace_archival_schedulertrace_archival_service.py:72)内部: 每分钟被唤醒,但 _should_run_trace_archival_scheduler(interval_seconds) 会按配置的间隔 决定这次是否真干活,不到点就返回 0。

一次归档 pass 分两段(sqlalchemy_store.py:5792archive_traces):

_plan_trace_archival → 选出要归档的 trace(在 DB session 内,只读)
|
| 按 cutoff 分组,让同样紧急度的 experiment 共用一次候选查询
| 两类来源:① archive_now 标签(人工催)② retention 到期(常规)
v
_execute_trace_archival_plan → 逐条搬(在 session 外,慢 I/O)

单条搬运是 _archive_trace_candidate:6161),四步且每步失败都清理

  1. 读出 trace_info 和 span 行 → _load_trace_archival_data
  2. 序列化成 OTLP protobuf → spans_to_traces_data_pbmlflow/tracing/otel/otel_archival.py:53)。 格式不合法时标记 MALFORMED_TRACE 并跳过,不重试(:6174-6180)。
  3. 上传到 <location>/<experiment_id>/traces/<trace_id>/artifacts:6183-6190)。
  4. 事务性 finalize:改 SPANS_LOCATIONARCHIVE_REPO、写 ARCHIVE_LOCATION 标签、删 span 行。

任何一步失败,都会调 _delete_unreferenced_archived_trace_payload 把刚上传的孤儿文件删掉:6194:6205:6222),然后把异常统一包装成 MlflowException 交给外层, 外层只记 warning、让这条 trace 保持可重试(:5950-5957)。

归档格式为什么选 OTLP TracesData?看 traces_data_pb_to_spans 的校验就明白了 (otel_archival.py:92):它强制要求恰好一个 ResourceSpans、一个 ScopeSpans、 所有 span 属于同一个 OTLP trace id。注释说明 MLflow 的 trace 模型是扁平 span 列表, 不保留分组结构,所以归档用一个规范化的单一形状(:105-106)。 选 OTLP 而非自定义格式的好处是:这份归档文件任何 OTel 生态工具都能读。

7.5 归档与删除的交互

_delete_tracessqlalchemy_store.py:4183)也要处理归档态。它把选中的 trace 分成两类

  • DB-backed 的:在事务里 DELETE,级联清掉从表。
  • 已归档的:留到事务外第二阶段处理,因为要删对象存储上的文件,慢 (:4235-4238 的注释直说「payload cleanup can happen outside the transaction」)。

注释还提到一个并发收敛技巧:文件已经不存在也算清理成功,这样两个并发删除者不会互相报错。

还有一个贯穿始终的并发字段:trace_info.db_payload_generationmodels.py:821)。 span 写入前会 _advance_db_payload_generations_for_db_span_writes 推进这个代数 (sqlalchemy_store.py:5346),归档 finalize 时校验代数没变才提交。 读侧对应 _refresh_transitioning_trace_snapshot——如果读完 trace_info 之后、读 span 之前 归档刚好完成,就重读一次元数据,免得返回空 span(:5527-5531 的注释)。


8. 生产在线评分:让 judge 一直跑着

8.1 全貌

这是把 05 章 的 judge 搬到生产的那一步。整条链路:

每分钟 register_periodic_tasks (server/jobs/utils.py:762)
|
v
run_online_scoring_scheduler (genai/scorers/job.py:430)
| 拉出所有 active online scorer,按 experiment 分组,shuffle
| 按 trace 级 / session 级拆成两种 job 提交
+-----------------------------+
v v
run_online_trace_scorer_job run_online_session_scorer_job
(job.py:77) (job.py:111)
| |
v v
OnlineTraceScoringProcessor OnlineSessionScoringProcessor
.process_traces() .process_sessions()
|
+-- ① checkpoint 算时间窗
+-- ② search_traces 拉这段时间的 TraceInfo
+-- ③ sampler 决定哪些 trace 跑哪些 scorer
+-- ④ 线程池并发打分,写回 assessment
+-- ⑤ 推进 checkpoint

两个 job 都带 exclusive=["experiment_id"]job.py:74:107), 保证同一个 experiment 不会有两个打分 job 并发——否则 checkpoint 会乱。

8.2 断点续跑:checkpoint 存在 experiment tag 上

OnlineTraceCheckpointManageronline/trace_checkpointer.py:39)没有专用的表, checkpoint 就是一条 experiment tag,key 是 MLFLOW_LATEST_ONLINE_SCORING_TRACE_CHECKPOINT, value 是 {"timestamp_ms": ..., "trace_id": ...} 的 JSON。

为什么要带 trace_id? 因为同一毫秒可能有多条 trace。恢复时按 (timestamp_ms, trace_id) 字典序判断「已处理」:

trace_infos = [
t for t in trace_infos
if not (t.timestamp_ms == checkpoint.timestamp_ms and t.trace_id <= checkpoint.trace_id)
]

出自 online/trace_processor.py:192-199。这和 §5.4 里 order by 兜底的 timestamp_ms DESC, request_id ASC 是同一套稳定序在两头呼应。

防卡死的兜底: calculate_time_windowtrace_checkpointer.py:73)算时间窗时, 下界取 max(checkpoint, now - MAX_LOOKBACK_MS)MAX_LOOKBACK_MS 是 1 小时 (online/constants.py:7)。docstring 写明用意:防止被一批持续失败的老 trace 永远卡住—— 超过 1 小时就直接跳过它们。

还有一个细节:即使这一轮一条 trace 都没选中,也必须推进 checkpoint, 否则下一轮又扫同一个空窗口(trace_processor.py:111-116)。

8.3 采样:dense sampling

OnlineScorerSampler.sampleonline/sampler.py:57)不是「每个 scorer 独立抛硬币」, 而是条件概率瀑布

假设有两个 scorer,采样率分别是 50%(记作 A)和 25%(记作 B):

策略结果
独立抛硬币12.5% 两个都跑、37.5% 只跑 A、12.5% 只跑 B、37.5% 都不跑
dense sampling25% 两个都跑、25% 只跑 A、50% 都不跑——「只跑 B」这种组合不存在

关键差别不是被打分的 trace 变多了(其实更少),而是被选中的 trace 覆盖得更全: 凡是 B 打过分的 trace,A 一定也打过。于是两个 scorer 的分数落在同一批 trace 上, 可以横向比较(docstring 在 sampler.py:61-67 明说了这个目的)。

上游 docstring 的举例与实现不符。 sampler.py:63-66 写的是「50% 两个都跑 / 25% 只跑第一个 / 25% 都不跑」。但按代码算:B 的条件概率是 rate / prev_rate = 0.25 / 0.5 = 0.5:91), 它只在 A 已命中的那 50% 里再筛掉一半,所以真实分布是 25% / 25% / 50%。以代码为准。

实现:按 sample_rate 降序排,逐个算条件概率 rate / prev_rate,一旦被拒就 break:89-101)。随机源不是 random(),而是 sha256(f"{entity_id}:{scorer.name}") 除以 2**256——确定性哈希, 同一条 trace 重跑得到同样的结果(:94-95);哈希输入里带了 scorer 名, 所以不同 scorer 拿到的取值互相独立。

8.4 采什么、怎么采

_build_scoring_taskstrace_processor.py:150)按 filter_string 分组拉 trace, 每个不同的 filter 只发一次查询。每组的 filter 都会被强制 AND 上一条:

EXCLUDE_EVAL_RUN_TRACES_FILTER = f"metadata.{TraceMetadataKey.SOURCE_RUN} IS NULL"

出自 online/constants.py:15——排除评估 run 产生的 trace,免得离线 evaluate 的产物 又被在线 judge 打一遍。

拉数据用 OnlineTraceLoader.fetch_trace_infos_in_rangeonline/trace_loader.py:91), 它内部翻页调 search_traces,order by 固定 ["timestamp_ms ASC", "request_id ASC"], 注释解释 request_id 是必要的 tie-breaker(:132-135)。单个 job 最多 500 条 (MAX_TRACES_PER_JOBconstants.py:10)。

拿到 TraceInfo 之后才去取全量 trace,OnlineTraceLoader.fetch_tracestrace_loader.py:18) 做二级 fallback:先 batch_get_traces,没返回的再按 artifact 仓下载 (_fetch_traces_from_artifact_repo:49)。这里有个正确的负向判断: 如果一条 trace 的 SPANS_LOCATIONTRACKING_STORE 却没被 batch 返回,说明它只是还没导完, 直接跳过而不是去 artifact 仓乱找(:66-76)。

8.5 执行与写回

_execute_scoringtrace_processor.py:259)用线程池并发跑,每条 trace 一个 future, 复用评估引擎的 _compute_eval_scores_log_assessments(见 04 章)。

导入是函数内懒加载的,注释写明理由:这两个模块会拉进 pandas,模块级导入会打破 skinny client (:247-250)。

失败也要留痕:_log_error_assessments:200)会给每个 scorer 造一条带 errorFeedback 写回 trace,让失败在 assessment 历史里可见,而不是静默消失。

8.6 session 级打分

OnlineSessionScoringProcessor.process_sessionsonline/session_processor.py:96)结构对称, 只是把「找 trace」换成 tracking_store.find_completed_sessions(就是 §5.5 那个), checkpoint 存的是 (last_trace_timestamp_ms, session_id),单 job 上限 100 个 session (MAX_SESSIONS_PER_JOB)。

多打了一件事:_clean_up_old_assessments:211)——成功写入新评价后, 清掉这个 session 上一轮留下的旧在线评价,避免同一 session 每轮都堆一层。

8.7 手动触发与批次划分

除了调度器,还有一条人工路径:invoke_scorer_jobjob.py:141)。 批次划分逻辑单独抽了出来,get_trace_batches_for_scorerjob.py:400):

scorer 类型分批规则
session 级按 session_id 分组,一个 session 一批(不能拆)
单轮MLFLOW_SERVER_SCORER_INVOKE_BATCH_SIZE 定长切

invoke_scorer_job 还有一个身份透传细节:把触发者的 username 写进 MLFLOW_TRACKING_USERNAME 环境变量,让下游 gateway 请求以正确的用户身份鉴权(:165-170)。

8.8 scorer 注册存在哪

三张表(models.py):

存什么行号
scorersscorer 身份:(experiment_id, scorer_name) 唯一,scorer_id 主键:2125
scorer_versions每次注册产生一个新版本,serialized_scorer 是 JSON 文本:2166
online_scoring_configs采样率 + filter_string,每个 scorer 最多一条:2226

「最多一条」不是数据库约束,而是服务端保证的——代码里出现三次 # Each scorer has at most one online configuration, guaranteed by the servermlflow/genai/scorers/registry.py:354:257:273)。

registry.py 里的 MlflowTrackingStore:191)是 scorer 的存储门面, _hydrate_scorer:209)负责把从数据库读出的 OnlineScoringConfig 装回 Scorer 对象的 _sampling_config 字段,让 scorer.start() / .update() 这些 API 有状态可读。

Databricks 环境走另一个实现 DatabricksStore:334),配置实体是 ScorerScheduleConfigmlflow/genai/scheduled_scorers.py:12)。这个 dataclass 本身 只有四个字段(scorer / 名字 / sample_rate / filter_string),执行完全在 Databricks 侧。


9. 收尾:GenAI 这条线怎么接回 MLflow 原来的骨架

前面八节讲的都是 trace 独有的机制。但 trace 并不是一个平行宇宙——它长在 MLflow 原有的 experiment / run / logged model 骨架上。三个接点:

9.1 同一套 experiment 命名空间

SqlTraceInfo.experiment_id 是指向 experiments.experiment_id 的外键, 且关系上带 cascade="all, delete-orphan"models.py:783-786)。 注释明确说这个级联是为了让 session.delete(experiment)mlflow gc 会用) 能先把 trace 行删干净。

也就是说:删 experiment 会删掉它下面所有 trace,和删掉所有 run 是同一个语义。 artifact 路径也复用 experiment 的 artifact_location,只是多套一层 /traces 子目录(§4.2)。

9.2 trace ↔ run:两条并存的关联

历史原因导致有两种表达,_apply_trace_filter_clauses 里用 OR 把它们并起来 (sqlalchemy_store.py:3684-3706):

方式存在哪谁写的
元数据trace_request_metadata 里 key = mlflow.sourceRuntrace 在 active run 里产生时自动带上
关联边entity_associations 里 TRACE→RUN显式调 link_traces_to_run

所以 mlflow.search_traces(run_id=...) 两种都能搜到。

9.3 trace ↔ logged model:靠元数据键

TraceMetadataKey.MODEL_ID = "mlflow.modelId"mlflow/tracing/constant.py:9)。

有意思的是,SqlAlchemyStore.search_traces 的签名收了 model_id 参数, 但函数体里根本没用它sqlalchemy_store.py:3713-3796,只出现在签名和 docstring 里)。 真正的转换发生在客户端:TracingClient.search_traces 把 model_id 改写成一句普通 filter

filter_string = f"request_metadata.`mlflow.modelId` = '{model_id}'"

mlflow/tracing/client.py:352)。只有配了 MLFLOW_TRACING_SQL_WAREHOUSE_ID 的 Databricks 路径才把 model_id 原样往下传。

9.4 所以整张图长这样

experiment ────────────────┬──────────────┬────────────────┐
│ │ │ │
v v v v
runs logged models traces scorers
│ │ │ │
│ entity_associations │ metadata │ online_scoring_configs
└───────── TRACE→RUN ───┴─ modelId ────┤

┌────────────────────┼────────────────────┐
v v v
trace_tags / spans / artifact assessments
trace_metadata / archive (judge 写回)
/ trace_metrics

一句话总结这一章:GenAI 的 trace 没有另起炉灶,它复用了 experiment 这个命名空间、 复用了 artifact_location 这套存储、复用了 tag/metadata 这套键值扩展点, 只额外加了一张主表、几张从表,和两个后台循环。


10. 巧妙之处(可借鉴的技术)

妙在哪位置
同一份 token usage 存两遍:JSON 原文进 metadata 表给人读,拆成数值行进 metrics 表给 SQL 聚合——避免在查询时解析 JSONsqlalchemy_store.py:3501-3509
乐观写 + IntegrityError 补救,而不是先 SELECT 再 INSERT:并发少一次往返,冲突时才付出重读代价sqlalchemy_store.py:3527
span 条件攒成单一子查询:源码用两行注释举反例说明分开 join 会误命中,教科书级的注释sqlalchemy_store.py:9582-9597
不做 DISTINCT 并写明前提:注释论证了「右表在 join key 上唯一」,并警告改代码的人别破坏它sqlalchemy_store.py:3775-3781
确定性采样:用 sha256(trace_id + scorer_name) 代替随机数,重跑结果一致,便于排查online/sampler.py:94-95
dense sampling 的条件概率瀑布:用覆盖广度换覆盖深度,让被选中的 trace 一定被更高采样率的 scorer 全打过,于是 scorer 之间可比online/sampler.py:89-101
checkpoint 的最大回看窗:1 小时兜底,防止被一批坏 trace 永久卡住online/trace_checkpointer.py:101
归档失败即清理孤儿文件,且「文件不存在也算成功」让并发删除收敛sqlalchemy_store.py:6194:4236
service.name 白名单:来自外部 OTLP 客户端的自由文本一律不落库,除非在冻结集合里otel_api.py:51
配置热加载但失败保旧:5 秒 TTL 缓存,解析失败时继续用上一份有效配置trace_archival_config.py:76-89
路由从 protobuf 反射生成.proto 里的 path 就是唯一真相,不会和代码里手写的路由漂移handlers.py:7146

11. 边界与局限(诚实版)

  • 只有 SqlAlchemyStore 是完整的。 FileStoresearch_traces 是全量扫描 + 内存过滤 (file_store.py:2187),且没实现 get_trace / batch_get_traces / log_spans。 OTLP 端点对它会返回 501。
  • 分页是 offset 式的。 token 就是个偏移量(search_utils.py:942), 深翻页在大表上会退化,且翻页过程中有新 trace 写入会导致结果漂移。
  • query_trace_metrics 的分页没写。 源码里是一句 # TODO: Implement pagination with page_tokensqlalchemy_store.py:4142), 当前恒返回 PagedList(data_points, None)
  • span 属性搜索是字符串匹配。 span.attributes.x 落到 spans.content 这个 Text 列上做 LIKE/RLIKE,源码自己挂了 TODO 说应该把属性单独抽表(sqlalchemy_store.py:9481)。
  • OTLP 摄取无部分成功。 一批里有一个 span 转换失败,整批 422(otel_api.py:111-113)。
  • UC Delta 归档不在本仓库。 mlflow/tracing/archival.py 只是对 databricks-agents 的转发, 这条路的实现细节从这个克隆里看不出来
  • 在线打分是「最终一致」的。 一分钟一轮 + 1 小时最大回看,意味着极端情况下老 trace 会被跳过不打分。
  • search_traces 的 model_id 在 SQL 后端是空参数(§9.3),依赖客户端改写, 直接调 store 层 API 时它不生效。

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

主题文件路径符号名
存储抽象:trace 接口全集mlflow/store/tracking/abstract_store.pyAbstractStore.start_trace / search_traces / archive_traces / supports_trace_archival
写入主路径mlflow/store/tracking/sqlalchemy_store.pySqlAlchemyStore.start_tracelog_spans_get_trace_artifact_location_tag
读取与三路分派mlflow/store/tracking/sqlalchemy_store.pyget_trace_get_tracebatch_get_traces_get_spans_with_trace_info
搜索构造mlflow/store/tracking/sqlalchemy_store.py_get_filter_clauses_for_search_traces_apply_trace_filter_clauses_get_orderby_clauses_for_search_traces_get_session_scoped_trace_ids
会话搜索mlflow/store/tracking/sqlalchemy_store.pyfind_completed_sessions_build_first_trace_filter_subquery_build_session_stats_subquery
相关性分析mlflow/store/tracking/sqlalchemy_store.pycalculate_trace_filter_correlation_build_trace_filter_subquery
删除与归档mlflow/store/tracking/sqlalchemy_store.py_delete_tracesarchive_traces_plan_trace_archival_archive_trace_candidate
评价 CRUDmlflow/store/tracking/sqlalchemy_store.pycreate_assessmentupdate_assessmentdelete_assessment_get_sql_assessment
关联边mlflow/store/tracking/sqlalchemy_store.pylink_traces_to_runlink_prompts_to_trace
表定义mlflow/store/tracking/dbmodels/models.pySqlTraceInfoSqlTraceTagSqlTraceMetadataSqlTraceMetricsSqlSpanMetricsSqlSpanSqlAssessmentsSqlIssueSqlEntityAssociationSqlScorerSqlOnlineScoringConfig
过滤语法定义mlflow/utils/search_utils.pySearchTraceUtils_IDENTIFIERS_ALTERNATE_IDENTIFIERSVALID_ASSESSMENT_COMPARATORSparse_start_offset_from_page_token
DataFrame 化与字段抽取mlflow/tracing/utils/search.pytraces_to_df_FieldParser_parse_fields
客户端分页包装mlflow/tracing/fluent.pysearch_tracessearch_sessions
model_id → filter 改写mlflow/tracing/client.pyTracingClient.search_traces
CLI 查询入口mlflow/cli/traces.pysearch--experiment-id / --filter-string / --extract-fields
Flask 端点mlflow/server/handlers.py_start_trace_v3_search_traces_v3_get_trace_batch_get_traces_create_assessment_update_assessmentget_trace_artifact_handler_query_trace_metricsget_service_endpoints
OTLP 摄取mlflow/server/otel_api.pyexport_traces_convert_otlp_json_ids_to_base64_KNOWN_SERVICE_NAMES
应用装配mlflow/server/fastapi_app.pycreate_fastapi_app
归档配置mlflow/tracing/trace_archival_config.pyTraceArchivalServerConfigget_trace_archival_server_configload_trace_archival_server_config
归档调度mlflow/tracing/trace_archival_service.pyrun_trace_archival_scheduler_resolve_scheduler_trace_archival_config
归档格式mlflow/tracing/otel/otel_archival.pyspans_to_traces_data_pbtraces_data_pb_to_spans
Databricks 归档转发mlflow/tracing/archival.pyenable_databricks_trace_archivaldisable_databricks_trace_archival
在线打分入口mlflow/genai/scorers/job.pyrun_online_trace_scorer_jobrun_online_session_scorer_jobrun_online_scoring_schedulerget_trace_batches_for_scorerinvoke_scorer_job
采样与断点mlflow/genai/scorers/online/OnlineScorerSampler.sampleOnlineTraceCheckpointManager.calculate_time_windowOnlineTraceScoringProcessor.process_tracesOnlineTraceLoader.fetch_tracesOnlineSessionScoringProcessor.process_sessions
scorer 注册存储mlflow/genai/scorers/registry.pyMlflowTrackingStore.upsert_online_scoring_config_hydrate_scorerDatabricksStore
Databricks 计划配置mlflow/genai/scheduled_scorers.pyScorerScheduleConfig
定时任务注册mlflow/server/jobs/utils.pyregister_periodic_tasks
关键常量mlflow/tracing/constant.pyTraceMetadataKeyTraceTagKeySpansLocationTokenUsageKey
store 选择mlflow/tracking/_tracking_service/utils.py_register_tracking_stores

继续读: 想知道这些 span 一开始怎么来的,看 01 章02 章;想知道 assessment 里那些分数怎么算出来的, 看 04 章05 章