从 MLflow 追踪构建延迟与错误监控

一条 trace 可以记录聊天请求经历了哪些函数和模型调用。要把这些信息变成可用的监控,还需要明确用户与会话归属、导出方式、采样比例、查询窗口和统计口径。否则,图上出现了 p95 或错误率,也未必能说明真实流量的健康状态。

本文依据 MLflow 官方 Cookbook Production Observability with MLflow Tracing完整整理,原页日期为 2026 年 3 月 18 日,未见个人作者署名,维护方为 MLflow Project。本文于 2026 年 10 月 5 日读取全文并核对最新在线 API;涉及的代码修订会明确说明。未安装软件、调用模型、发送请求或创建监控任务。

聊天请求产生带标签的trace,经采样与异步队列写入MLflow,再按时间窗查询,计算已完成请求的延迟、错误和令牌统计;采样与丢弃会影响统计覆盖。
原创技术示意图,依据原文与 API 文档绘制;不包含实测指标或虚构仪表盘。

先建立最小追踪链路

原文依赖 mlflow 和 openai,并使用已在本机 http://127.0.0.1:5000 运行的跟踪服务。安装示例如下;原文未固定包版本,实际项目应记录并锁定验证过的版本,而不是把任意最新组合当成已验证环境。

python -m pip install mlflow openai

OpenAI 客户端需要有效凭据,模型调用可能收费。凭据应通过受控环境变量或部署的密钥机制提供,不要写入源代码、trace 标签或日志。下面把原文固定模型名改成从 OPENAI_MODEL 读取,避免让文章替读者决定可用模型和费用。

编辑修订:异步和采样环境变量移到 MLflow 导入与初始化之前。先用 100% 采样验证链路;之后才按需要调整生产采样。跟踪服务地址保留为回环地址示例,不代表可将无认证的服务暴露到公网。

import os

os.environ["MLFLOW_ENABLE_ASYNC_TRACE_LOGGING"] = "true"
os.environ["MLFLOW_ASYNC_TRACE_LOGGING_MAX_WORKERS"] = "10"
os.environ["MLFLOW_ASYNC_TRACE_LOGGING_MAX_QUEUE_SIZE"] = "1000"
os.environ["MLFLOW_TRACE_SAMPLING_RATIO"] = "1.0"

import mlflow
import openai

mlflow.set_tracking_uri("http://127.0.0.1:5000")
experiment = mlflow.set_experiment("production-chatbot-demo")
mlflow.openai.autolog()
client = openai.OpenAI()
model = os.environ["OPENAI_MODEL"]

mlflow.openai.autolog() 记录受支持的 OpenAI 调用,外层再用 @mlflow.trace 包装应用函数。这样一条请求能同时看到应用入口与模型调用,而不只是一段孤立的 API 日志。

给请求补上用户、会话与部署上下文

原文先展示基本聊天函数,再加入用户与会话字段。下面将两步合并为最终函数:tags 保存环境与版本,metadata 保存请求创建时确定的用户和会话。mlflow.trace.user 与 mlflow.trace.session 是用于界面分组的保留键;tags 可以后续修改,metadata 设计用于固定信息。

@mlflow.trace
def support_chatbot(
    user_message: str,
    user_id: str,
    session_id: str,
    conversation_history: list[dict] | None = None,
) -> str:
    mlflow.update_current_trace(
        tags={
            "environment": "production",
            "app_version": "2.1.0",
            "request_type": "support",
        },
        metadata={
            "mlflow.trace.user": user_id,
            "mlflow.trace.session": session_id,
        },
    )
    messages = [{
        "role": "system",
        "content": "You are a helpful customer support agent. "
                   "Be concise and actionable.",
    }]
    # 编辑补充:历史应由服务端可信存储提供,不允许覆盖系统角色。
    for item in conversation_history or []:
        role = item.get("role")
        content = item.get("content")
        if role not in {"user", "assistant"} or not isinstance(content, str):
            raise ValueError("Invalid conversation history")
        messages.append({"role": role, "content": content})
    messages.append({"role": "user", "content": user_message})
    response = client.chat.completions.create(
        model=model, messages=messages
    )
    return response.choices[0].message.content or ""

