跳到主要内容

数据截至 (上游 commit 36c7a7f6eca6)

追踪运行时:span 从产生到落盘的完整管道

这章讲什么: 01 章讲了 @mlflow.trace 怎么把一次函数调用变成一棵 span 树——那是数据长什么样。这一章讲这些对象在运行期怎么被搬走:谁负责造 span、谁决定留不留、半成品 trace 存在哪、什么时候打包、走哪条线程、最后送到哪个后端。服务端收到之后怎么落库、怎么被搜,是 06 章的事。


1. 这条管道要解决的四个麻烦

MLflow 的追踪底座是 OpenTelemetry(简称 OTel,业界标准的分布式追踪 SDK)。但"直接用 OTel"会立刻撞上四个麻烦,管道的每一层基本都是在化解其中之一。

麻烦具体表现这一层的答案
抢全局单例OTel 的 TracerProvider 是进程级单例。PromptFlow、Snowpark、FastAPI 自动埋点都可能已经设过它自己维护一个隔离的 provider(§3)
span 是散的OTel 只按 span 逐个回调,没有"一棵 trace"的概念InMemoryTraceManager 攒树(§5)
写后端很慢每个 span 都同步发一次 HTTP,会把用户的推理链路拖垮processor + exporter 两级异步(§6、§7)
目的地不止一个本地实验、Databricks UC 表、model serving 推理表、OTLP collector按配置组装 processor 列表(§8)

2. 全景:七道关卡

先看整条流水线。从左到右是一个 span 的生命周期,中间那条向下的支线是"半成品暂存"。

@mlflow.trace ① ② ③ ④ ⑤
用户代码 ──► tracer ──► 采样判定 ──► span processor ──► exporter ──► 目的地
start_span() provider 留 or 丢 on_start/on_end 两条路径 实验 / UC 表
(惰性、隔离) │ / 推理表 / OTLP

┌───────────┴───────────┐
▼ ▼
InMemoryTraceManager AsyncTraceExportQueue
(半成品 trace 的暂存处) (后台线程池,真正发请求)

部件与职责一览:

部件干什么文件
_TracerProviderWrapper决定用 MLflow 私有 provider 还是 OTel 全局单例mlflow/tracing/provider.py:121
_MlflowSampler按比例决定一条 trace 留不留,支持按调用点覆盖mlflow/tracing/sampling.py:14
InMemoryTraceManager单例聚合中枢,存"还没结束的 trace"mlflow/tracing/trace_manager.py:57
BaseMlflowSpanProcessorspan 开始/结束时的钩子,把 span 挂进 tracemlflow/tracing/processor/base_mlflow.py:176
MlflowV3SpanExporter把 span 和整棵 trace 发给 tracking servermlflow/tracing/export/mlflow_v3.py:77
AsyncTraceExportQueue有界队列 + 线程池,异步执行发送任务mlflow/tracing/export/async_export_queue.py:38
UserTraceDestinationRegistry目的地的优先级解析(context-local → 全局 → 环境变量)mlflow/tracing/destination.py:24

主线走一遍(不进代码): 用户调 @mlflow.trace 的函数 → 拿 tracer(此时才初始化整条管道)→ 采样器说"留" → OTel 造 span 并回调 on_start,processor 顺手在 trace manager 里建 trace/挂 span → 函数返回,span end() → 回调 on_end,processor 把 span 交给 exporter → exporter 分两路:单个 span 增量上报、根 span 结束时 pop_trace 打包整棵树 → 丢进异步队列,后台线程发 HTTP。


3. 关卡①:tracer provider 的隔离与惰性初始化

3.1 为什么坚持不污染全局 provider

这不是我的推断,模块 docstring 直说了(mlflow/tracing/provider.py:1-8):

每个追踪操作必须走这个模块而不是直接用 OTel API,因为 MLflow 需要控制 provider 的初始化,并保证不干扰同一进程里为别的目的使用 OpenTelemetry 的库(例如 PromptFlow、Snowpark)。

OTel 的 trace.set_tracer_provider() 是"一次性、全局、先到先得"的。如果 MLflow 抢了它,用户原本发往自家 collector 的业务 trace 就全跑到 MLflow 实验里去了;反过来如果别人先抢了,MLflow 就一条 trace 也收不到。

3.2 两种模式,一个开关

MLFLOW_USE_DEFAULT_TRACER_PROVIDER 默认 Truemlflow/environment_variables.py:973)。名字有点绕——它为真时用的是 MLflow 自己的隔离 provider,为假时才去用 OTel 全局单例。

隔离模式(默认,env=true)统一模式(env=false)
provider 存在哪self._isolated_tracer_provider 私有字段OTel 的 trace._TRACER_PROVIDER
context 用哪个MLflow 私有的 mlflow_runtime_contextOTel 的 context_api
别的库已设过 provider互不影响把 MLflow 的 processor 加进已有 provider
适合谁大多数人想让 MLflow span 和原生 OTel span 拼成同一棵树的人

两种模式的切换点集中在一个门面类里(_TracerProviderWrapperprovider.py:121),get / set / once 三个成员各自按开关分叉:

