链路追踪日志

一、一个抓狂的下午

教务处那边反馈说"偶尔回复特别慢,等好几秒才出结果"。

我打开日志,看到的是这个:

[supervisor] 开始
[supervisor] 结束, 耗时1.2s
[route] 开始
[route] 结束, 耗时5.1s        ← 慢在这里
[aggregator] 开始
[aggregator] 结束, 耗时0.3s

没了。就这些。route agent 那 5 秒到底花在哪——是 LLM 生成慢?检索慢了?网络抖了一下?完全不知道。

这就是只有"节点级日志"的后果。知道哪个 Agent 慢,但不知道为什么慢。要修都没方向。

我后来写了一套 Span 级别的追踪,把问题从"永远找不到"变成了"一条 SQL 10 秒定位"。下面是把整个过程拆开来讲。

二、先从数据结构开始

不管什么追踪系统,核心就两个东西:Trace 和 Span。

Trace 是一次完整的用户请求。Span 是请求里的一个操作——调了一次 LLM、执行了一次工具、跑了一个 Agent 节点。Span 之间可以嵌套,父子关系自动形成一棵树。

from dataclasses import dataclass, field
import uuid, time

@dataclass
class Span:
    span_id: str
    parent_id: str | None       # null 说明这是顶层 Span
    name: str                    # 比如 "llm_call" / "tool_exec" / "agent_node"
    agent: str                   # 属于哪个 Agent: supervisor / route / knowledge_graph
    start_time: float
    end_time: float | None = None
    status: str = "running"

    # 每种类型的 Span 存不同的 metadata
    # LLM Span 存: {"model": "qwen-max", "input_tokens": 1200, "output_tokens": 350}
    # Tool Span 存: {"tool_name": "vector_search", "args": {...}, "result_count": 5}
    # Node Span 存: {"node_name": "route_agent", "input_keys": ["destination"]}
    metadata: dict = field(default_factory=dict)


@dataclass
class Trace:
    trace_id: str
    session_id: str
    user_input: str
    spans: list[Span]
    created_at: float

数据模型不复杂——两个 dataclass 就够。关键是把 Span 的 metadata 设计成字典而不是固定字段——因为 LLM 调用和工具执行的元数据完全不同,固定字段套不下。

这么做的好处是查数据的时候可以按维度筛选:想看所有 LLM 调用按耗时排序、想看某个 Agent 下面的所有 tool 执行、想看过去 24 小时哪种错误最多——全是 SQL 的 WHERE 条件。

三、怎么让嵌套的调用自动关联

这是整个追踪系统里最核心的问题。假设一次请求的执行路径是这样的:

supervisor 调 LLM(第 1 个 Span)
  ├── route agent 调 LLM(第 2 个 Span,父 Span 是 supervisor)
  │     └── 向量检索 tool(第 3 个 Span,父 Span 是 route)
  ├── knowledge_graph agent 调 Neo4j(第 4 个 Span)
  └── aggregator 调 LLM(第 5 个 Span)

第 2 个 Span 怎么知道自己的父 Span 是 supervisor 而不是 knowledge_graph?Agent 是并行跑的,在 ThreadPoolExecutor 的不同线程里,怎么保证各自的 Span 不串?

答案是用 ContextVar——Python 标准库自带的上下文变量,线程安全,不需要传参,自动在当前线程内传递:

from contextvars import ContextVar

# ContextVar 是线程/协程安全的——每个线程有自己独立的副本
_current_trace: ContextVar = ContextVar("trace")  # 当前请求的 Trace
_current_span: ContextVar = ContextVar("span")    # 当前正在写入的 Span

ContextVar 的原理(这个理解很重要):你把它想象成一个线程本地存储的字典。线程 A 调用 _current_span.set(x),只有线程 A 能看到 x;线程 B 同时调用 _current_span.set(y),看到的是 y。两个线程互不干扰。这就是为什么并行 Agent 不会把 Span 写串。

有了 ContextVar,TraceContext 的实现就很简单了:

class TraceContext:
    """每个请求创建一个 TraceContext,管整个请求生命周期"""

    def __init__(self, session_id: str, user_input: str):
        self.trace = Trace(
            trace_id=uuid.uuid4().hex[:12],   # 12 位随机 hex,够用且短
            session_id=session_id,
            user_input=user_input,
            spans=[],
            created_at=time.time(),
        )
        # 把这个 TraceContext 设为当前线程的"活跃 Trace"
        # _token 保存旧的 ContextVar 值,close 时要恢复
        self._token = _current_trace.set(self)

    def start_span(self, name: str, agent: str, span_type: str,
                   metadata: dict = None) -> Span:
        # 关键:从 ContextVar 获取当前的活跃 Span 作为父 Span
        # 如果当前没有活跃 Span,说明这是顶层 Span,parent_id 为 None
        parent = _current_span.get(None)

        span = Span(
            span_id=uuid.uuid4().hex[:8],
            parent_id=parent.span_id if parent else None,
            name=name,
            agent=agent,
            start_time=time.time(),
            metadata=metadata or {},
        )
        span.metadata["type"] = span_type
        self.trace.spans.append(span)
        return span

    def close(self):
        # 恢复 ContextVar 的旧值——防止内存泄漏
        _current_trace.reset(self._token)

start_span 这个方法每次被调用时,自动从 _current_span 取出当前的活跃 Span 作为父 Span——不需要手动传参,不需要管当前在第几层嵌套。

接下来是一个 with 语句的上下文管理器,让埋点只需要两行代码:

from contextlib import contextmanager

@contextmanager
def traced_span(name: str, agent: str, span_type: str, **metadata):
    """包一层 with,自动计时 + 异常捕获 + 父子关联"""
    ctx = _current_trace.get()
    if ctx is None:
        # 没有活跃的 TraceContext(比如在测试环境),直接裸跑函数
        yield None
        return

    span = ctx.start_span(name, agent, span_type, metadata)
    token = _current_span.set(span)   # 把当前 Span 推到"活跃 Span"栈

    try:
        yield span                     # 执行被追踪的代码
        span.status = "success"
    except Exception as e:
        span.status = "error"
        span.metadata["error"] = str(e)
        raise                          # 日志要记,异常也要继续抛
    finally:
        span.end_time = time.time()
        _current_span.reset(token)     # 恢复之前的活跃 Span

用法是这样的:

with traced_span("llm_call", "supervisor", "llm",
                 model="qwen-max") as span:
    response = llm.invoke(messages)
    # span 对象可以继续往里塞 metadata
    span.metadata["input_tokens"] = response.usage.input_tokens
    span.metadata["output_tokens"] = response.usage.output_tokens

两行——一行 with,一行捞结果塞 metadata。对业务代码的侵入就这么多。

四、在真实 Agent 节点里怎么埋

有了 traced_span,埋点就是在关键调用外面套一层 with。以 supervisor 节点为例,我把所有有意义的操作都包了 Span:

def supervisor_node(state: AgentState) -> dict:
    ctx = _current_trace.get()

    # 整个 supervisor 节点包在最外层
    with traced_span("agent_node", "supervisor", "node",
                     node="supervisor") as node_span:

        # 组装上下文(这一步很快,不需要 Span)
        context = build_context(state)

        # LLM 调用 —— 这是最需要追踪的环节
        with traced_span("llm_call", "supervisor", "llm",
                         model="qwen-max") as llm_span:
            response = llm.invoke(
                [SystemMessage(content=SUPERVISOR_PROMPT + context)]
                + state["messages"]
            )
            # response.usage 里有 token 数,顺手记下来
            llm_span.metadata.update({
                "input_tokens": response.usage.input_tokens,
                "output_tokens": response.usage.output_tokens,
                "prompt_preview": context[:100],   # 截前 100 字,debug 用
            })

        # 解析 PLAN —— 这是一个纯逻辑操作,也记
        with traced_span("parse_plan", "supervisor", "logic") as parse_span:
            plan = extract_plan(response.content)
            parse_span.metadata["plan"] = plan

        node_span.metadata["agents_to_call"] = plan.get("agents", [])

    return build_state_update(response, plan)