上面的生产标签只是沿用原文查询条件,实验名已明确为演示专用。实际应用应由部署配置设置环境和版本,用户标识宜使用内部匿名化标识。自动追踪可能记录消息内容与历史,因此仅把 user_id 改成匿名编号还不够:还应控制输入、输出的敏感信息、访问权限及保留范围。

原文直接把历史字典追加到消息列表。若历史来自不可信客户端,这会允许其伪造高权限角色。本稿增加角色与类型检查;真正的会话历史仍应从服务端可信记录读取,不能把客户端自报的 assistant 消息当成可信历史。普通用户文本中的提示注入也不会因此自动消失。

import uuid

# 会调用模型,可能产生费用;本文未执行。
answer = support_chatbot(
    user_message="How do I reset my password?",
    user_id="user-demo-0001",
    session_id=uuid.uuid4().hex,
)
print(answer)

执行者可在 MLflow 界面中确认应用 span 与模型调用 span 是否同属一条 trace。若脚本马上转入查询,应使用后文的 flush=True 等待当前进程尚未写出的异步 trace;不能把短暂查不到结果直接当成埋点失败。

异步导出减少等待,但并非零开销或无损保证

原文说明 OSS MLflow 默认同步记录,异步模式把导出放到后台线程。工作线程数和队列长度用于控制并发与缓冲能力,原文示例分别为 10 和 1000。当队列在高负载下满了,新 trace 会被丢弃并产生警告,而非让应用一直等待。

原文还说程序正常退出时会自动 flush,长时间运行的服务持续在后台导出,通常无需每次请求后手动刷新。这不能推出进程被强杀、机器崩溃或队列过载时也绝不丢失。采集 span、序列化、排队和后台发送仍需要资源,不能沿用原文结语中“响应时间完全不受追踪开销影响”的绝对表述。

采样按请求类型管理,统计也要保留同样的分层

若验证完成后需要减少记录量,可在下次初始化前把全局 MLFLOW_TRACE_SAMPLING_RATIO 改为 "0.1",表示采样约 10% 的 trace。原文还给出函数级覆盖:账单咨询路径设为 1.0,FAQ 设为 0.05。

@mlflow.trace(sampling_ratio_override=1.0)
def process_billing_request(user_id: str, action: str):
    mlflow.update_current_trace(
        tags={"request_type": "billing"},
        metadata={"mlflow.trace.user": user_id},
    )
    response = client.chat.completions.create(
        model=model,
        messages=[
            {"role": "system", "content": "Handle billing inquiries."},
            {"role": "user", "content": action},
        ],
    )
    return response.choices[0].message.content or ""

@mlflow.trace(sampling_ratio_override=0.05)
def handle_faq(question: str):
    mlflow.update_current_trace(tags={"request_type": "faq"})
    response = client.chat.completions.create(
        model=model,
        messages=[
            {"role": "system", "content": "Answer common questions briefly."},
            {"role": "user", "content": question},
        ],
    )
    return response.choices[0].message.content or ""

这些函数处理的是咨询示例,并未实现任何实际账单操作。只有标上 environment=production 的请求会进入后面的生产筛选;若这些路径也要计入同一个监控,需在真实服务中一致地补齐环境、版本与会话字段。

10%、100%、5% 三种采样率混在一个结果集中,业务类型占比就会变化。直接计算总体错误率或分位数,会把被重点采集的类型放大。即便单一随机采样可用于估计,仍有采样误差和丢弃风险;采样 trace 的令牌求和也只是已记录样本的总量,不是全流量用量或账单。用于 SLO 的口径需要完整计数或经过论证的估计方法。

生成演示流量,再按明确范围检索

原文用 20 个模拟用户、8 个问题和 30 次调用填充数据。下面保留这一步,并把原来的静默 except Exception: pass 改为计数失败,避免把全部调用失败误认成一批正常样本。仍然不输出潜在含敏感信息的异常详情。

import random