# mlflow/tracing/provider.py:146-158(节选)
def get(self) -> TracerProvider:
if MLFLOW_USE_DEFAULT_TRACER_PROVIDER.get():
return self._isolated_tracer_provider
return trace.get_tracer_provider()

注意 set() 在统一模式下是直接写 trace._TRACER_PROVIDER = tracer_providerprovider.py:191)——故意绕过 OTel 自己的 once 标志,否则 set_destination() 之后的更新会被 OTel 静默忽略。

统一模式还有一段"不覆盖、只追加"的逻辑:如果检测到已有一个 SDK 的 TracerProvider(比如 Uvicorn/FastAPI 自动埋点设的),MLflow 只把自己的 processor add_span_processor 进去,然后直接 return,不换掉人家的 provider(_initialize_tracer_providerprovider.py:656-674)。

3.3 惰性初始化:Once 与两个标志位

管道不在 import mlflow 时初始化,而是等第一次真的要 tracer 时才建。入口是 _get_tracerprovider.get_or_init_tracerprovider.py:193-195):

def get_or_init_tracer(self, module_name: str) -> trace.Tracer:
self.once.do_once(_initialize_tracer_provider)
return self.get().get_tracer(module_name)

Once 是从 OTel 抄过来的小工具(mlflow/tracing/utils/once.py:6):先无锁读 _done 快速返回,未完成才拿锁、双重检查、执行、置位。"执行一次"和"阻塞其他并发调用者直到执行完"两件事一起给。

为什么要两个 Once 实例(_isolated_tracer_provider_once_global_provider_init_onceprovider.py:135-140)?代码注释给了理由:统一模式下不能借用 OTel 的 _TRACER_PROVIDER_SET_ONCE,否则外部库先设过全局 provider 时,MLflow 会以为"已经初始化过了",自己的 processor 一个都装不上。

3.4 隔离的随机数发生器

OTel 默认的 RandomIdGenerator 用的是 Python 全局 random 模块。用户代码里一句 random.seed(42)(跑实验时很常见),就会让重跑产生一模一样的 trace/span ID,后端撞主键。

_IsolatedRandomIdGeneratorprovider.py:90-118)的做法是持有一个私有 random.Random() 实例,与全局状态无关。配套还有一个容易被忽略的细节:

# mlflow/tracing/provider.py:75-85(节选)
if hasattr(os, "register_at_fork"):
def _reseed():
for r in tuple(_private_random_generators):
r.seed()
os.register_at_fork(after_in_child=_reseed)

fork 出来的子进程会继承父进程的随机数状态,多 worker 的服务(Gunicorn 一类)会集体生成同一串 ID。这里在 fork 后钩子里给每个私有实例重新播种。tuple(...) 那一层是为了在迭代时给集合拍快照,避免并发修改。

这个能力默认关着MLFLOW_TRACE_USE_ISOLATED_RANDOM_ID_GENERATOR 默认 Falseenvironment_variables.py:1006),撞 ID 了才开。


4. 关卡②:采样

4.1 它要解决的小问题

生产流量大的时候,不是每条请求都值得存一份完整 trace。要能说"只留 10%",还要能说"这个函数特别重要,它的 trace 全留"。

4.2 ParentBased 包一层 _MlflowSampler

采样器在 provider 初始化时装进 TracerProvider_get_trace_samplerprovider.py:733-754):

# mlflow/tracing/provider.py:635
return ParentBased(root=_MlflowSampler(sampling_ratio))

ParentBased 是 OTel 自带的组合器,它把"整条 trace 要不要留"的决定只做一次

span 类型决定者结果
根 span(无 parent)_MlflowSampler按比例掷骰子
子 span,父被采样ParentBased一律留
子 span,父被丢弃ParentBased一律丢

这保证不会出现"一棵树只留下一半枝干"的残缺 trace。

一个值得注意的事实:MLFLOW_TRACE_SAMPLING_RATIO 默认值是 1.0 而不是 Noneenvironment_variables.py:953),所以 sampling_ratio is not None 恒成立——默认情况下采样器也是装上的,只是比例为 1(全留)。返回 None(用 OTel 默认采样)只发生在用户把比例设成 0~1 之外的非法值时(provider.py:747-752)。这一点很关键,因为下面的按调用点覆盖要靠这个采样器在位才生效。

4.3 按调用点覆盖比例

_MlflowSampler.should_sample 每次都先看一个 ContextVar(sampling.py:37-42):

override = _SAMPLING_RATIO_OVERRIDE.get()
if override is not None:
sampler = TraceIdRatioBased(override)
return sampler.should_sample(...)
return self._default_sampler.should_sample(...)

写这个 ContextVar 的是 fluent 层的一个上下文管理器(_set_sampling_ratio_overridemlflow/tracing/fluent.py:74-82):

# 示意,非源码:@mlflow.trace 支持逐函数指定采样率
@mlflow.trace(sampling_ratio_override=1.0) # 全局设了 0.01,这个函数仍然全采
def critical_step(x):
return model(x)