三层 Span:最外层是 agent_node 代表整个 supervisor,里面 llm_call 代表 LLM 调用(最耗时、最花钱的环节),parse_plan 代表解析逻辑。查数据的时候可以一层层下钻——先看 agent_node 的耗时,如果慢,点进去看是 llm_call 慢还是 parse_plan 慢。

再看一个有工具调用的 Agent——知识图谱节点:

def knowledge_graph_agent(state: AgentState) -> dict:
    with traced_span("agent_node", "knowledge_graph", "node") as node_span:

        # Span 1: 图数据库查询 —— 这是外部 I/O,需要追踪延迟
        with traced_span("tool_exec", "knowledge_graph", "tool",
                         tool="neo4j_query") as tool_span:
            graph_result = neo4j.query(
                "MATCH (t:Teacher {name:$name})-[r:teaches]->(c:Course) RETURN c.name",
                {"name": state["entity_name"]}
            )
            tool_span.metadata["result_count"] = len(graph_result)
            tool_span.metadata["query_ms"] = graph_result.elapsed_ms

        # Span 2: LLM 把图结果整理成自然语言
        with traced_span("llm_call", "knowledge_graph", "llm") as llm_span:
            response = llm.invoke(format_graph_context(graph_result))
            llm_span.metadata["input_tokens"] = response.usage.input_tokens

        node_span.metadata["entities_found"] = len(graph_result)

    return {"kg_result": response.content}

细到这个程度后,出问题时就能定位到"是 Neo4j 查询慢了"还是"LLM 整理慢了"——不再是一个笼统的"knowledge_graph 耗时 3 秒"。

五、并行 Agent 的 Span 怎么关联

四个 Agent 并行跑——route、knowledge_graph、vector_search、info_agent。每个都有自己的子 Span。怎么在 Span 树里体现出"这四个是并行的"?

我加了一个"并行扇出"节点,把四个 Agent 的 Span 全挂在它下面:

supervisor (agent_node)
  ├── llm_call
  ├── parse_plan
  └── fanout_parallel (parallel)         ← 新增的父 Span,标记为并行
        ├── agent_node (route)
        │     ├── llm_call
        │     └── vector_search
        ├── agent_node (knowledge_graph)
        │     ├── neo4j_query
        │     └── llm_call
        ├── agent_node (vector_search)
        │     └── milvus_search
        └── agent_node (info_agent)
              └── llm_call

实现:在并行执行函数外面套一层 traced_span

def execute_parallel_agents(agent_tasks: dict) -> dict:
    # 包一层 parallel 类型的 Span,agent 写 supervisor
    with traced_span("fanout_parallel", "supervisor", "parallel",
                     agents=list(agent_tasks.keys())) as fanout_span:

        results = {}
        with ThreadPoolExecutor(max_workers=len(agent_tasks)) as executor:
            futures = {executor.submit(fn): name
                       for name, fn in agent_tasks.items()}

            # as_completed:谁先跑完先收谁——不用等最慢的那个
            for future in as_completed(futures):
                name = futures[future]
                try:
                    results[name] = future.result(timeout=30)
                except Exception as e:
                    results[name] = {"error": str(e)}

        # 汇总信息
        fanout_span.metadata["agent_count"] = len(agent_tasks)
        fanout_span.metadata["success_count"] = sum(
            1 for r in results.values() if "error" not in r
        )

    return results

每个 Agent 内部自己执行的 traced_span 会自动把 parent_id 指向当前线程的活跃 Span。而每个 Agent 在执行时,活跃 Span 就是 fanout_parallel——因为 ContextVar 在每个线程里有独立副本,线程 A 的 _current_span 和线程 B 的不冲突,但它们的共同父 Span(fanout_parallel)是在主线程创建后传过来的。

一个关键细节parent_id 存的是 fanout_parallelspan_id,这是一个字符串。子线程里通过 _current_span.get() 拿到的就是 fanout_parallel 这个 Span 对象,所以 parent_id 能正确填上。这个过程看起来"自动",但实际上全靠 ContextVar 的线程本地存储机制保证的。

六、存和查

追踪数据只在内存里没用。我直接用 SQLite 存——最省事的方案,不需要装数据库、不需要配连接、一个文件搞定。