user_ids = [f"user-{i:04d}" for i in range(20)]
questions = [
    "How do I reset my password?", "Can I export my data?",
    "What's the API rate limit?", "How do I add team members?",
    "My integration isn't working.", "How do I cancel my subscription?",
    "Where are the API docs?", "How do I enable SSO?",
]
failed_calls = 0
for _ in range(30):
    try:
        support_chatbot(
            user_message=random.choice(questions),
            user_id=random.choice(user_ids),
            session_id=uuid.uuid4().hex,
        )
    except Exception:
        failed_calls += 1
print("Failed demo calls:", failed_calls)

每次循环生成新会话 ID,表示独立请求演示,并不模拟连续多轮会话。异常能被追踪装饰器捕获并标记,但被采样掉或导出丢失的请求仍可能查不到。本文没有执行这 30 次模型调用。

查询可按 tags、metadata、状态和时间过滤。以下函数把范围固定到本次演示实验,并显式刷新当前进程异步写入。最新在线 API 推荐 locations 替代已弃用的 experiment_ids;原文过滤语法中的 trace.status 与 trace.execution_time_ms 仍按查询文档保留。

import time

def search_demo(filter_string=None):
    return mlflow.search_traces(
        locations=[experiment.experiment_id],
        filter_string=filter_string,
        return_type="list",
        flush=True,
    )

all_traces = search_demo()
prod_traces = search_demo("tag.environment = 'production'")
error_traces = search_demo("trace.status = 'ERROR'")

one_hour_ago = int((time.time() - 3600) * 1000)
recent_traces = search_demo(f"trace.timestamp_ms > {one_hour_ago}")

user_traces = search_demo(
    "metadata.`mlflow.trace.user` = 'user-0001'"
)
slow_errors = search_demo(
    "tag.environment = 'production' "
    "AND trace.status = 'ERROR' "
    "AND trace.execution_time_ms > 5000"
)

for name, rows in [
    ("all", all_traces), ("production", prod_traces),
    ("errors", error_traces), ("last_hour", recent_traces),
    ("user_0001", user_traces), ("slow_errors", slow_errors),
]:
    print(name, len(rows))

这些条件中的字符串都是固定示例,时间值由数值计算而来,没有把不可信用户输入直接拼接进过滤表达式。若开发成公共查询接口,应验证可查询字段和值,不要直接插入任意文本。大量数据应使用底层客户端分页或分窗口聚合;把所有历史 trace 一次装进内存只适合有限规模演示。

明确状态类型、分位数方法和令牌缺失

最新 TraceInfo API 将 request_metadata 和 status 标为弃用别名,建议分别用 trace_metadata 与 state;state 是 TraceState 枚举。持续时间可用 execution_duration,单位毫秒。以下代码据此修订,而不是把枚举直接当作字符串比较。

import json
import math
from collections import defaultdict
from mlflow.entities import TraceState

def nearest_rank(values, q):
    ordered = sorted(values)
    if not ordered:
        return None
    return ordered[max(0, math.ceil(q * len(ordered)) - 1)]

def summarize_traces(traces):
    completed = [
        t for t in traces
        if t.info.state in (TraceState.OK, TraceState.ERROR)
    ]
    latencies = [
        t.info.execution_duration for t in completed
        if t.info.execution_duration is not None
    ]
    errors = sum(t.info.state == TraceState.ERROR for t in completed)
    total = len(completed)
    per_user = defaultdict(list)
    input_tokens = output_tokens = 0
    token_records = missing_or_invalid_token_records = 0

    for t in completed:
        meta = t.info.trace_metadata
        uid = meta.get("mlflow.trace.user", "unknown")
        if t.info.execution_duration is not None:
            per_user[uid].append(t.info.execution_duration)
        raw = meta.get("mlflow.trace.tokenUsage")
        try:
            usage = json.loads(raw) if raw else None
            if not isinstance(usage, dict):
                raise ValueError("Token usage absent")
            inputs = usage["input_tokens"]
            outputs = usage["output_tokens"]
            if (type(inputs) is not int or type(outputs) is not int
                    or inputs < 0 or outputs < 0):
                raise ValueError("Invalid token counts")
        except (ValueError, TypeError, KeyError):
            missing_or_invalid_token_records += 1
            continue
        input_tokens += inputs
        output_tokens += outputs
        token_records += 1

    return {
        "queried": len(traces),
        "completed": total,
        "errors": errors,
        "error_rate_pct": 100 * errors / total if total else None,
        "latency_count": len(latencies),
        "p50_ms": nearest_rank(latencies, 0.50),
        "p95_ms": nearest_rank(latencies, 0.95),
        "p99_ms": nearest_rank(latencies, 0.99),
        "recorded_input_tokens": input_tokens,
        "recorded_output_tokens": output_tokens,
        "recorded_total_tokens": input_tokens + output_tokens,
        "token_records": token_records,
        "missing_or_invalid_token_records": missing_or_invalid_token_records,
        "per_user_mean_ms": {
            uid: {"mean": sum(v) / len(v), "count": len(v)}
            for uid, v in sorted(per_user.items())
        },
    }