装饰器包装函数时把它套在最外层——同步/异步/生成器四种包装路径都套了(fluent.py:352fluent.py:361fluent.py:457fluent.py:484),所以无论被装饰的是什么形态的函数,进函数体前 ContextVar 都已就位。合法性校验(必须在 0~1)在装饰器入口做(fluent.py:241-244)。

用 ContextVar 而不是全局变量的好处:并发的协程/线程各有各的覆盖值,互不串味。


5. 关卡③:聚合中枢 InMemoryTraceManager

5.1 它要解决的小问题

OTel 只会一个一个地回调你"某 span 开始了""某 span 结束了"。但 MLflow 要发给后端的是一整棵 trace(带 TraceInfo、token 用量汇总、请求/响应预览)。中间这段"还没凑齐"的状态得有地方放。

5.2 两张表 + 一个双向映射

InMemoryTraceManager 是懒汉双重检查的单例(trace_manager.py:65-71),内部只有两个容器(trace_manager.py:73-79):

otel_trace_id (128 位 int,OTel 造的)

│ _otel_id_to_mlflow_trace_id:普通 dict

"tr-<32 位 hex>" (MLflow trace_id)

│ _traces:带超时的缓存(见 5.3)

_Trace ──► info : TraceInfo(状态、时间、metadata、tags)
├─► span_dict : {span_id: LiveSpan} ← 用 dict 而非列表,便于按 id 取
└─► prompts : [PromptVersion]

为什么要两套 ID?OTel 内部只认整数 trace ID,MLflow 对外的 ID 是 "tr-" + hexgenerate_mlflow_trace_id_from_otel_trace_idmlflow/tracing/utils/__init__.py:539-549)。管道里两种 ID 都会出现:processor 的 on_start 拿到的是 OTel span,exporter 要 pop 的也是按 OTel trace ID,而用户和后端看到的全是 tr-xxx。所以正反两个方向都得能查。

核心 API 一张表:

方法何时被调干什么
register_trace根 span 开始(trace_manager.py:81_Trace,登记双向映射
register_span每个 span 开始(trace_manager.py:102LiveSpan 塞进 span_dict
get_trace到处(trace_manager.py:138上下文管理器,持锁 yield,保证读改写原子
get_mlflow_trace_id_from_otel_idexporter(trace_manager.py:171OTel ID → MLflow ID
has_open_spansexporter(trace_manager.py:177树里还有没有没结束的 span
pop_trace根 span 导出时(trace_manager.py:195弹出并冻结成不可变 Trace

pop_trace 里那句 to_mlflow_trace()可变 → 不可变的分水岭trace_manager.py:29-36):每个 LiveSpanto_immutable_span() 变成 Span,之后再改也影响不到已导出的数据。

注意 _Trace 只是内部表示,注释说得很直白:用 dict[str, Span] 而不是 TraceData 就是为了能按 span_id 随机访问(trace_manager.py:19-20)。

5.3 带超时的缓存:防止未结束的 trace 泄漏内存

如果一个根 span 永远不 end() 会怎样? 它会一直躺在 _traces 里。一个卡死的 agent 循环 + 高 QPS,内存就慢慢涨没了。

get_trace_cache_with_timeout()mlflow/tracing/utils/timeout.py:28-49)按配置返回两种缓存:

条件返回行为
设了 MLFLOW_TRACE_TIMEOUT_SECONDSMlflowTraceTimeoutCache后台线程定期扫,超时的 trace 主动结束并以 ERROR 状态导出
没设(默认 Nonecachetools.TTLCache过期就静默丢弃,ttl 默认 3600 秒、maxsize 默认 1000

第一种是这里的精华。MlflowTraceTimeoutCachetimeout.py:129)自己维护一条按过期时间排序的链表__setitem__ 时把新 key 挂到链表尾部并记 expirestimeout.py:154-169)。扫描时从头走,遇到第一个没过期的就立刻返回——O(过期条数) 而不是 O(缓存大小)_get_expired_tracestimeout.py:228-245)。

过期处理不是简单删除,而是"给它体面收尾"(expiretimeout.py:200-226):

root_span.set_status(SpanStatusCode.ERROR)
root_span.add_event(SpanEvent.from_exception(MlflowTracingException(msg)))
root_span.end() # 调 end() 会触发正常的导出流程

于是卡死的 trace 也能在 UI 里看到,还带一条"超时了,可以调大 MLFLOW_TRACE_TIMEOUT_SECONDS"的说明。后面那句 if request_id in self: del self[request_id] 是兜底:正常情况下 end() 触发的导出会 pop_trace 把它带走,万一出错了就强删。

超时值改了怎么办?链表是按旧超时构建的,没法就地调整。_check_timeout_updatetrace_manager.py:212-229)在每次 register_trace 时检查,值变了就整个换一个新缓存,并明确警告"这个操作会丢弃当前所有进行中的 trace"。诚实但粗暴。


6. 关卡④:processor 钩子

6.1 它站在什么位置

OTel 的 SpanProcessor 是官方留给你的两个钩子:span 开始时 on_start、结束时 on_end。MLflow 的所有 processor 都继承自 BaseMlflowSpanProcessorprocessor/base_mlflow.py:176),它同时继承 OtelMetricsMixinSimpleSpanProcessor——后者意味着默认路径是"结束一个就同步 export 一个"