import sqlite3, json

class TraceCollector:
    def __init__(self, db_path="data/traces.db"):
        self.conn = sqlite3.connect(db_path, check_same_thread=False)
        # WAL 模式允许读写并发——多个线程同时写不会锁
        self.conn.execute("PRAGMA journal_mode=WAL")
        self._init_tables()

    def _init_tables(self):
        self.conn.executescript("""
            CREATE TABLE IF NOT EXISTS traces (
                trace_id TEXT PRIMARY KEY,
                session_id TEXT NOT NULL,
                user_input TEXT,
                span_count INTEGER DEFAULT 0,
                total_ms REAL,
                created_at REAL NOT NULL
            );
            CREATE TABLE IF NOT EXISTS spans (
                span_id TEXT PRIMARY KEY,
                trace_id TEXT NOT NULL REFERENCES traces(trace_id),
                parent_id TEXT,
                name TEXT NOT NULL,
                agent TEXT NOT NULL,
                span_type TEXT NOT NULL,
                status TEXT NOT NULL,
                duration_ms REAL,
                metadata TEXT DEFAULT '{}',
                created_at REAL NOT NULL
            );
            -- 这几个索引覆盖了最常用的查询模式
            CREATE INDEX IF NOT EXISTS idx_spans_trace ON spans(trace_id);
            CREATE INDEX IF NOT EXISTS idx_spans_agent ON spans(agent);
            CREATE INDEX IF NOT EXISTS idx_spans_type ON spans(span_type);
            CREATE INDEX IF NOT EXISTS idx_spans_created ON spans(created_at);
        """)
        self.conn.commit()

    def save_trace(self, trace: Trace):
        """请求结束时调用,一次性写入整个 Trace 的所有 Span"""
        total_ms = 0
        for span in trace.spans:
            span_ms = 0
            if span.end_time:
                span_ms = (span.end_time - span.start_time) * 1000

            total_ms += span_ms
            self.conn.execute(
                """INSERT INTO spans (span_id, trace_id, parent_id, name,
                   agent, span_type, status, duration_ms, metadata, created_at)
                   VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""",
                (span.span_id, trace.trace_id, span.parent_id,
                 span.name, span.agent, span.metadata.get("type", ""),
                 span.status, span_ms,
                 json.dumps(span.metadata, ensure_ascii=False),
                 trace.created_at)
            )

        self.conn.execute(
            """INSERT INTO traces (trace_id, session_id, user_input,
               span_count, total_ms, created_at) VALUES (?, ?, ?, ?, ?, ?)""",
            (trace.trace_id, trace.session_id, trace.user_input,
             len(trace.spans), total_ms, trace.created_at)
        )
        self.conn.commit()

注意一个细节save_trace 是在请求结束后统一调一次,不是每个 Span 调一次。这么做的原因是减少磁盘 I/O——一次请求可能有 15-20 个 Span,如果每个 Span 都单独写一次 SQLite,就是 15-20 次 fsync。SQLite 的 WAL 模式虽然快,但频繁 fsync 还是慢。在请求结束时批量写入,一次 commit,磁盘友好。

查数据的部分我封装了一个查询类,把最常用的几种查法写了上去:

class TraceQueries:
    def __init__(self, collector: TraceCollector):
        self.db = collector.conn

    def slowest_llm_calls(self, limit=10):
        """找最近最慢的 LLM 调用——排查延迟问题第一个用的查询"""
        return self.db.execute("""
            SELECT agent, duration_ms,
                   json_extract(metadata, '$.model') as model,
                   json_extract(metadata, '$.input_tokens') as tokens_in,
                   json_extract(metadata, '$.output_tokens') as tokens_out
            FROM spans
            WHERE span_type = 'llm' AND status = 'success'
            ORDER BY duration_ms DESC LIMIT ?
        """, (limit,)).fetchall()

    def agent_latency_breakdown(self, trace_id: str):
        """看一次请求里每个 Agent 的耗时分布"""
        return self.db.execute("""
            SELECT agent, name, span_type, duration_ms
            FROM spans
            WHERE trace_id = ? AND name = 'agent_node'
            ORDER BY duration_ms DESC
        """, (trace_id,)).fetchall()

    def error_rate_by_agent(self, hours=24):
        """过去 N 小时各 Agent 的报错比例"""
        cutoff = time.time() - hours * 3600
        return self.db.execute("""
            SELECT agent,
                   COUNT(*) as total,
                   SUM(CASE WHEN status = 'error' THEN 1 ELSE 0 END) as errors,
                   ROUND(100.0 * SUM(CASE WHEN status = 'error' THEN 1 ELSE 0 END)
                         / COUNT(*), 1) as error_pct
            FROM spans
            WHERE created_at > ?
            GROUP BY agent
            ORDER BY error_pct DESC
        """, (cutoff,)).fetchall()

    def daily_token_cost(self):
        """按天统计 token 消耗——看钱花哪了"""
        return self.db.execute("""
            SELECT date(created_at, 'unixepoch') as day,
                   SUM(json_extract(metadata, '$.input_tokens')) as total_input,
                   SUM(json_extract(metadata, '$.output_tokens')) as total_output
            FROM spans
            WHERE span_type = 'llm'
            GROUP BY day ORDER BY day DESC LIMIT 7
        """).fetchall()

json_extract 是 SQLite 的内置函数,能从 metadata 那个 JSON 字段里直接提一个 key 出来,不用把整个 JSON 读到 Python 再解析。这个函数让"按 metadata 里的某个字段过滤/排序"变得一脚 SQL 的事。

七、请求入口——把TraceContext串起来

def handle_user_request(session_id: str, user_input: str) -> dict:
    # 请求进来,先创建 Trace
    ctx = TraceContext(session_id=session_id, user_input=user_input)

    try:
        # 整个请求包在 root Span 里
        with traced_span("request_root", "gateway", "node",
                         session=session_id) as root_span:
            result = agent_graph.invoke({"user_input": user_input, ...})
            root_span.metadata["output_length"] = len(str(result))

        # 请求处理完了,一次性写入所有 Span
        collector = TraceCollector()
        collector.save_trace(ctx.trace)
        return result

    finally:
        # 无论如何要 close——释放 ContextVar
        ctx.close()

八、回到开头那个 5 秒的 bug

有了这套追踪之后再遇到"偶尔慢"的问题,直接查最慢的 LLM 调用:

SELECT agent, duration_ms, json_extract(metadata, '$.input_tokens') as tokens
FROM spans
WHERE span_type = 'llm' AND duration_ms > 3000
ORDER BY duration_ms DESC LIMIT 5;

结果:

agent              duration_ms   tokens
────────────────────────────────────────
route              5120          12800    ← token 爆了
knowledge_graph    1200          1800
supervisor         980           2100

route agent 某次 LLM 调用的 input token 飙到了 12800,正常应该是 2000 左右。回去查那次请求的上下文,发现是向量检索条件太宽,捞了几十个 chunk 全塞进 prompt 把 token 撑爆了。修起来就一行代码——把向量检索的 top_k 从 20 改到 10。

没有追踪之前,这个问题我只能靠猜。加 Span 之后,一条 SQL 掏出来,10 秒定位。

九、这套方案的成本

SQLite 一个文件,不需要 ELK、Jaeger、Prometheus 这些东西。每天几百次请求,每个请求存十几 KB,一个月也就几百 MB。对于个人项目或者小团队完全够了。

如果量大了 SQLite 扛不住,换 PostgreSQL 只需要改 TraceCollector 的实现(接口不变),Span 数据模型不用动。


整件事下来其实就三板斧:

  1. Span 数据模型——Trace 套 Span,Span 有类型、有 agent、有 metadata 字典。metadata 是开放的,LLM 调用和工具执行各自塞不同字段
  2. ContextVar 自动传上下文——traced_span with 语句包关键代码,父子 Span 自动关联。对业务代码的侵入就是两行
  3. SQLite 存 + SQL 查——两个表,四个索引。请求结束批量写。查最慢的 LLM、查报错率、查 token 消耗都是 SQL 一句

三百行代码不到,但出 bug 的时候省的时间是真的。