summary = summarize_traces(recent_traces)
print(summary)

分位数采用最近秩法,索引为 ceil(q*n)-1;这是本稿明确选择的方法,与原文直接使用 int(n*q) 的索引有区别。不同统计工具可能采用插值定义,应在仪表盘说明方法。只有 30 条甚至采样后更少的 trace 时,p99 接近样本最大值,不足以稳定刻画真实尾延迟。

错误率分母只计入已完成的 OK 或 ERROR trace,运行中和未指定状态另留在 queried 与 completed 的差异里。令牌字段缺失、无效或不受集成支持时,不把缺失悄悄解释为零用量;已记录总数应连同覆盖条数一起报告。每用户平均延迟也保留样本数,避免一个只有一次请求的平均值被误读为稳定表现。

把统计封装成一次健康检查

下面保留原文“最近时间窗、错误率阈值、p95 阈值”的监控流程,明确窗口上下界,并复用前面的统计口径。它只打印提示,不连接任何通知服务,也没有被本文安排定时执行。

def check_production_health(
    lookback_minutes=30,
    error_rate_threshold=5.0,
    p95_latency_threshold_ms=5000,
):
    if lookback_minutes <= 0:
        raise ValueError("lookback_minutes must be positive")
    now_ms = int(time.time() * 1000)
    cutoff_ms = now_ms - int(lookback_minutes * 60 * 1000)
    traces = search_demo(
        f"trace.timestamp_ms > {cutoff_ms} "
        f"AND trace.timestamp_ms <= {now_ms} "
        "AND tag.environment = 'production'"
    )
    result = summarize_traces(traces)
    if result["completed"] == 0:
        print("No completed traces; check traffic and export coverage.")
        return result
    if result["error_rate_pct"] > error_rate_threshold:
        print("ALERT: observed error rate exceeds threshold")
    if (result["p95_ms"] is not None
            and result["p95_ms"] > p95_latency_threshold_ms):
        print("ALERT: observed p95 exceeds threshold")
    print(result)
    return result

# 只会查询和打印;未接通知系统。本文未执行。
check_production_health(
    lookback_minutes=60,
    error_rate_threshold=5.0,
    p95_latency_threshold_ms=5000,
)

“窗口中没有已完成 trace”不等于服务健康:可能没有流量,也可能是采样、导出故障或权限问题。真实运行时,阈值应按业务类型、样本量和基线设置,还需决定告警去重、恢复通知与路由;原文提到可以用 cron、Airflow 等调度,这只是下一步集成方向。

原文展示的 30 条 trace、3.3% 错误率、3400 毫秒 p95 及令牌总量都是示例输出,本文没有把它们当成实测结果。完成本流程后,得到的是可查询的追踪上下文和一次指标检查的实现基础;完整生产监控仍需验证采样覆盖、数据权限、导出可靠性和通知交付。

继续核对接口时可参考 Search Traces 与 search_traces API。原文还指向 RAG 评估与自定义 LLM 评审教程,它们处理质量评估,与本文的延迟和错误监控是不同任务。

原始版权与许可

MLflow 项目的 Apache License 2.0 与原始版权声明保留如下;正文已标明中文整理和代码修订。 官方许可来源