6.2 on_start:建树

# mlflow/tracing/processor/base_mlflow.py:186-201(节选逻辑)
trace_id = self._trace_manager.get_mlflow_trace_id_from_otel_id(span.context.trace_id)
if not trace_id and span.parent is not None:
return # 非根 span 却查不到 trace → 多半已超时被清掉
if span.parent is None:
trace_info = self._start_trace(span) # 子类实现:造 TraceInfo 并 register_trace
trace_id = trace_info.trace_id
InMemoryTraceManager.get_instance().register_span(create_mlflow_span(span, trace_id))

三件事,顺序不能乱:先认 trace(根 span 就现建一个)、再把 OTel span 包成 LiveSpan、最后挂进 tracecreate_mlflow_spanLiveSpan.__init__ 会把 MLflow trace_id 写进 span 属性 SpanAttributeKey.REQUEST_IDmlflow/entities/span.py:676)——这一步是后面 on_end 能反查 trace 的前提。

_start_trace 是留给子类的抽象方法(base_mlflow.py:238-239),各变体的差别几乎全在这里。

6.3 on_end:计数器 + 临界区 + 分流

on_end 本身只做一件事:给一个"在飞的 on_end"计数器加减base_mlflow.py:241-256)。

with self._pending_on_end_condition:
self._pending_on_end_count += 1
try:
self._on_end_impl(span)
finally:
with self._pending_on_end_condition:
self._pending_on_end_count -= 1
if self._pending_on_end_count == 0:
self._pending_on_end_condition.notify_all()

这个计数器是为 flush 服务的。没有它会有一个竞态:用户调 mlflow.flush_trace_async_logging() 时,某个 span 正走在 on_end 半路上、还没进 BatchSpanProcessor 的队列,flush 信号先发给了后台线程 → 这个 span 被漏掉。所以 flush_all_batch_processors 的第一步就是等所有 processor 的计数器归零(base_mlflow.py:87-100):

flush_all_batch_processors()

├─ 0. 对每个 processor:wait_for(pending_on_end_count == 0) ← 等在飞的 on_end 落地
├─ 1. processor.force_flush() span 队列 ──► exporter.export()
├─ 2. exporter._async_queue.flush() 发送任务队列 ──► tracking store
└─ 3. terminate 时才 shutdown + 把 _batch_delegate 置 None

注释还解释了为什么这样不会死锁:wait_for 总是先求值谓词再阻塞,所以即使 notify_all 在进入 wait_for 之前就发生了(计数已经是 0),也会立刻返回。

真正的活在 _on_end_implbase_mlflow.py:258-283)里,两段:

第一段——临界区。_deduplication_lock,再用 get_trace 上下文管理器持 trace manager 的锁,若是根 span 就回填 trace info(结束时间、状态、token 用量、cost、输入输出预览,见 _update_trace_infobase_mlflow.py:344-384),并立刻更新"最近活跃 trace ID",这样即使走批量模式 mlflow.get_trace() 也能拿到对的 ID。

关于这把锁要说句实话:它的名字(_deduplication_lockbase_mlflow.py:200-203)说的是"span 名去重",但在本 commit 里它保护的临界区只做 trace info 的读改写;SDK 侧唯一一处给重名 span 加数字后缀的代码在 TraceData.intermediate_outputsmlflow/entities/trace_data.py:55-65),并不在这把锁里。LiveSpan 记了 _original_name 说是"供 span 日志阶段去重用"(mlflow/entities/span.py:697-702),但真正的后缀逻辑不在 SDK 侧。名字比实现走得快了一步。

第二段——分流。

if self._batch_delegate is not None and not maybe_get_request_id(is_evaluate=True):
self._batch_delegate.on_end(span) # 进 BatchSpanProcessor 队列,后台线程批量发
else:
super().on_end(span) # SimpleSpanProcessor:当场同步 export

maybe_get_request_id(is_evaluate=True)mlflow/tracing/utils/__init__.py:480-494)判断"当前是否跑在 MLflow 评估里"。评估期强制绕过批量模式,因为评估引擎打完分要立刻按 trace 读回结果,等 5 秒一批的调度就乱套了。同样的判断在 exporter 的 _should_log_async 里又出现一次(export/mlflow_v3.py:401-407)——两级流水线都要在评估期退化成同步。

批量委托只在两个条件同时满足时才建(provider.py:920):

use_batch_processor=MLFLOW_USE_BATCH_SPAN_PROCESSOR.get() and exporter._is_async_enabled

它是标准的 OTel BatchSpanProcessor,参数来自 MLflow 的环境变量(_create_batch_span_processorbase_mlflow.py:164-173):批大小默认 5、调度间隔默认 5000 毫秒,队列上限取 max(批大小, 2048)——因为 OTel 规定 max_export_batch_size <= max_queue_size,否则直接抛 ValueError

6.4 那个 WeakSet 注册表

_batch_processor_registry 是个 weakref.WeakSetbase_mlflow.py:64)。注释说明了它存在的理由:set_destination() 会造一个全新的 tracer provider,旧 processor 被孤立,但它的 BatchSpanProcessor 后台线程还活着、队列里还压着 span。 用注册表才能把这些"孤儿"一起 flush 掉。用弱引用则保证被 GC 的 processor 自动出列,不会无限膨胀。

6.5 四个变体

Processor用在哪与基类的关键差别文件
MlflowV3SpanProcessor默认(本地/远程 tracking server)_start_trace 生成 tr- 前缀 ID + 实验位置processor/mlflow_v3.py:15
DatabricksUCTableSpanProcessor目的地是 Unity CatalogID 用 v4 格式、schema 版本写 4、根 span 结束时把 user/session 回填成 span 属性processor/uc_table.py:23
InferenceTableSpanProcessorDatabricks model serving不继承基类,直接继承 SimpleSpanProcessor;trace ID 关联 Databricks request ID,流式场景从 Flask 请求头兜底取processor/inference_table.py:48
OtelSpanProcessor配了 OTLP 端点继承 BatchSpanProcessor;双写模式下不注册 trace(交给 MLflow processor),把用户 tag 摊成 span 属性以便过 OTLPprocessor/otel.py:20

OtelSpanProcessor 那个 tag 摊平很有意思(processor/otel.py:70-79):OTLP 协议里没有"trace 级 tag"这种东西,所以非 mlflow. 前缀的 tag 被逐个写成 SpanAttributeKey.TRACE_TAG_PREFIX + key 的 span 属性,服务端再还原成 tag 行。


7. 关卡⑤:导出

7.1 一次 export,两条路径

MlflowV3SpanExporter.exportexport/mlflow_v3.py:103-115)只有四行,但分出两条完全不同的路:

exporter.export(spans)

├─► _export_spans_incrementally(spans) 每个 span 单独上报(增量、能在 UI 里边跑边看)
│ └─ 按 experiment_id 分组 → 异步/同步 _log_spans()

└─► _export_traces(spans) 只处理根 span:pop 整棵树,一次性 start_trace
└─ 异步/同步 _log_trace()
_export_spans_incrementally_export_traces
处理哪些 span全部只有 span._parent is None 的根 span
后端接口log_spans(experiment_id, spans)start_trace(trace_info) + 上传 trace data
何时跳过tracking URI 是 databricks,或后端不支持树里还有没结束的 span(见 7.2)
失败后果_store_supports_log_spans 置 False,之后不再试打 warning,trace 丢失

增量路径的降级很克制(_log_spansexport/mlflow_v3.py:267-301):后端抛 NotImplementedError(旧 FileStore)或返回 501 时,把 _store_supports_log_spans 永久置 False,之后 export() 第一行的 if 就直接跳过这条路,不会每个 span 刷一条警告。

分组时的 experiment_id 来源有个坑,注释专门写了(_collect_mlflow_spans_for_exportexport/mlflow_v3.py:147-190):不能现场调 get_experiment_id_for_trace(),因为那玩意读的是线程本地 ContextVar,而这段代码可能跑在 BatchSpanProcessor 的 worker 线程里,ContextVar 是空的。所以要从 on_start 时(在原线程)已写进 TraceInfo 的值里取。

7.2 延后导出:后台线程的 span 还没跑完

这是整个 exporter 里最不显然的一段(_export_tracesexport/mlflow_v3.py:191-225)。

问题场景: 根函数返回了(根 span 结束),但它 spawn 出去的后台线程还在跑,还会产生子 span。如果这时就 pop_trace,OTel→MLflow 的 ID 映射被删掉,后来那些迟到的 span 在 _collect_mlflow_spans_for_export 里查不到 trace,直接被丢。

解法: 根 span 到达时先问 manager.has_open_spans()trace_manager.py:177-185,遍历 span_dict 看有没有 end_time_ns is None 的),有就把根 span 押在 _deferred_root_spans 里不导出;等下一批 span 进来时再回头检查、补导出。

根 span 结束

├─ has_open_spans(trace_id)? ──是──► 押进 _deferred_root_spans,本轮不导
│ │
│ 下一次 export 进来时回头查 ──► 都结束了 → _do_export_trace
└─ 否 ──► _do_export_trace(pop_trace → 打包 → 送队列)

锁的顺序也被注释点名了:先在 _deferred_lock拷一份 key 列表,然后出锁再调 has_open_spans——因为后者要拿 InMemoryTraceManager._lock,两把锁嵌套着拿会有死锁风险(export/mlflow_v3.py:200-210)。

shutdown() 时会把还押着的根 span 无条件导出(export/mlflow_v3.py:409-417),避免"那个后台 span 永远不结束"导致整条 trace 静默消失。

7.3 AsyncTraceExportQueue:有界队列 + 背压

_should_log_async() 说要异步时,导出动作被包成 Task 丢进队列(export/mlflow_v3.py:256-263)。这个队列的设计有三个刻意选择:

其一:满了就丢,绝不阻塞。

# mlflow/tracing/export/async_export_queue.py:68-71
try:
# Do not block if the queue is full, it will block the main application
self._queue.put(task, block=False)
except queue_Full:
... # 30 秒最多警告一次

追踪是旁路功能,宁可丢 trace 也不能拖慢用户的推理。队列上限默认 1000(MLFLOW_ASYNC_TRACE_LOGGING_MAX_QUEUE_SIZE)。警告做了 30 秒节流,免得刷屏。

其二:消费者线程手动限流线程池。 _dispatch_taskasync_export_queue.py:92-124)在提交前先看在飞任务数:

if len(self._active_tasks) >= self._max_workers:
_, self._active_tasks = wait(self._active_tasks, return_when=FIRST_COMPLETED)

注释讲得很清楚:ThreadPoolExecutor 内部队列是无界的,如果消费者线程只管往里塞,self._queue 那个 maxsize 就形同虚设,内存照涨。所以只在有空闲 worker 时才取下一个任务,把积压留在有界队列里。默认 10 个 worker。

其三:三条退路。

情况处理位置
队列已被 flush(terminate=True) 关停当场同步执行任务(否则 _stop_event 永不清除,会死等)async_export_queue.py:56-61
提交线程池失败(解释器正在关闭)在当前线程直接执行async_export_queue.py:117-124
进程退出atexit 回调里 flush(terminate=True)async_export_queue.py:52152-164

flush(terminate=False) 还会在排空后把线程重新拉起来async_export_queue.py:184-187),所以用户可以在脚本中间随时调 mlflow.flush_trace_async_logging()mlflow/tracking/fluent.py:1003)而不影响后续追踪。

7.4 SpanBatcher:给 UC 表用的攒批器

注意区分:§6.3 的 BatchSpanProcessor 攒的是 OTel span,这里的 SpanBatcherexport/span_batcher.py:18)攒的是已经转好的 MLflow span,只被 UC 表 exporter 用(export/uc_table.py:27-31)。

两个触发条件,谁先到算谁(_worker_loopspan_batcher.py:69-80):

worker 线程循环:
_worker_awaken.wait(max_interval)

├─ 返回 True → 是被 add_span 叫醒的(攒够了) → _consume_batch(flush_all=False)
└─ 返回 False → 超时了(间隔到) → _consume_batch(flush_all=True)

_consume_batch 出队后还要按 location 二次分组span_batcher.py:92-100),因为队列里可能混着发往不同 UC 表的 span,而后端接口一次只收一个位置。

批大小 ≤ 1 时不起线程、直接透传(span_batcher.py:57-59span_batcher.py:41),省掉一整条后台线程。


8. 关卡⑥:送到哪里

8.1 目的地的三级优先级

UserTraceDestinationRegistry.get()destination.py:29-35)的顺序写得很短:

context-local 值(set_destination(context_local=True) 设的,按协程/线程隔离)
↓ 没有
全局值(set_destination() 设的)
↓ 没有
MLFLOW_TRACING_DESTINATION 环境变量

环境变量的解析用了 match 模式(destination.py:52-80):一段是实验 ID、两段是 catalog.schema、三段的 UC 表前缀明确拒绝并告诉你该用 set_experiment(trace_location=...)

set_destination() 设完值会立刻重建整条管道_initialize_tracer_provider()provider.py:542)——注意它绕过了 Once,所以是无条件重建。这也是 §6.4 那个 WeakSet 注册表存在的原因。

8.2 processor 列表怎么组装

_get_span_processorsprovider.py:794-904)是整个模块最需要照着读的一段。从上往下是优先级,命中并 return 即停:

_get_span_processors(disabled)

├─ disabled=True ─────────────────────────────► [] → 装 NoOpTracerProvider

├─ 有用户目的地?(含"活跃实验绑了 UC 位置"的自动解析)
│ ├─ UC 表 / UC schema → DatabricksUCTableSpanProcessor
│ └─ MLflow 实验 → MlflowV3SpanProcessor
│ └─ 除非「OTLP 已配 且 开了双写」,否则 ──────────────► return

├─ OTLP 端点已配置? → OtelSpanProcessor
│ └─ 没开双写,或列表里已有别的 processor ──────────► return

└─ 兜底:在 Databricks model serving 里 → InferenceTableSpanProcessor
否则 → MlflowV3SpanProcessor(tracking_uri)

几个容易踩空的点:

  • 返回空列表不是"报错",而是装 NoOpTracerProviderprovider.py:632-635),全链路静默失效。
  • 没设显式目的地时,会去查活跃实验是不是绑在 UC 位置上(_resolve_experiment_uc_locationprovider.py:757-791),查到就当成用户目的地写回注册表。
  • _get_span_processor()(单数,provider.py:586-612)在多 processor 场景下显式扫 BaseMlflowSpanProcessor 类型,而不是取索引 0——双写模式下 OTLP processor 可能排在前面。

8.3 OTLP 双写