Copyright 2018 Databricks, Inc.  All rights reserved.

				Apache License
                           Version 2.0, January 2004
                        http://www.apache.org/licenses/

   TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION

   1. Definitions.

      "License" shall mean the terms and conditions for use, reproduction,
      and distribution as defined by Sections 1 through 9 of this document.

      "Licensor" shall mean the copyright owner or entity authorized by
      the copyright owner that is granting the License.

      "Legal Entity" shall mean the union of the acting entity and all
      other entities that control, are controlled by, or are under common
      control with that entity. For the purposes of this definition,
      "control" means (i) the power, direct or indirect, to cause the
      direction or management of such entity, whether by contract or
      otherwise, or (ii) ownership of fifty percent (50%) or more of the
      outstanding shares, or (iii) beneficial ownership of such entity.

      "You" (or "Your") shall mean an individual or Legal Entity
      exercising permissions granted by this License.

      "Source" form shall mean the preferred form for making modifications,
      including but not limited to software source code, documentation
      source, and configuration files.

      "Object" form shall mean any form resulting from mechanical
      transformation or translation of a Source form, including but
      not limited to compiled object code, generated documentation,
      and conversions to other media types.

      "Work" shall mean the work of authorship, whether in Source or
      Object form, made available under the License, as indicated by a
      copyright notice that is included in or attached to the work
      (an example is provided in the Appendix below).

      "Derivative Works" shall mean any work, whether in Source or Object
      form, that is based on (or derived from) the Work and for which the
      editorial revisions, annotations, elaborations, or other modifications
      represent, as a whole, an original work of authorship. For the purposes
      of this License, Derivative Works shall not include works that remain
      separable from, or merely link (or bind by name) to the interfaces of,
      the Work and Derivative Works thereof.

      "Contribution" shall mean any work of authorship, including
      the original version of the Work and any modifications or additions
      to that Work or Derivative Works thereof, that is intentionally
      submitted to Licensor for inclusion in the Work by the copyright owner
      or by an individual or Legal Entity authorized to submit on behalf of
      the copyright owner. For the purposes of this definition, "submitted"
      means any form of electronic, verbal, or written communication sent
      to the Licensor or its representatives, including but not limited to
      communication on electronic mailing lists, source code control systems,
      and issue tracking systems that are managed by, or on behalf of, the
      Licensor for the purpose of discussing and improving the Work, but
      excluding communication that is conspicuously marked or otherwise
      designated in writing by the copyright owner as "Not a Contribution."

      "Contributor" shall mean Licensor and any individual or Legal Entity
      on behalf of whom a Contribution has been received by Licensor and
      subsequently incorporated within the Work.

   2. Grant of Copyright License. Subject to the terms and conditions of
      this License, each Contributor hereby grants to You a perpetual,
      worldwide, non-exclusive, no-charge, royalty-free, irrevocable
      copyright license to reproduce, prepare Derivative Works of,
      publicly display, publicly perform, sublicense, and distribute the
      Work and such Derivative Works in Source or Object form.

   3. Grant of Patent License. Subject to the terms and conditions of
      this License, each Contributor hereby grants to You a perpetual,
      worldwide, non-exclusive, no-charge, royalty-free, irrevocable
      (except as stated in this section) patent license to make, have made,
      use, offer to sell, sell, import, and otherwise transfer the Work,
      where such license applies only to those patent claims licensable
      by such Contributor that are necessarily infringed by their
      Contribution(s) alone or by combination of their Contribution(s)
      with the Work to which such Contribution(s) was submitted. If You
      institute patent litigation against any entity (including a
      cross-claim or counterclaim in a lawsuit) alleging that the Work
      or a Contribution incorporated within the Work constitutes direct
      or contributory patent infringement, then any patent licenses
      granted to You under this License for that Work shall terminate
      as of the date such litigation is filed.

   4. Redistribution. You may reproduce and distribute copies of the
      Work or Derivative Works thereof in any medium, with or without
      modifications, and in Source or Object form, provided that You
      meet the following conditions:

      (a) You must give any other recipients of the Work or
          Derivative Works a copy of this License; and

      (b) You must cause any modified files to carry prominent notices
          stating that You changed the files; and

      (c) You must retain, in the Source form of any Derivative Works
          that You distribute, all copyright, patent, trademark, and
          attribution notices from the Source form of the Work,
          excluding those notices that do not pertain to any part of
          the Derivative Works; and

      (d) If the Work includes a "NOTICE" text file as part of its
          distribution, then any Derivative Works that You distribute must
          include a readable copy of the attribution notices contained
          within such NOTICE file, excluding those notices that do not
          pertain to any part of the Derivative Works, in at least one
          of the following places: within a NOTICE text file distributed
          as part of the Derivative Works; within the Source form or
          documentation, if provided along with the Derivative Works; or,
          within a display generated by the Derivative Works, if and
          wherever such third-party notices normally appear. The contents
          of the NOTICE file are for informational purposes only and
          do not modify the License. You may add Your own attribution
          notices within Derivative Works that You distribute, alongside
          or as an addendum to the NOTICE text from the Work, provided
          that such additional attribution notices cannot be construed
          as modifying the License.

      You may add Your own copyright statement to Your modifications and
      may provide additional or different license terms and conditions
      for use, reproduction, or distribution of Your modifications, or
      for any such Derivative Works as a whole, provided Your use,
      reproduction, and distribution of the Work otherwise complies with
      the conditions stated in this License.

   5. Submission of Contributions. Unless You explicitly state otherwise,
      any Contribution intentionally submitted for inclusion in the Work
      by You to the Licensor shall be under the terms and conditions of
      this License, without any additional terms or conditions.
      Notwithstanding the above, nothing herein shall supersede or modify
      the terms of any separate license agreement you may have executed
      with Licensor regarding such Contributions.

   6. Trademarks. This License does not grant permission to use the trade
      names, trademarks, service marks, or product names of the Licensor,
      except as required for reasonable and customary use in describing the
      origin of the Work and reproducing the content of the NOTICE file.

   7. Disclaimer of Warranty. Unless required by applicable law or
      agreed to in writing, Licensor provides the Work (and each
      Contributor provides its Contributions) on an "AS IS" BASIS,
      WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or
      implied, including, without limitation, any warranties or conditions
      of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A
      PARTICULAR PURPOSE. You are solely responsible for determining the
      appropriateness of using or redistributing the Work and assume any
      risks associated with Your exercise of permissions under this License.

   8. Limitation of Liability. In no event and under no legal theory,
      whether in tort (including negligence), contract, or otherwise,
      unless required by applicable law (such as deliberate and grossly
      negligent acts) or agreed to in writing, shall any Contributor be
      liable to You for damages, including any direct, indirect, special,
      incidental, or consequential damages of any character arising as a
      result of this License or out of the use or inability to use the
      Work (including but not limited to damages for loss of goodwill,
      work stoppage, computer failure or malfunction, or any and all
      other commercial damages or losses), even if such Contributor
      has been advised of the possibility of such damages.

   9. Accepting Warranty or Additional Liability. While redistributing
      the Work or Derivative Works thereof, You may choose to offer,
      and charge a fee for, acceptance of support, warranty, indemnity,
      or other liability obligations and/or rights consistent with this
      License. However, in accepting such obligations, You may act only
      on Your own behalf and on Your sole responsibility, not on behalf
      of any other Contributor, and only if You agree to indemnify,
      defend, and hold each Contributor harmless for any liability
      incurred by, or claims asserted against, such Contributor by reason
      of your accepting any such warranty or additional liability.

   END OF TERMS AND CONDITIONS
   APPENDIX: How to apply the Apache License to your work.

      To apply the Apache License to your work, attach the following
      boilerplate notice, with the fields enclosed by brackets "[]"
      replaced with your own identifying information. (Don't include
      the brackets!)  The text should be enclosed in the appropriate
      comment syntax for the file format. We also recommend that a
      file or class name and description of purpose be included on the
      same "printed page" as the copyright notice for easier
      identification within third-party archives.

   Copyright [yyyy] [name of copyright owner]

   Licensed under the Apache License, Version 2.0 (the "License");
   you may not use this file except in compliance with the License.
   You may obtain a copy of the License at

       http://www.apache.org/licenses/LICENSE-2.0

   Unless required by applicable law or agreed to in writing, software
   distributed under the License is distributed on an "AS IS" BASIS,
   WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
   See the License for the specific language governing permissions and
   limitations under the License.
© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容