should_use_otlp_exporter()utils/otlp.py:38-42)要求两件事同时成立:OTel 标准环境变量给了端点, MLFLOW_ENABLE_OTLP_EXPORTER 为真(默认真)。端点解析遵守 OTel 规范:OTEL_EXPORTER_OTLP_TRACES_ENDPOINT 原样用,OTEL_EXPORTER_OTLP_ENDPOINT 要拼上 /v1/tracesutils/otlp.py:89-104)。

MLFLOW_TRACE_ENABLE_OTLP_DUAL_EXPORT 默认 False——默认情况下配了 OTLP 就是"替代"MLflow 导出,不是"追加"。开了双写才两边都发。

双写打开时有两处防重复:

防什么怎么防位置
指标重复上报OTLP processor 的 export_metrics 强制置 False,只让 MLflow processor 报provider.py:870-872
trace 重复注册OtelSpanProcessor._should_register_traces 置 False,不碰 trace managerprocessor/otel.py:46

9. 关卡⑦:开关与跨进程

9.1 disable / enable / trace_disabled

三个 API 共用一套机制(provider.py:925 / 829 / 869):

API做什么关键细节
disable()disabled=True 重建 provider(得到 NoOpTracerProvider),再手动把 once._done 置 True_done 是为了阻止下次 _get_tracer 又把它初始化回来
enable()正常重建 + 置 _done已启用且 _done 为真时只打一行 info
trace_disabled装饰器,进函数前 disable()finallyenable()内部用,防止 MLflow 自己的调用(如 judge 请求)产生 trace

is_tracing_enabled()provider.py:1089-1106)的判断有个反直觉的分支:if not provider.once._done: return True——还没初始化过就算"启用",因为默认状态就是开的。

trace_disabled 的错误处理值得学(provider.py:1026-1067):它只捕获 MlflowTracingException(来自 disable/enable 本身),用 is_func_called 标志确保即使开关操作炸了,被装饰的原函数也一定会被执行一次、且只执行一次。追踪出问题绝不能连累业务逻辑。

reset()provider.py:1070-1086)比 disable() 更彻底:装 NoOp → 把 once._done 翻回 False(这样下次操作会重新初始化)→ 清用户目的地 → 重置配置。

9.2 跨进程:W3C traceparent 传播

场景:客户端起一个 span,通过 HTTP 调远端 agent,希望远端产生的 span 挂进同一棵 trace

标准做法是 W3C Trace Context:把 trace 上下文序列化成一个 traceparent 请求头。MLflow 给了配对的两个 API(mlflow/tracing/distributed/__init__.py):

# 示意,非源码:客户端
with mlflow.start_span("client-root"):
headers = get_tracing_context_headers_for_http_request() # {"traceparent": "00-<trace>-<span>-01"}
requests.post(url, headers=headers)

# 示意,非源码:服务端
with set_tracing_context_from_http_request_headers(dict(request.headers)):
with mlflow.start_span("server-handler"): # 自动挂到客户端那棵树上
...

注入端(distributed/__init__.py:16-68)直接用 OTel 官方的 TraceContextTextMapPropagator().inject,context 从 get_current_context() 取——这是隔离模式必须的,因为 MLflow 的 span 活在私有 runtime context 里,不在 OTel 全局 context 里(provider.py:242-250)。

提取端(distributed/__init__.py:89-190)比想象中多做了两件事:

  1. 大小写补丁。 Flask 会把请求头首字母大写成 Traceparent,而 OTel 的 propagator 只认全小写,所以先手动改键(distributed/__init__.py:146-151)。这种细节不踩一次是想不到的。
  2. 注册一个"占位 trace"。 远端进程的 trace manager 里并没有这棵树,而 processor 的 on_start 查不到 trace 就会丢弃非根 span(§6.2)。所以这里先用 is_remote_trace=True 建一个空壳 TraceInfo 顶上(distributed/__init__.py:175-183)。

第 2 点还有个反例保护:如果本地已经有一棵同 ID 的真 trace(同进程分布式 + 异步导出还没落地),就不注册占位,否则会把真数据盖掉(注释在 distributed/__init__.py:169-173)。退出时也只清理"确实由自己注册的"占位(registered_dummy_trace 标志,distributed/__init__.py:189-190)。

远端产生的 trace 需要后端支持才能入库;不支持时导出侧会给一条明确的升级提示而不是静默失败(export/mlflow_v3.py:233-238)。


10. 巧妙之处(可以抄的)

① 惰性 + Once + 双标志位。 "永不在 import 时做重活"是 SDK 的基本礼貌;用两个独立的 Once 区分隔离/统一模式,让"外部库已抢走全局 provider"这种情况不至于让自己彻底装不上钩子(provider.py:135-146)。

② 私有随机源 + fork 重播种。 一句 random.seed(42) 或一次 fork 就能让 ID 撞车,而这类 bug 极难复现。把随机源私有化并挂 os.register_at_fork 是根治(provider.py:77-118)。

③ 有界队列 + 手动限流线程池。 只给队列设上限而不管线程池的内部队列,等于没设上限——这个观察写在注释里(async_export_queue.py:96-101),值得任何写异步上报的人抄一遍。

④ "在飞计数器"和 flush 握手。 两级流水线的 flush 必须先等 on_end 全部落地再排空队列,否则永远有一个尾巴漏掉(base_mlflow.py:87-100)。

⑤ 按过期时间排序的链表。 超时扫描做成 O(过期条数),让"每秒扫一次"这种默认配置也不心疼(timeout.py:228-245)。

⑥ 超时不是丢弃,是标 ERROR 后正常导出。 用户最需要看到的恰恰是卡死的那条 trace(timeout.py:200-226)。

⑦ 评估期两级同时退化成同步。 processor 和 exporter 各自独立判断 maybe_get_request_id(is_evaluate=True)base_mlflow.py:280export/mlflow_v3.py:404),任一级漏判都会让评估拿不到 trace。


11. 边界与坑

  • 默认丢,不默认阻塞。 队列满、后端挂、超时——一律丢 trace 保业务。要"不丢"就得自己调大 MLFLOW_ASYNC_TRACE_LOGGING_MAX_QUEUE_SIZE 或关掉异步。
  • 改超时会丢掉所有在途 trace。 _check_timeout_update 直接换缓存实例,只给一条 warning(trace_manager.py:212-229)。
  • _traces 默认上限 1000 条、TTL 3600 秒。 超出按缓存策略淘汰,静默。
  • 统一模式下的隔离 ID 生成器不生效。 环境变量注释明说:全局模式下若检测到已有 provider,MLFLOW_TRACE_USE_ISOLATED_RANDOM_ID_GENERATOR 无效(environment_variables.py:1001-1003)。
  • MlflowTraceTimeoutCache 不是线程安全的,类 docstring 自己声明"只能在单例上下文里用"(timeout.py:129-134);线程安全靠外面 InMemoryTraceManager._lock 兜。
  • 批量导出对多数 exporter 无效。 环境变量注释写明:目前只有 UC 表 exporter 支持 span 攒批,其他 exporter 立即导出(environment_variables.py:1238-1239)。
  • Databricks tracking URI 下跳过增量 span 上报export/mlflow_v3.py:124-129),只走整棵 trace 那条路。
  • Databricks notebook 里异步导出默认关,除非显式设环境变量(export/mlflow_v3.py:386-399)。

12. 代码地图

主题文件符号名
provider 门面与双模式mlflow/tracing/provider.py_TracerProviderWrapper / get_or_init_tracer / get_context_api
惰性初始化一次mlflow/tracing/utils/once.pyOnce.do_once
provider 组装mlflow/tracing/provider.py_initialize_tracer_provider / _get_span_processors / _get_mlflow_span_processor
隔离随机 IDmlflow/tracing/provider.py_IsolatedRandomIdGenerator / _reseed
采样器mlflow/tracing/sampling.py_MlflowSampler / _SAMPLING_RATIO_OVERRIDE
按调用点覆盖采样率mlflow/tracing/fluent.py_set_sampling_ratio_override
聚合中枢mlflow/tracing/trace_manager.pyInMemoryTraceManager / register_trace / pop_trace / has_open_spans
超时缓存mlflow/tracing/utils/timeout.pyget_trace_cache_with_timeout / MlflowTraceTimeoutCache / _get_expired_traces
processor 基类mlflow/tracing/processor/base_mlflow.pyBaseMlflowSpanProcessor / on_start / _on_end_impl / flush_all_batch_processors
processor 变体mlflow/tracing/processor/MlflowV3SpanProcessor / DatabricksUCTableSpanProcessor / InferenceTableSpanProcessor / OtelSpanProcessor
指标上报 mixinmlflow/tracing/processor/otel_metrics_mixin.pyOtelMetricsMixin.record_metrics_for_span
主 exportermlflow/tracing/export/mlflow_v3.pyMlflowV3SpanExporter / _export_spans_incrementally / _export_traces / _should_log_async
异步队列mlflow/tracing/export/async_export_queue.pyAsyncTraceExportQueue / Task / _dispatch_task / flush
span 攒批mlflow/tracing/export/span_batcher.pySpanBatcher / _worker_loop / _consume_batch
目的地解析mlflow/tracing/destination.pyUserTraceDestinationRegistry / _get_trace_location_from_env
OTLP 配置mlflow/tracing/utils/otlp.pyshould_use_otlp_exporter / get_otlp_exporter / _get_otlp_traces_endpoint
开关mlflow/tracing/provider.pydisable / enable / trace_disabled / reset / is_tracing_enabled
跨进程传播mlflow/tracing/distributed/__init__.pyget_tracing_context_headers_for_http_request / set_tracing_context_from_http_request_headers
公开 flush 入口mlflow/tracking/fluent.pyflush_trace_async_logging

接着读: span 里的属性和 LiveSpan 结构见 01 章;第三方 SDK 怎么被自动接进这条管道见 03 章;评估期为什么要把这条管道强制退化成同步见 04 章;这些 trace 到了服务端怎么存、怎么被搜、怎么被持续打分见 06 章