AI PRO·Harness Day 8 Observability层设计

作者:


一、引言:可观测性是Agent的”黑匣子”

2024年3月,某头部AI公司的一次内部事故复盘会上,工程师们面对一个棘手的问题:他们的客服Agent在凌晨2点到4点之间,突然开始向用户推荐竞品产品。日志显示Agent正常运行,没有报错,没有异常调用,一切看起来”正常”——直到用户投诉涌入。

事后花了72小时才定位到根因:一个RAG检索器的知识库在凌晨被错误更新,引入了一批含竞品信息的文档。Agent没有做错任何事——它忠实地基于检索到的内容回答了用户的问题。问题在于,没有人能实时看到”Agent检索了什么”和”Agent基于什么在回答”

这就是缺乏可观测性的代价。

在传统软件工程中,可观测性(Observability)已经是一个成熟的概念——日志(Logs)、指标(Metrics)、追踪(Traces)三大支柱支撑起了整个DevOps和SRE体系。但在AI Agent的世界里,可观测性面临着前所未有的挑战:

传统服务的行为是确定性的。 给定相同的输入,一个REST API总会返回相同的输出(或者至少是可预测的输出)。但Agent的行为是非确定性的——同一个用户问题,Agent可能走完全不同的推理路径,调用不同的工具,使用不同的上下文,最终给出不同的回答。

传统服务的”状态”是有限的。 一个微服务的状态可以用几个数据库字段和一组配置参数完整描述。但Agent的”状态”包括了:当前的推理链、历史对话、检索到的上下文、工具调用的结果、LLM的temperature设置、甚至system prompt中的一个微小措辞变化。状态空间是近乎无限的。

传统服务的”错误”是明确的。 HTTP 500就是错误,超时就是错误。但Agent的”错误”是模糊的——Agent没有崩溃,没有报错,但它给出了一个有害的回答、或者在不必要的时候调用了昂贵的API、或者陷入了无效的推理循环。这些是语义层面的错误,传统的健康检查根本检测不到。

`

传统服务的可观测性:

请求 → [服务A] → [服务B] → [数据库] → 响应

↓ ↓ ↓

日志/指标 日志/指标 日志/指标

(确定性行为,明确的错误边界)

Agent的可观测性:

用户输入 → [推理] → [检索] → [工具调用] → [再推理] → [生成] → 输出

↓ ↓ ↓ ↓ ↓

“想什么” “找到什么” “做了什么” “怎么想的” “说了什么”

(非确定性,语义级错误,无限状态空间)

`

在六层架构中,Observability层是唯一一个横切所有其他层的存在。它不直接参与Agent的推理、记忆、上下文或工具执行,但它是所有这些层的”眼睛”和”耳朵”:

`

┌─────────────────────────────────────┐

│ 6. Observability(可观测性) │ ← 本篇 ★

├─────────────────────────────────────┤

│ 5. Eval(评估) │

├─────────────────────────────────────┤

│ 4. Memory(记忆) │

├─────────────────────────────────────┤

│ 3. Tool(工具) │

├─────────────────────────────────────┤

│ 2. Context(上下文) │

├─────────────────────────────────────┤

│ 1. Prompt(提示词) │

└─────────────────────────────────────┘

↑ ↑ ↑ ↑ ↑ ↑

横切所有层,采集全链路信号

`

本文将从可观测性的三大支柱出发,深入Agent特有的观测指标、分布式追踪方案、异常检测机制、成本监控策略,最后通过一个实战项目搭建完整的Agent监控Dashboard。


二、可观测性三大支柱(日志、指标、追踪)

2.1 三大支柱概览

可观测性的三大支柱——日志(Logs)、指标(Metrics)、追踪(Traces)——在Agent场景下各有独特的变化:

支柱 传统服务 Agent服务

|——|———|———-|

日志 请求/响应、错误堆栈 推理链、工具调用、上下文快照
指标 QPS、延迟、错误率 Token消耗、任务完成率、循环次数
追踪 跨服务调用链 推理→检索→工具→生成的决策链

三者的关系可以用一个比喻来理解:指标告诉你”什么出了问题”,日志告诉你”出了什么问题”,追踪告诉你”问题出在哪里”。 在Agent场景中,还需要加上第四句:追踪还告诉你”Agent为什么做出这个决定”。

2.2 日志(Logs):Agent的”思维日记”

Agent的日志远比传统服务复杂。一条完整的Agent日志应该包含以下维度:

`python

import time

import uuid

import json

from dataclasses import dataclass, field, asdict

from typing import Any, Optional

from enum import Enum

class LogLevel(Enum):

DEBUG = “debug”

INFO = “info”

WARN = “warn”

ERROR = “error”

class SpanType(Enum):

INFERENCE = “inference” # LLM推理

TOOL_CALL = “tool_call” # 工具调用

RETRIEVAL = “retrieval” # 知识检索

PLANNING = “planning” # 规划分解

SYNTHESIS = “synthesis” # 结果综合

@dataclass

class AgentLogEntry:

“””Agent日志条目——远比传统日志丰富”””

# 基础字段

timestamp: float = field(default_factory=time.time)

level: LogLevel = LogLevel.INFO

message: str = “”

# 身份标识

trace_id: str = field(default_factory=lambda: str(uuid.uuid4()))

span_id: str = field(default_factory=lambda: str(uuid.uuid4()))

parent_span_id: Optional[str] = None

session_id: str = “”

user_id: str = “”

# Agent特有字段

span_type: SpanType = SpanType.INFERENCE

model: str = “”

prompt_tokens: int = 0

completion_tokens: int = 0

total_cost_usd: float = 0.0

# 推理过程

system_prompt_hash: str = “” # system prompt的指纹(非明文)

user_input: str = “”

model_output: str = “”

reasoning_steps: list[str] = field(default_factory=list)

# 工具调用

tool_name: Optional[str] = None

tool_input: Optional[dict] = None

tool_output: Optional[str] = None

tool_duration_ms: Optional[float] = None

tool_status: Optional[str] = None # success / error / timeout

# 上下文

retrieved_docs: list[dict] = field(default_factory=list)

context_window_used: int = 0

context_window_max: int = 0

# 元数据

metadata: dict[str, Any] = field(default_factory=dict)

def to_json(self) -> str:

return json.dumps(asdict(self), default=str, ensure_ascii=False)

`

这里有一个关键的设计决策:日志应该记录什么粒度的内容?

最细粒度的做法是记录LLM的完整输入输出——每一轮对话的完整prompt和response。但这会带来巨大的存储开销,而且存在隐私风险(用户输入可能包含敏感信息)。一个实用的策略是分层记录

`python

class LogGranularity(Enum):

“””日志粒度等级”””

TRACE = “trace” # 完整输入输出(开发环境)

DETAIL = “detail” # 摘要 + 关键字段(staging)

SUMMARY = “summary” # 仅统计指标(生产环境)

class AgentLogger:

def __init__(self, granularity: LogGranularity = LogGranularity.DETAIL):

self.granularity = granularity

self.entries: list[AgentLogEntry] = []

def log_inference(self, entry: AgentLogEntry):

if self.granularity == LogGranularity.TRACE:

# 记录完整的prompt和response

self.entries.append(entry)

elif self.granularity == LogGranularity.DETAIL:

# 记录摘要:输入的前200字符 + 输出的前200字符

entry.user_input = entry.user_input[:200] + “…”

entry.model_output = entry.model_output[:200] + “…”

# 只记录检索文档的ID,不记录全文

entry.retrieved_docs = [

{“doc_id”: d.get(“id”), “score”: d.get(“score”)}

for d in entry.retrieved_docs

]

self.entries.append(entry)

elif self.granularity == LogGranularity.SUMMARY:

# 只记录统计指标

summary = AgentLogEntry(

timestamp=entry.timestamp,

level=entry.level,

message=f”[{entry.span_type.value}] tokens={entry.prompt_tokens + entry.completion_tokens}”,

trace_id=entry.trace_id,

span_id=entry.span_id,

session_id=entry.session_id,

span_type=entry.span_type,

model=entry.model,

prompt_tokens=entry.prompt_tokens,

completion_tokens=entry.completion_tokens,

total_cost_usd=entry.total_cost_usd,

)

self.entries.append(summary)

`

分层记录的核心思想是:开发时看得见一切,生产时只留统计量,出问题时可以临时调高粒度。

2.3 指标(Metrics):Agent的”体检报告”

指标是可观测性中最适合自动化的部分。传统服务关注的指标是QPS、延迟、错误率;Agent需要关注一套全新的指标体系(下一章详述),但底层的采集和存储机制是通用的。

推荐使用OpenTelemetry作为指标采集的标准框架:

`python

from opentelemetry import metrics

from opentelemetry.sdk.metrics import MeterProvider

from opentelemetry.sdk.metrics.export import (

ConsoleMetricExporter,

PeriodicExportingMetricReader,

)

# 初始化指标采集

reader = PeriodicExportingMetricReader(

ConsoleMetricExporter(), export_interval_millis=10000

)

provider = MeterProvider(metric_readers=[reader])

metrics.set_meter_provider(provider)

meter = metrics.get_meter(“agent-observability”)

# 定义Agent指标

token_counter = meter.create_counter(

“agent.tokens.total”,

description=”Total tokens consumed”,

unit=”tokens”,

)

tool_duration_histogram = meter.create_histogram(

“agent.tool.duration”,

description=”Tool execution duration”,

unit=”ms”,

)

task_completion_gauge = meter.create_up_down_counter(

“agent.task.completions”,

description=”Task completion counter”,

unit=”tasks”,

)

active_sessions_gauge = meter.create_up_down_counter(

“agent.sessions.active”,

description=”Currently active sessions”,

unit=”sessions”,

)

`

2.4 追踪(Traces):Agent的”推理X光片”

追踪是三大支柱中对Agent最有价值的一个。传统服务的追踪记录的是”请求经过了哪些服务”,Agent的追踪记录的是”决策经过了哪些步骤”。

一次完整的Agent执行,追踪结构如下:

`

Trace: user_query_abc123 (总耗时 4.2s)

├── Span: planning (0.8s)

│ ├── LLM Call: claude-sonnet-4-20250514 (prompt: 1200 tokens, completion: 300 tokens)

│ └── Output: “需要分3步执行: 1. 搜索文档 2. 分析结果 3. 生成报告”

├── Span: step_1_retrieval (1.1s)

│ ├── Vector Search: knowledge_base (top_k=5, 返回5条)

│ └── Rerank: cross-encoder (5→3条)

├── Span: step_2_tool_call (1.5s)

│ ├── Tool: search_database (query=”SELECT …”, duration=800ms)

│ └── Tool: call_api (endpoint=”/v2/analyze”, duration=700ms)

├── Span: step_3_synthesis (0.8s)

│ ├── LLM Call: claude-sonnet-4-20250514 (prompt: 3500 tokens, completion: 800 tokens)

│ └── Output: “根据分析结果…”

└── Metrics:

├── Total Tokens: 5800 (prompt: 4700, completion: 1100)

├── Total Cost: $0.0174

├── Tool Calls: 2 (all successful)

└── LLM Calls: 2

`

用代码实现这种追踪结构:

`python

from contextlib import contextmanager

import time

import uuid

class Span:

“””追踪中的一个跨度”””

def __init__(self, name: str, span_type: str, parent: “Span | None” = None):

self.name = name

self.span_type = span_type

self.span_id = str(uuid.uuid4())[:8]

self.parent = parent

self.children: list[Span] = []

self.start_time: float = 0

self.end_time: float = 0

self.attributes: dict = {}

self.events: list[dict] = []

self.status: str = “ok”

@property

def duration_ms(self) -> float:

return (self.end_time – self.start_time) * 1000

class AgentTracer:

“””Agent分布式追踪器”””

def __init__(self):

self.current_span: Span | None = None

self.traces: dict[str, list[Span]] = {}

def start_trace(self, trace_id: str, name: str, span_type: str) -> Span:

span = Span(name, span_type)

self.current_span = span

self.traces.setdefault(trace_id, []).append(span)

return span

@contextmanager

def trace(self, name: str, span_type: str, **attributes):

“””上下文管理器,自动管理span的生命周期”””

parent = self.current_span

span = Span(name, span_type, parent)

span.attributes.update(attributes)

if parent:

parent.children.append(span)

self.current_span = span

span.start_time = time.time()

try:

yield span

except Exception as e:

span.status = “error”

span.events.append({

“name”: “exception”,

“message”: str(e),

“type”: type(e).__name__,

})

raise

finally:

span.end_time = time.time()

self.current_span = parent

def print_trace(self, trace_id: str):

“””打印追踪树”””

spans = self.traces.get(trace_id, [])

if not spans:

print(“No trace found”)

return

def _print_span(span: Span, indent: int = 0):

prefix = “│ ” * indent + “├── ” if indent > 0 else “”

status_icon = “✓” if span.status == “ok” else “✗”

attrs = “, “.join(f”{k}={v}” for k, v in span.attributes.items())

print(f”{prefix}[{status_icon}] {span.name} ({span.duration_ms:.0f}ms) {attrs}”)

for child in span.children:

_print_span(child, indent + 1)

for span in spans:

_print_span(span)

`


三、Agent特有指标(Token消耗、工具调用、循环次数、任务完成率)

传统服务监控的核心指标是”四个黄金信号”(Google SRE):延迟、流量、错误率、饱和度。Agent需要一套全新的指标体系。

3.1 Agent四大指标类别

我把Agent的指标分为四大类别:成本指标、质量指标、效率指标、安全指标

`python

from dataclasses import dataclass, field

from enum import Enum

class MetricCategory(Enum):

COST = “cost” # 成本

QUALITY = “quality” # 质量

EFFICIENCY = “efficiency” # 效率

SAFETY = “safety” # 安全

@dataclass

class AgentMetrics:

“””一次Agent执行的完整指标集”””

# === 成本指标 ===

prompt_tokens: int = 0 # 输入token数

completion_tokens: int = 0 # 输出token数

total_cost_usd: float = 0.0 # 总成本(美元)

cache_hit_tokens: int = 0 # 缓存命中的token数

cache_miss_tokens: int = 0 # 缓存未命中的token数

# === 质量指标 ===

task_completed: bool = False # 任务是否完成

task_success: bool = False # 任务是否成功(用户确认)

hallucination_score: float = 0.0 # 幻觉检测分数 (0-1)

relevance_score: float = 0.0 # 回答相关性分数 (0-1)

user_rating: int = 0 # 用户评分 (1-5)

# === 效率指标 ===

total_duration_ms: float = 0.0 # 总耗时

llm_calls: int = 0 # LLM调用次数

tool_calls: int = 0 # 工具调用次数

tool_errors: int = 0 # 工具调用失败次数

loop_count: int = 0 # 推理循环次数

retry_count: int = 0 # 重试次数

first_token_latency_ms: float = 0.0 # 首token延迟

tokens_per_second: float = 0.0 # 生成速度

# === 安全指标 ===

guardrail_triggers: int = 0 # 安全护栏触发次数

pii_detected: bool = False # 是否检测到PII

sensitive_tool_blocks: int = 0 # 敏感工具被阻止次数

max_tool_permission: str = “” # 最高工具权限等级

@property

def cache_hit_rate(self) -> float:

total = self.cache_hit_tokens + self.cache_miss_tokens

return self.cache_hit_tokens / total if total > 0 else 0.0

@property

def tool_error_rate(self) -> float:

return self.tool_errors / self.tool_calls if self.tool_calls > 0 else 0.0

@property

def efficiency_score(self) -> float:

“””效率分数:用最少的token和工具调用完成任务”””

if not self.task_completed:

return 0.0

# 越少的token和工具调用完成任务,效率越高

token_penalty = min(self.prompt_tokens + self.completion_tokens, 50000) / 50000

tool_penalty = min(self.tool_calls, 20) / 20

loop_penalty = min(self.loop_count, 10) / 10

return 1.0 – (token_penalty * 0.4 + tool_penalty * 0.3 + loop_penalty * 0.3)

`

3.2 Token消耗追踪

Token是Agent的”燃料”,也是最大的成本来源。精细化的Token追踪是成本控制的基础:

`python

class TokenTracker:

“””Token消耗追踪器”””

# 模型定价(每百万token,美元)

PRICING = {

“gpt-4o”: {“input”: 2.50, “output”: 10.00, “cached_input”: 1.25},

“claude-sonnet-4-20250514”: {“input”: 3.00, “output”: 15.00, “cached_input”: 0.30},

“gpt-4o-mini”: {“input”: 0.15, “output”: 0.60, “cached_input”: 0.075},

“deepseek-v3”: {“input”: 0.27, “output”: 1.10, “cached_input”: 0.07},

}

def __init__(self):

self.records: list[dict] = []

def record(

self,

model: str,

prompt_tokens: int,

completion_tokens: int,

cached_tokens: int = 0,

trace_id: str = “”,

span_id: str = “”,

purpose: str = “”, # “planning”, “synthesis”, “tool_selection” 等

):

pricing = self.PRICING.get(model, {“input”: 5.0, “output”: 15.0, “cached_input”: 2.5})

# 计算成本

regular_input = prompt_tokens – cached_tokens

input_cost = (regular_input * pricing[“input”] + cached_tokens * pricing[“cached_input”]) / 1_000_000

output_cost = completion_tokens * pricing[“output”] / 1_000_000

total_cost = input_cost + output_cost

record = {

“model”: model,

“prompt_tokens”: prompt_tokens,

“completion_tokens”: completion_tokens,

“cached_tokens”: cached_tokens,

“input_cost_usd”: input_cost,

“output_cost_usd”: output_cost,

“total_cost_usd”: total_cost,

“trace_id”: trace_id,

“span_id”: span_id,

“purpose”: purpose,

“timestamp”: time.time(),

}

self.records.append(record)

return total_cost

def get_session_summary(self) -> dict:

if not self.records:

return {“total_cost”: 0, “total_tokens”: 0}

total_cost = sum(r[“total_cost_usd”] for r in self.records)

total_prompt = sum(r[“prompt_tokens”] for r in self.records)

total_completion = sum(r[“completion_tokens”] for r in self.records)

total_cached = sum(r[“cached_tokens”] for r in self.records)

# 按用途分组

by_purpose: dict[str, float] = {}

for r in self.records:

purpose = r.get(“purpose”, “unknown”)

by_purpose[purpose] = by_purpose.get(purpose, 0) + r[“total_cost_usd”]

return {

“total_cost_usd”: round(total_cost, 6),

“total_prompt_tokens”: total_prompt,

“total_completion_tokens”: total_completion,

“total_cached_tokens”: total_cached,

“cache_hit_rate”: total_cached / total_prompt if total_prompt > 0 else 0,

“llm_calls”: len(self.records),

“cost_by_purpose”: by_purpose,

}

`

3.3 循环次数监控

Agent的推理循环(loop)是一个需要重点关注的指标。循环次数过多通常意味着Agent在”原地打转”:

`python

class LoopDetector:

“””推理循环检测器”””

def __init__(self, max_loops: int = 10, similarity_threshold: float = 0.85):

self.max_loops = max_loops

self.similarity_threshold = similarity_threshold

self.loop_history: list[dict] = []

self.warning_callbacks: list = []

def on_warning(self, callback):

self.warning_callbacks.append(callback)

def record_iteration(self, iteration: dict):

“””

记录一次推理迭代

iteration: {

“reasoning”: str, # 推理内容

“tool_call”: str, # 调用的工具

“tool_input”: dict, # 工具输入

“output”: str, # 输出结果

}

“””

self.loop_history.append({

**iteration,

“timestamp”: time.time(),

“iteration_number”: len(self.loop_history) + 1,

})

# 检查是否超过最大循环次数

if len(self.loop_history) >= self.max_loops:

self._trigger_warning(“max_loops_exceeded”, {

“current_loops”: len(self.loop_history),

“max_loops”: self.max_loops,

})

# 检查是否陷入重复循环

if len(self.loop_history) >= 3:

recent = self.loop_history[-3:]

if self._is_repeating(recent):

self._trigger_warning(“repeating_pattern”, {

“pattern_length”: 3,

“recent_tools”: [r.get(“tool_call”) for r in recent],

})

def _is_repeating(self, iterations: list[dict]) -> bool:

“””检测最近的迭代是否在重复”””

tools = [i.get(“tool_call”) for i in iterations]

inputs = [str(i.get(“tool_input”, {})) for i in iterations]

# 检查是否调用了相同的工具

if len(set(tools)) == 1 and tools[0] is not None:

# 还需要检查输入是否相似

if self._text_similarity(inputs[0], inputs[1]) > self.similarity_threshold:

return True

return False

def _text_similarity(self, a: str, b: str) -> float:

“””简单的文本相似度计算”””

if not a or not b:

return 0.0

set_a = set(a.split())

set_b = set(b.split())

if not set_a or not set_b:

return 0.0

intersection = set_a & set_b

union = set_a | set_b

return len(intersection) / len(union)

def _trigger_warning(self, warning_type: str, details: dict):

for callback in self.warning_callbacks:

callback(warning_type, details)

@property

def current_loop_count(self) -> int:

return len(self.loop_history)

def reset(self):

self.loop_history.clear()

`

3.4 任务完成率

任务完成率是衡量Agent有效性的核心指标,但它比听起来要复杂得多:

`python

class TaskTracker:

“””任务完成率追踪器”””

def __init__(self):

self.tasks: list[dict] = []

def start_task(self, task_id: str, task_type: str, description: str):

self.tasks.append({

“task_id”: task_id,

“task_type”: task_type,

“description”: description,

“status”: “in_progress”,

“start_time”: time.time(),

“end_time”: None,

“outcome”: None,

“failure_reason”: None,

“metrics”: {},

})

def complete_task(self, task_id: str, success: bool, reason: str = “”):

task = self._find_task(task_id)

if task:

task[“status”] = “completed” if success else “failed”

task[“end_time”] = time.time()

task[“outcome”] = “success” if success else “failure”

task[“failure_reason”] = reason if not success else None

def add_metrics(self, task_id: str, metrics: dict):

task = self._find_task(task_id)

if task:

task[“metrics”].update(metrics)

def _find_task(self, task_id: str) -> dict | None:

for task in self.tasks:

if task[“task_id”] == task_id:

return task

return None

def get_completion_stats(self) -> dict:

if not self.tasks:

return {“total”: 0, “completion_rate”: 0, “success_rate”: 0}

completed = [t for t in self.tasks if t[“status”] != “in_progress”]

successful = [t for t in completed if t[“outcome”] == “success”]

# 按任务类型分组

by_type: dict[str, dict] = {}

for task in self.tasks:

t = task[“task_type”]

if t not in by_type:

by_type[t] = {“total”: 0, “success”: 0, “failed”: 0}

by_type[t][“total”] += 1

if task[“outcome”] == “success”:

by_type[t][“success”] += 1

elif task[“outcome”] == “failure”:

by_type[t][“failed”] += 1

# 失败原因统计

failure_reasons: dict[str, int] = {}

for task in completed:

if task[“failure_reason”]:

r = task[“failure_reason”]

failure_reasons[r] = failure_reasons.get(r, 0) + 1

return {

“total_tasks”: len(self.tasks),

“completed”: len(completed),

“successful”: len(successful),

“completion_rate”: len(completed) / len(self.tasks),

“success_rate”: len(successful) / len(completed) if completed else 0,

“by_type”: by_type,

“failure_reasons”: failure_reasons,

“avg_duration_ms”: self._avg_duration(completed),

}

def _avg_duration(self, tasks: list[dict]) -> float:

durations = [

(t[“end_time”] – t[“start_time”]) * 1000

for t in tasks

if t[“end_time”] and t[“start_time”]

]

return sum(durations) / len(durations) if durations else 0

`


四、分布式追踪(LangSmith、Langfuse、Helicone)

4.1 为什么需要分布式追踪

当Agent从一个简单的”LLM + 工具”进化到多Agent协作、多步骤推理、RAG增强的复杂系统时,单靠日志已经无法看清全貌。分布式追踪提供了端到端的可视化能力,让你能看到每一次推理的完整”决策树”。

2024-2026年,Agent追踪领域形成了三个主流方案:LangSmith(LangChain生态)、Langfuse(开源)、Helicone(专注LLM代理层)。它们各有侧重,但在核心概念上高度一致。

4.2 Langfuse:开源追踪的标杆

Langfuse是目前最成熟的开源Agent追踪方案,核心概念包括:Trace(一次完整交互)、Span(一个子步骤)、Generation(一次LLM调用)、Event(一个标记事件)。

`python

# pip install langfuse

from langfuse import Langfuse

from langfuse.decorators import observe, langfuse_context

# 初始化

langfuse = Langfuse(

public_key=”pk-…”,

secret_key=”sk-…”,

host=”https://cloud.langfuse.com”, # 或自部署地址

)

@observe(as_type=”generation”)

def call_llm(prompt: str, model: str = “claude-sonnet-4-20250514”) -> str:

“””被追踪的LLM调用”””

# Langfuse自动记录输入、输出、模型、token数

response = client.chat.completions.create(

model=model,

messages=[{“role”: “user”, “content”: prompt}],

)

# 更新Langfuse上下文

langfuse_context.update_current_observation(

usage={

“input”: response.usage.prompt_tokens,

“output”: response.usage.completion_tokens,

“total”: response.usage.total_tokens,

},

model=model,

)

return response.choices[0].message.content

@observe(as_type=”span”)

def retrieve_documents(query: str, top_k: int = 5) -> list[dict]:

“””被追踪的检索操作”””

results = vector_store.search(query, top_k=top_k)

# 记录检索元数据

langfuse_context.update_current_observation(

metadata={

“query”: query,

“top_k”: top_k,

“results_count”: len(results),

“scores”: [r[“score”] for r in results],

}

)

return results

@observe() # 最外层自动创建Trace

def agent_run(user_input: str) -> str:

“””被追踪的完整Agent执行”””

# Step 1: 规划

plan = call_llm(f”请规划以下任务的执行步骤:{user_input}”)

# Step 2: 检索

docs = retrieve_documents(user_input)

# Step 3: 综合

context = “n”.join([d[“content”] for d in docs])

answer = call_llm(

f”基于以下信息回答问题。nn上下文:{context}nn问题:{user_input}”

)

return answer

`

4.3 LangSmith:LangChain生态的追踪

LangSmith与LangChain/LangGraph深度集成,如果你已经在使用LangChain生态,LangSmith是最自然的选择:

`python

# pip install langsmith

import os

os.environ[“LANGCHAIN_TRACING_V2”] = “true”

os.environ[“LANGCHAIN_API_KEY”] = “ls-…”

os.environ[“LANGCHAIN_PROJECT”] = “my-agent”

from langsmith import traceable

from langchain_openai import ChatOpenAI

from langchain.agents import AgentExecutor, create_openai_tools_agent

@traceable(

name=”agent_planner”,

run_type=”chain”,

tags=[“planning”, “v1”],

metadata={“model”: “gpt-4o”, “version”: “1.0”},

)

def plan_task(user_input: str) -> list[str]:

llm = ChatOpenAI(model=”gpt-4o”)

response = llm.invoke(

f”将以下任务分解为子步骤:{user_input}”

)

steps = response.content.split(“n”)

return [s.strip() for s in steps if s.strip()]

@traceable(

name=”knowledge_retrieval”,

run_type=”retriever”,

)

def retrieve(query: str) -> list[dict]:

# 模拟检索

return [

{“page_content”: “相关文档片段1”, “metadata”: {“source”: “doc1.pdf”}},

{“page_content”: “相关文档片段2”, “metadata”: {“source”: “doc2.pdf”}},

]

@traceable(name=”agent_pipeline”)

def run_agent(question: str) -> str:

steps = plan_task(question)

docs = retrieve(question)

# … 后续处理

return “最终回答”

`

4.4 Helicone:代理层追踪

Helicone采用不同的架构——它作为LLM API的反向代理,无需修改代码即可追踪所有请求:

`python

# Helicone方式:只修改API base URL,零代码侵入

import openai

client = openai.OpenAI(

api_key=”your-api-key”,

base_url=”https://oai.helicone.ai/v1″,

default_headers={

“Helicone-Auth”: “Bearer helicone-api-key”,

“Helicone-Property-Session”: “session-abc123”,

“Helicone-Property-Task-Type”: “code-generation”,

},

)

# 正常使用,Helicone自动捕获所有请求/响应

response = client.chat.completions.create(

model=”gpt-4o”,

messages=[{“role”: “user”, “content”: “Hello!”}],

# Helicone通过header自动关联元数据

)

`

4.5 三种方案对比

特性 LangSmith Langfuse Helicone

|——|———-|———|———|

部署方式 SaaS SaaS / 自部署 SaaS / 自部署
开源 是(MIT) 部分开源
侵入性 中(需用装饰器) 中(需用装饰器) 低(代理层)
LangChain集成 原生 一般
非LangChain支持
Prompt版本管理 内置 内置
评估框架 内置 内置 基础
成本追踪 好(核心功能)
自定义Dashboard 有限 灵活 灵活
推荐场景 LangChain重度用户 需要自部署/开源 零侵入快速接入

我的建议: 如果你刚起步,选Langfuse——开源、可自部署、社区活跃。如果你已经深度使用LangChain生态,LangSmith是最顺畅的选择。如果你追求零侵入,Helicone最快上手。


五、异常检测与告警

5.1 Agent异常分类

Agent的异常与传统服务不同,需要一套专门的分类体系:

`python

from enum import Enum

from dataclasses import dataclass

class AnomalySeverity(Enum):

LOW = “low” # 关注即可

MEDIUM = “medium” # 需要调查

HIGH = “high” # 需要立即处理

CRITICAL = “critical” # 需要立即干预

class AnomalyType(Enum):

# 性能异常

HIGH_LATENCY = “high_latency”

TOKEN_EXPLOSION = “token_explosion” # token使用量突然飙升

LOOP_DETECTED = “loop_detected” # 检测到无效循环

# 质量异常

HALLUCINATION = “hallucination” # 检测到幻觉

OFF_TOPIC = “off_topic” # 回答偏离主题

CONTRADICTION = “contradiction” # 回答自相矛盾

# 安全异常

PII_LEAK = “pii_leak” # PII信息泄露

INJUNCTION_DETECTED = “injection_detected” # 检测到prompt注入

UNAUTHORIZED_TOOL = “unauthorized_tool” # 未授权的工具调用

HARMFUL_OUTPUT = “harmful_output” # 有害内容输出

# 系统异常

COST_SPIKE = “cost_spike” # 成本突然飙升

API_DEGRADATION = “api_degradation” # 上游API质量下降

CONTEXT_OVERFLOW = “context_overflow” # 上下文溢出

@dataclass

class Anomaly:

type: AnomalyType

severity: AnomalySeverity

description: str

trace_id: str

span_id: str

metadata: dict

timestamp: float

acknowledged: bool = False

resolved: bool = False

`

5.2 基于规则的异常检测

最简单也最有效的异常检测方式是基于规则:

`python

import time

from collections import deque

class AnomalyDetector:

“””基于规则的Agent异常检测器”””

def __init__(self):

self.rules: list[dict] = []

self.anomalies: list[Anomaly] = []

self.metric_history: dict[str, deque] = {}

# 加载默认规则

self._load_default_rules()

def _load_default_rules(self):

self.rules = [

{

“name”: “token_explosion”,

“type”: AnomalyType.TOKEN_EXPLOSION,

“severity”: AnomalySeverity.HIGH,

“condition”: lambda m: (

m.get(“total_tokens”, 0) > 50000

),

“message”: “Token使用量超过50K阈值: {total_tokens}”,

},

{

“name”: “loop_detection”,

“type”: AnomalyType.LOOP_DETECTED,

“severity”: AnomalySeverity.HIGH,

“condition”: lambda m: (

m.get(“loop_count”, 0) > 8

),

“message”: “推理循环次数异常: {loop_count}”,

},

{

“name”: “cost_spike”,

“type”: AnomalyType.COST_SPIKE,

“severity”: AnomalySeverity.CRITICAL,

“condition”: lambda m: (

m.get(“session_cost_usd”, 0) > 1.0

),

“message”: “单次会话成本超过$1: ${session_cost_usd:.4f}”,

},

{

“name”: “high_latency”,

“type”: AnomalyType.HIGH_LATENCY,

“severity”: AnomalySeverity.MEDIUM,

“condition”: lambda m: (

m.get(“total_duration_ms”, 0) > 30000

),

“message”: “Agent响应时间超过30秒: {total_duration_ms}ms”,

},

{

“name”: “tool_error_rate”,

“type”: AnomalyType.API_DEGRADATION,

“severity”: AnomalySeverity.MEDIUM,

“condition”: lambda m: (

m.get(“tool_calls”, 0) >= 3 and

m.get(“tool_errors”, 0) / m.get(“tool_calls”, 1) > 0.5

),

“message”: “工具错误率超过50%: {tool_errors}/{tool_calls}”,

},

{

“name”: “context_overflow_risk”,

“type”: AnomalyType.CONTEXT_OVERFLOW,

“severity”: AnomalySeverity.HIGH,

“condition”: lambda m: (

m.get(“context_window_used”, 0) /

max(m.get(“context_window_max”, 1), 1) > 0.9

),

“message”: “上下文使用率超过90%”,

},

{

“name”: “empty_output”,

“type”: AnomalyType.OFF_TOPIC,

“severity”: AnomalySeverity.MEDIUM,

“condition”: lambda m: (

m.get(“task_completed”, True) and

len(m.get(“model_output”, “x”)) < 10

),

“message”: “Agent输出过短,可能未正常完成任务”,

},

]

def evaluate(self, metrics: dict, trace_id: str = “”, span_id: str = “”) -> list[Anomaly]:

“””评估一组指标,返回触发的异常”””

triggered = []

for rule in self.rules:

try:

if rule“condition”:

anomaly = Anomaly(

type=rule[“type”],

severity=rule[“severity”],

description=rule[“message”].format(**metrics),

trace_id=trace_id,

span_id=span_id,

metadata=metrics,

timestamp=time.time(),

)

triggered.append(anomaly)

self.anomalies.append(anomaly)

except (KeyError, ZeroDivisionError):

continue # 指标缺失时跳过该规则

return triggered

`

5.3 基于统计的异常检测

对于更微妙的异常(如成本缓慢上升、延迟逐渐恶化),需要基于统计的方法:

`python

import statistics

from collections import deque

class StatisticalAnomalyDetector:

“””基于统计的异常检测——Z-score方法”””

def __init__(self, window_size: int = 100, z_threshold: float = 3.0):

self.window_size = window_size

self.z_threshold = z_threshold

self.windows: dict[str, deque] = {}

def add_metric(self, name: str, value: float) -> dict | None:

“””

添加一个指标值,返回异常信息(如果检测到)

使用Z-score方法:如果新值偏离均值超过z_threshold个标准差,视为异常

“””

if name not in self.windows:

self.windows[name] = deque(maxlen=self.window_size)

window = self.windows[name]

# 需要足够的历史数据

if len(window) >= 10:

mean = statistics.mean(window)

stdev = statistics.stdev(window)

if stdev > 0:

z_score = abs(value – mean) / stdev

if z_score > self.z_threshold:

return {

“metric”: name,

“value”: value,

“mean”: mean,

“stdev”: stdev,

“z_score”: z_score,

“threshold”: self.z_threshold,

“is_anomaly”: True,

}

window.append(value)

return None

def get_baseline(self, name: str) -> dict | None:

“””获取指标的基线统计”””

window = self.windows.get(name)

if not window or len(window) < 5:

return None

return {

“metric”: name,

“count”: len(window),

“mean”: statistics.mean(window),

“median”: statistics.median(window),

“stdev”: statistics.stdev(window) if len(window) > 1 else 0,

“min”: min(window),

“max”: max(window),

“p95”: sorted(window)[int(len(window) * 0.95)] if len(window) > 1 else window[0],

}

`

5.4 告警系统

检测到异常后,需要一个灵活的告警系统:

`python

from abc import ABC, abstractmethod

class AlertChannel(ABC):

“””告警通道基类”””

@abstractmethod

def send(self, anomaly: Anomaly) -> bool:

pass

class SlackAlertChannel(AlertChannel):

def __init__(self, webhook_url: str):

self.webhook_url = webhook_url

def send(self, anomaly: Anomaly) -> bool:

severity_emoji = {

AnomalySeverity.LOW: “ℹ️”,

AnomalySeverity.MEDIUM: “⚠️”,

AnomalySeverity.HIGH: “🚨”,

AnomalySeverity.CRITICAL: “🔴”,

}

message = {

“text”: (

f”{severity_emoji[anomaly.severity]} “

f”Agent异常告警 [{anomaly.severity.value.upper()}]n”

f”类型: {anomaly.type.value}n”

f”描述: {anomaly.description}n”

f”Trace: {anomaly.trace_id}”

)

}

# requests.post(self.webhook_url, json=message)

print(f”[SLACK] {message[‘text’]}”)

return True

class WebhookAlertChannel(AlertChannel):

def __init__(self, url: str, headers: dict = None):

self.url = url

self.headers = headers or {}

def send(self, anomaly: Anomaly) -> bool:

payload = {

“type”: anomaly.type.value,

“severity”: anomaly.severity.value,

“description”: anomaly.description,

“trace_id”: anomaly.trace_id,

“timestamp”: anomaly.timestamp,

“metadata”: anomaly.metadata,

}

# requests.post(self.url, json=payload, headers=self.headers)

print(f”[WEBHOOK] {payload}”)

return True

class AlertManager:

“””告警管理器——路由异常到正确的通道”””

def __init__(self):

self.channels: dict[AnomalySeverity, list[AlertChannel]] = {

AnomalySeverity.LOW: [],

AnomalySeverity.MEDIUM: [],

AnomalySeverity.HIGH: [],

AnomalySeverity.CRITICAL: [],

}

self.suppression_cache: dict[str, float] = {}

self.suppression_window: float = 300 # 5分钟内相同异常只告警一次

def add_channel(self, severity: AnomalySeverity, channel: AlertChannel):

self.channels[severity].append(channel)

def alert(self, anomaly: Anomaly):

# 抑制重复告警

suppression_key = f”{anomaly.type.value}:{anomaly.severity.value}”

now = time.time()

if suppression_key in self.suppression_cache:

if now – self.suppression_cache[suppression_key] < self.suppression_window:

return # 抑制

self.suppression_cache[suppression_key] = now

# 发送到对应严重级别的所有通道

channels = self.channels.get(anomaly.severity, [])

for channel in channels:

try:

channel.send(anomaly)

except Exception as e:

print(f”告警发送失败: {e}”)

`


六、成本监控与优化

6.1 成本是Agent的”生死线”

成本监控在Agent场景中的重要性远超传统服务。原因很简单:传统服务的边际成本趋近于零,但Agent的每一次推理都有真金白银的token费用。

一个中等规模的客服Agent,每天处理1万次对话,每次对话平均5轮推理,每轮消耗2000 token——按GPT-4o定价计算:

`

日均token消耗:10,000 × 5 × 2,000 = 100,000,000 tokens

日均成本:~$150

月均成本:~$4,500

年均成本:~$54,000

`

如果Agent效率低下(循环多、prompt冗余、不必要的工具调用),成本可能翻2-3倍。没有成本监控的Agent,就是一个没有油量表的飞机。

6.2 成本监控架构

`python

from dataclasses import dataclass, field

from datetime import datetime, timedelta

from collections import defaultdict

@dataclass

class CostBudget:

“””成本预算配置”””

daily_limit_usd: float = 10.0

monthly_limit_usd: float = 200.0

per_session_limit_usd: float = 0.50

per_request_limit_usd: float = 0.10

# 告警阈值(百分比)

warning_threshold: float = 0.7 # 70%时告警

critical_threshold: float = 0.9 # 90%时紧急告警

class CostMonitor:

“””成本监控器”””

def __init__(self, budget: CostBudget):

self.budget = budget

self.costs_by_date: dict[str, float] = defaultdict(float)

self.costs_by_session: dict[str, float] = defaultdict(float)

self.costs_by_model: dict[str, float] = defaultdict(float)

self.costs_by_feature: dict[str, float] = defaultdict(float)

self.cost_history: list[dict] = []

def record_cost(

self,

amount_usd: float,

session_id: str = “”,

model: str = “”,

feature: str = “”,

metadata: dict = None,

):

today = datetime.now().strftime(“%Y-%m-%d”)

self.costs_by_date[today] += amount_usd

self.costs_by_session[session_id] += amount_usd

self.costs_by_model[model] += amount_usd

self.costs_by_feature[feature] += amount_usd

self.cost_history.append({

“timestamp”: datetime.now().isoformat(),

“amount_usd”: amount_usd,

“session_id”: session_id,

“model”: model,

“feature”: feature,

“metadata”: metadata or {},

})

# 检查预算

return self._check_budget(session_id)

def _check_budget(self, session_id: str) -> list[str]:

“””检查预算,返回告警列表”””

alerts = []

today = datetime.now().strftime(“%Y-%m-%d”)

# 检查日预算

daily_cost = self.costs_by_date[today]

daily_ratio = daily_cost / self.budget.daily_limit_usd

if daily_ratio >= 1.0:

alerts.append(f”🚨 日预算已耗尽: ${daily_cost:.4f} / ${self.budget.daily_limit_usd:.2f}”)

elif daily_ratio >= self.budget.critical_threshold:

alerts.append(f”🔴 日预算即将耗尽: {daily_ratio:.0%}”)

elif daily_ratio >= self.budget.warning_threshold:

alerts.append(f”⚠️ 日预算使用: {daily_ratio:.0%}”)

# 检查会话预算

session_cost = self.costs_by_session.get(session_id, 0)

if session_cost > self.budget.per_session_limit_usd:

alerts.append(f”🚨 会话预算超限 [{session_id[:8]}]: ${session_cost:.4f}”)

return alerts

def get_dashboard_data(self) -> dict:

“””获取Dashboard数据”””

today = datetime.now().strftime(“%Y-%m-%d”)

# 计算本月成本

now = datetime.now()

month_start = now.replace(day=1).strftime(“%Y-%m-%d”)

monthly_cost = sum(

cost for date, cost in self.costs_by_date.items()

if date >= month_start

)

# 最近7天趋势

last_7_days = []

for i in range(7):

date = (now – timedelta(days=i)).strftime(“%Y-%m-%d”)

last_7_days.append({

“date”: date,

“cost_usd”: self.costs_by_date.get(date, 0),

})

last_7_days.reverse()

return {

“today_cost_usd”: self.costs_by_date[today],

“today_budget_remaining_usd”: max(0, self.budget.daily_limit_usd – self.costs_by_date[today]),

“monthly_cost_usd”: monthly_cost,

“monthly_budget_remaining_usd”: max(0, self.budget.monthly_limit_usd – monthly_cost),

“cost_by_model”: dict(self.costs_by_model),

“cost_by_feature”: dict(self.costs_by_feature),

“top_expensive_sessions”: sorted(

self.costs_by_session.items(), key=lambda x: x[1], reverse=True

)[:10],

“last_7_days_trend”: last_7_days,

}

`

6.3 成本优化策略

监控成本是第一步,优化成本才是目的。以下是经过实战验证的策略:

`python

class CostOptimizer:

“””成本优化器——基于观测数据的自动优化”””

def __init__(self, token_tracker: TokenTracker):

self.token_tracker = token_tracker

def analyze_and_recommend(self) -> list[dict]:

“””分析Token使用模式,生成优化建议”””

recommendations = []

records = self.token_tracker.records

if not records:

return recommendations

# 1. 检测是否有不必要的大模型调用

expensive_calls = [

r for r in records

if r[“model”] in (“gpt-4o”, “claude-sonnet-4-20250514”) and

r[“prompt_tokens”] < 100 and r["completion_tokens"] < 100

]

if expensive_calls:

savings = len(expensive_calls) * 0.002 # 估算节省

recommendations.append({

“type”: “model_downgrade”,

“description”: f”发现 {len(expensive_calls)} 次低复杂度调用使用了昂贵模型”,

“recommendation”: “对简单任务(分类、格式化)使用gpt-4o-mini或deepseek-v3”,

“estimated_savings_usd”: savings,

“confidence”: “high”,

})

# 2. 检测prompt是否过于冗长

long_prompts = [r for r in records if r[“prompt_tokens”] > 8000]

if long_prompts:

recommendations.append({

“type”: “prompt_optimization”,

“description”: f”发现 {len(long_prompts)} 次调用的prompt超过8000 token”,

“recommendation”: “精简system prompt,使用上下文压缩,只传必要信息”,

“estimated_savings_usd”: sum(r[“input_cost_usd”] * 0.3 for r in long_prompts),

“confidence”: “medium”,

})

# 3. 检测缓存机会

cache_miss_rate = 1 – (

sum(r[“cached_tokens”] for r in records) /

max(sum(r[“prompt_tokens”] for r in records), 1)

)

if cache_miss_rate > 0.7:

recommendations.append({

“type”: “cache_optimization”,

“description”: f”缓存命中率仅 {(1-cache_miss_rate):.0%},有大量重复prompt”,

“recommendation”: “启用Prompt Caching,将稳定的system prompt放在消息开头”,

“estimated_savings_usd”: cache_miss_rate * 0.4 * sum(r[“input_cost_usd”] for r in records),

“confidence”: “high”,

})

# 4. 检测冗余的LLM调用

tool_selection_calls = [

r for r in records if r.get(“purpose”) == “tool_selection”

]

if len(tool_selection_calls) > len(records) * 0.3:

recommendations.append({

“type”: “reduce_planning_calls”,

“description”: “工具选择调用占比过高”,

“recommendation”: “合并规划和执行步骤,减少中间推理轮次”,

“estimated_savings_usd”: len(tool_selection_calls) * 0.001,

“confidence”: “medium”,

})

return recommendations

`


七、实战:搭建Agent监控Dashboard

现在让我们把前面所有组件串起来,构建一个完整的Agent监控Dashboard。

7.1 整体架构

`

┌──────────────────────────────────────────────────────────────────┐

│ Agent Monitoring Dashboard │

├──────────────────────────────────────────────────────────────────┤

│ │

│ ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ │

│ │ Agent执行引擎 │ │ Token追踪器 │ │ 异常检测器 │ │

│ │ │ │ │ │ │ │

│ │ 推理/工具 │──│ 记录消耗 │──│ 检测异常 │ │

│ │ 调用/检索 │ │ 计算成本 │ │ 触发告警 │ │

│ └──────┬──────┘ └──────┬──────┘ └──────┬──────┘ │

│ │ │ │ │

│ ▼ ▼ ▼ │

│ ┌─────────────────────────────────────────────────┐ │

│ │ AgentObservability │ │

│ │ (统一的可观测性编排层) │ │

│ │ │ │

│ │ ● 汇聚所有信号 │ │

│ │ ● 关联trace/span/metrics │ │

│ │ ● 触发告警 │ │

│ │ ● 生成Dashboard数据 │ │

│ └─────────────────────────────────────────────────┘ │

│ │ │

│ ▼ │

│ ┌─────────────────────────────────────┐ │

│ │ Dashboard API │ │

│ │ GET /overview → 总览数据 │ │

│ │ GET /traces → 追踪列表 │ │

│ │ GET /traces/:id → 追踪详情 │ │

│ │ GET /costs → 成本报告 │ │

│ │ GET /anomalies → 异常列表 │ │

│ │ GET /metrics → 时序指标 │ │

│ └─────────────────────────────────────┘ │

│ │ │

│ ▼ │

│ ┌─────────────────────────────────────┐ │

│ │ Frontend Dashboard (Web UI) │ │

│ └─────────────────────────────────────┘ │

└──────────────────────────────────────────────────────────────────┘

`

7.2 核心实现

`python

import time

import json

from dataclasses import dataclass, field, asdict

from typing import Any, Optional

from collections import defaultdict

from http.server import HTTPServer, BaseHTTPRequestHandler

from urllib.parse import urlparse, parse_qs

class AgentObservability:

“””Agent可观测性——统一编排层”””

def __init__(self, config: dict = None):

config = config or {}

# 初始化各子系统

self.token_tracker = TokenTracker()

self.task_tracker = TaskTracker()

self.loop_detector = LoopDetector(

max_loops=config.get(“max_loops”, 10),

)

self.anomaly_detector = AnomalyDetector()

self.statistical_detector = StatisticalAnomalyDetector(

window_size=config.get(“stat_window”, 100),

z_threshold=config.get(“z_threshold”, 3.0),

)

self.cost_monitor = CostMonitor(

budget=CostBudget(

daily_limit_usd=config.get(“daily_budget”, 10.0),

monthly_limit_usd=config.get(“monthly_budget”, 200.0),

)

)

self.alert_manager = AlertManager()

# 追踪存储

self.traces: dict[str, list[dict]] = defaultdict(list)

self.active_spans: dict[str, dict] = {}

# 指标时间序列

self.metrics_timeseries: dict[str, list[dict]] = defaultdict(list)

# 配置告警通道

self._setup_alerts(config)

# 循环检测回调

self.loop_detector.on_warning(self._on_loop_warning)

def _setup_alerts(self, config: dict):

“””配置告警通道”””

slack_url = config.get(“slack_webhook”)

if slack_url:

self.alert_manager.add_channel(

AnomalySeverity.HIGH,

SlackAlertChannel(slack_url),

)

self.alert_manager.add_channel(

AnomalySeverity.CRITICAL,

SlackAlertChannel(slack_url),

)

webhook_url = config.get(“alert_webhook”)

if webhook_url:

for severity in AnomalySeverity:

self.alert_manager.add_channel(

severity,

WebhookAlertChannel(webhook_url),

)

def _on_loop_warning(self, warning_type: str, details: dict):

“””循环检测警告回调”””

anomaly = Anomaly(

type=AnomalyType.LOOP_DETECTED,

severity=AnomalySeverity.HIGH,

description=f”检测到{warning_type}: {json.dumps(details)}”,

trace_id=details.get(“trace_id”, “”),

span_id=””,

metadata=details,

timestamp=time.time(),

)

self.anomaly_detector.anomalies.append(anomaly)

self.alert_manager.alert(anomaly)

def record_llm_call(

self,

model: str,

prompt_tokens: int,

completion_tokens: int,

cached_tokens: int = 0,

trace_id: str = “”,

purpose: str = “”,

latency_ms: float = 0,

):

“””记录一次LLM调用”””

# Token追踪

cost = self.token_tracker.record(

model=model,

prompt_tokens=prompt_tokens,

completion_tokens=completion_tokens,

cached_tokens=cached_tokens,

trace_id=trace_id,

purpose=purpose,

)

# 成本监控

alerts = self.cost_monitor.record_cost(

amount_usd=cost,

model=model,

feature=purpose,

)

# 触发成本告警

for alert_msg in alerts:

anomaly = Anomaly(

type=AnomalyType.COST_SPIKE,

severity=AnomalySeverity.CRITICAL if “🚨” in alert_msg else AnomalySeverity.HIGH,

description=alert_msg,

trace_id=trace_id,

span_id=””,

metadata={“cost_usd”: cost, “model”: model},

timestamp=time.time(),

)

self.alert_manager.alert(anomaly)

# 统计异常检测

total_tokens = prompt_tokens + completion_tokens

stat_anomaly = self.statistical_detector.add_metric(“tokens_per_call”, total_tokens)

if stat_anomaly:

anomaly = Anomaly(

type=AnomalyType.TOKEN_EXPLOSION,

severity=AnomalySeverity.MEDIUM,

description=f”Token使用量统计异常: z-score={stat_anomaly[‘z_score’]:.2f}”,

trace_id=trace_id,

span_id=””,

metadata=stat_anomaly,

timestamp=time.time(),

)

self.anomaly_detector.anomalies.append(anomaly)

# 记录时间序列

self.metrics_timeseries[“token_usage”].append({

“timestamp”: time.time(),

“value”: total_tokens,

“model”: model,

})

self.metrics_timeseries[“cost”].append({

“timestamp”: time.time(),

“value”: cost,

“model”: model,

})

return cost

def get_overview(self) -> dict:

“””获取Dashboard总览数据”””

cost_data = self.cost_monitor.get_dashboard_data()

task_stats = self.task_tracker.get_completion_stats()

token_summary = self.token_tracker.get_session_summary()

recent_anomalies = [

asdict(a) for a in self.anomaly_detector.anomalies[-20:]

]

return {

“cost”: cost_data,

“tasks”: task_stats,

“tokens”: token_summary,

“anomalies”: recent_anomalies,

“active_traces”: len(self.traces),

}

def get_trace_detail(self, trace_id: str) -> dict | None:

“””获取追踪详情”””

spans = self.traces.get(trace_id)

if not spans:

return None

return {

“trace_id”: trace_id,

“span_count”: len(spans),

“total_duration_ms”: sum(s.get(“duration_ms”, 0) for s in spans),

“total_tokens”: sum(

s.get(“prompt_tokens”, 0) + s.get(“completion_tokens”, 0)

for s in spans

),

“total_cost_usd”: sum(s.get(“cost_usd”, 0) for s in spans),

“spans”: spans,

}

`

7.3 Dashboard API Server

`python

class DashboardHandler(BaseHTTPRequestHandler):

“””Dashboard HTTP API”””

observability: AgentObservability = None # 类变量,由外部注入

def do_GET(self):

parsed = urlparse(self.path)

path = parsed.path

params = parse_qs(parsed.query)

routes = {

“/api/overview”: self._handle_overview,

“/api/traces”: self._handle_traces,

“/api/costs”: self._handle_costs,

“/api/anomalies”: self._handle_anomalies,

“/api/tokens”: self._handle_tokens,

}

handler = routes.get(path)

if handler:

data = handler(params)

self._json_response(200, data)

else:

self._json_response(404, {“error”: “Not found”})

def _handle_overview(self, params):

return self.observability.get_overview()

def _handle_traces(self, params):

limit = int(params.get(“limit”, [“50”])[0])

traces = []

for trace_id, spans in list(self.observability.traces.items())[-limit:]:

traces.append({

“trace_id”: trace_id,

“span_count”: len(spans),

“total_tokens”: sum(

s.get(“prompt_tokens”, 0) + s.get(“completion_tokens”, 0)

for s in spans

),

“total_cost_usd”: sum(s.get(“cost_usd”, 0) for s in spans),

})

return {“traces”: traces}

def _handle_costs(self, params):

return self.observability.cost_monitor.get_dashboard_data()

def _handle_anomalies(self, params):

severity_filter = params.get(“severity”, [None])[0]

anomalies = self.observability.anomaly_detector.anomalies

if severity_filter:

anomalies = [a for a in anomalies if a.severity.value == severity_filter]

return {

“anomalies”: [asdict(a) for a in anomalies[-100:]],

“total”: len(anomalies),

}

def _handle_tokens(self, params):

return self.observability.token_tracker.get_session_summary()

def _json_response(self, status: int, data: dict):

self.send_response(status)

self.send_header(“Content-Type”, “application/json”)

self.send_header(“Access-Control-Allow-Origin”, “*”)

self.end_headers()

self.wfile.write(json.dumps(data, default=str, ensure_ascii=False).encode())

`

7.4 完整使用示例

`python

def demo_observability():

“””可观测性系统演示”””

# 1. 初始化

obs = AgentObservability(config={

“daily_budget”: 5.0,

“monthly_budget”: 100.0,

“max_loops”: 10,

“slack_webhook”: “https://hooks.slack.com/xxx”,

})

# 2. 模拟Agent执行

trace_id = “trace-demo-001”

# 模拟规划阶段

cost = obs.record_llm_call(

model=”gpt-4o”,

prompt_tokens=1200,

completion_tokens=300,

trace_id=trace_id,

purpose=”planning”,

latency_ms=800,

)

print(f”规划阶段成本: ${cost:.6f}”)

# 模拟检索阶段(无LLM调用,只记录追踪)

obs.traces[trace_id].append({

“span_type”: “retrieval”,

“duration_ms”: 200,

“docs_retrieved”: 5,

})

# 模拟工具调用

obs.traces[trace_id].append({

“span_type”: “tool_call”,

“tool_name”: “search_database”,

“duration_ms”: 500,

“status”: “success”,

})

# 模拟综合阶段

cost = obs.record_llm_call(

model=”gpt-4o”,

prompt_tokens=3500,

completion_tokens=800,

cached_tokens=1200,

trace_id=trace_id,

purpose=”synthesis”,

latency_ms=1200,

)

print(f”综合阶段成本: ${cost:.6f}”)

# 3. 记录任务完成

obs.task_tracker.start_task(“task-001”, “qa”, “用户问题:…”)

obs.task_tracker.complete_task(“task-001”, success=True)

obs.task_tracker.add_metrics(“task-001”, {“trace_id”: trace_id})

# 4. 获取Dashboard数据

overview = obs.get_overview()

print(“n=== Dashboard Overview ===”)

print(json.dumps(overview, indent=2, default=str, ensure_ascii=False))

# 5. 启动Dashboard API(可选)

# DashboardHandler.observability = obs

# server = HTTPServer((“0.0.0.0”, 8080), DashboardHandler)

# print(“Dashboard API running at http://localhost:8080”)

# server.serve_forever()

if __name__ == “__main__”:

demo_observability()

`


八、常见陷阱(过度日志、性能开销)

8.1 陷阱一:过度日志——”一切皆可观测”的幻觉

最常见的陷阱是试图记录一切。我见过一个生产环境的Agent系统,每天产生超过500GB的日志——因为它们记录了每一轮推理的完整prompt和response,包括重复的system prompt。

后果: 存储成本飙升,日志查询变慢(在海量日志中找到有用信息变成大海捞针),而且大量日志中混杂了敏感的用户数据,带来合规风险。

正确做法:

`python

class SmartLogger:

“””智能日志器——根据场景自动调整粒度”””

def __init__(self):

self.default_granularity = LogGranularity.DETAIL

self.override_rules: list[dict] = []

def add_rule(self, condition, granularity: LogGranularity):

“””添加覆盖规则:满足条件时使用指定粒度”””

self.override_rules.append({

“condition”: condition,

“granularity”: granularity,

})

def get_granularity(self, context: dict) -> LogGranularity:

“””根据上下文决定日志粒度”””

# 1. 错误发生时自动升级到TRACE级别

if context.get(“has_error”):

return LogGranularity.TRACE

# 2. 高成本操作记录DETAIL级别

if context.get(“total_tokens”, 0) > 10000:

return LogGranularity.DETAIL

# 3. 低token消耗的简单操作只记SUMMARY

if context.get(“total_tokens”, 0) < 500:

return LogGranularity.SUMMARY

# 4. 检查覆盖规则

for rule in self.override_rules:

if rule“condition”:

return rule[“granularity”]

return self.default_granularity

`

关键原则:正常路径记摘要,异常路径记详情,调试模式记一切。

8.2 陷阱二:性能开销——观测本身拖慢了Agent

可观测性组件不应该成为Agent的性能瓶颈。但我见过因为同步写日志导致Agent延迟增加30%的案例。

正确做法:

`python

import asyncio

from collections import deque

from threading import Thread

class AsyncObservabilityWriter:

“””异步可观测性写入器——核心原则:采集同步,写入异步”””

def __init__(self, flush_interval: float = 5.0, batch_size: int = 100):

self.buffer: deque = deque(maxlen=10000)

self.flush_interval = flush_interval

self.batch_size = batch_size

self._running = False

self._worker_thread: Thread | None = None

# 性能统计

self.stats = {

“events_buffered”: 0,

“events_flushed”: 0,

“events_dropped”: 0,

“avg_flush_latency_ms”: 0,

}

def start(self):

“””启动后台写入线程”””

self._running = True

self._worker_thread = Thread(target=self._flush_loop, daemon=True)

self._worker_thread.start()

def stop(self):

“””停止并刷新剩余数据”””

self._running = False

if self._worker_thread:

self._worker_thread.join(timeout=10)

self._flush_all()

def emit(self, event: dict):

“””

采集事件——这个方法必须极快(微秒级)

它只做一件事:把事件放入内存队列

“””

if len(self.buffer) >= self.buffer.maxlen:

self.stats[“events_dropped”] += 1

return # 队列满时丢弃(宁可丢数据也不能阻塞Agent)

self.buffer.append({

**event,

“_emitted_at”: time.time(),

})

self.stats[“events_buffered”] += 1

def _flush_loop(self):

“””后台刷新循环”””

while self._running:

time.sleep(self.flush_interval)

self._flush_all()

def _flush_all(self):

“””批量刷新缓冲区中的事件”””

if not self.buffer:

return

batch = []

while self.buffer and len(batch) < self.batch_size:

batch.append(self.buffer.popleft())

if batch:

start = time.time()

self._write_batch(batch)

latency = (time.time() – start) * 1000

self.stats[“events_flushed”] += len(batch)

# 更新平均延迟

n = self.stats[“events_flushed”]

old_avg = self.stats[“avg_flush_latency_ms”]

self.stats[“avg_flush_latency_ms”] = old_avg + (latency – old_avg) / n

def _write_batch(self, batch: list[dict]):

“””实际写入——子类实现具体存储逻辑”””

# 可以写入文件、发送到Kafka、写入ClickHouse等

pass

`

关键原则:采集在当前线程(微秒级),写入在后台线程(毫秒级),绝不阻塞Agent主流程。

8.3 陷阱三:只观测LLM,不观测工具

很多团队只追踪LLM调用(因为LLM最贵),却忽略了工具调用的可观测性。但Agent的问题往往出在工具层——工具超时、返回错误数据、参数传递错误。

正确做法: 工具调用的追踪粒度应该与LLM调用一致:

`python

def tracked_tool_call(tool_name: str, tool_func, *args, **kwargs):

“””带追踪的工具调用包装器”””

start_time = time.time()

tool_input = {

“args”: [str(a)[:200] for a in args], # 截断长参数

“kwargs”: {k: str(v)[:200] for k, v in kwargs.items()},

}

try:

result = tool_func(*args, **kwargs)

duration_ms = (time.time() – start_time) * 1000

tool_output = str(result)[:500] # 截断长输出

# 记录成功的工具调用

log_entry = {

“span_type”: “tool_call”,

“tool_name”: tool_name,

“tool_input”: tool_input,

“tool_output”: tool_output,

“tool_duration_ms”: duration_ms,

“tool_status”: “success”,

“output_length”: len(str(result)),

}

# 如果工具返回为空或异常短,标记警告

if len(str(result).strip()) < 5:

log_entry[“warning”] = “tool_returned_empty”

return result

except TimeoutError:

duration_ms = (time.time() – start_time) * 1000

log_entry = {

“span_type”: “tool_call”,

“tool_name”: tool_name,

“tool_input”: tool_input,

“tool_duration_ms”: duration_ms,

“tool_status”: “timeout”,

}

raise

except Exception as e:

duration_ms = (time.time() – start_time) * 1000

log_entry = {

“span_type”: “tool_call”,

“tool_name”: tool_name,

“tool_input”: tool_input,

“tool_duration_ms”: duration_ms,

“tool_status”: “error”,

“error_type”: type(e).__name__,

“error_message”: str(e)[:200],

}

raise

`

8.4 陷阱四:忽视用户反馈回路

纯技术指标(token、延迟、工具调用)无法衡量Agent的真正价值——用户是否满意。没有用户反馈回路的可观测性系统是”半盲”的。

`python

class FeedbackCollector:

“””用户反馈收集器”””

def __init__(self):

self.feedback_store: list[dict] = []

def collect(

self,

session_id: str,

trace_id: str,

rating: int, # 1-5

feedback_type: str, # “thumbs_up”, “thumbs_down”, “star_rating”

comment: str = “”,

category: str = “”, # “helpful”, “accurate”, “harmful”, “irrelevant”

):

self.feedback_store.append({

“session_id”: session_id,

“trace_id”: trace_id,

“rating”: rating,

“feedback_type”: feedback_type,

“comment”: comment,

“category”: category,

“timestamp”: time.time(),

})

def correlate_with_traces(self) -> list[dict]:

“””将用户反馈与追踪数据关联,找出低分对话的共同模式”””

low_rated = [f for f in self.feedback_store if f[“rating”] <= 2]

high_rated = [f for f in self.feedback_store if f[“rating”] >= 4]

return {

“low_rated_count”: len(low_rated),

“high_rated_count”: len(high_rated),

“satisfaction_rate”: len(high_rated) / max(len(self.feedback_store), 1),

“common_complaints”: self._extract_patterns(low_rated),

}

def _extract_patterns(self, feedbacks: list[dict]) -> dict:

“””从低分反馈中提取常见模式”””

categories = defaultdict(int)

for f in feedbacks:

if f[“category”]:

categories[f[“category”]] += 1

return dict(sorted(categories.items(), key=lambda x: x[1], reverse=True))

`

8.5 陷阱五:观测数据没有闭环

采集了大量数据却从不回看,等于没有采集。观测数据必须形成闭环:采集 → 分析 → 发现 → 优化 → 验证 → 采集…

`python

class ObservabilityFeedbackLoop:

“””可观测性闭环——从数据到行动”””

def __init__(self, observability: AgentObservability):

self.obs = observability

self.optimization_history: list[dict] = []

def weekly_review(self) -> dict:

“””每周自动生成优化报告”””

overview = self.obs.get_overview()

report = {

“period”: “last_7_days”,

“summary”: {

“total_cost_usd”: overview[“cost”][“monthly_cost_usd”],

“task_success_rate”: overview[“tasks”][“success_rate”],

“total_anomalies”: len(overview[“anomalies”]),

},

“top_issues”: [],

“recommendations”: [],

}

# 分析最常见的异常

anomaly_types = defaultdict(int)

for a in overview[“anomalies”]:

anomaly_types[a[“type”]] += 1

for atype, count in sorted(anomaly_types.items(), key=lambda x: x[1], reverse=True)[:5]:

report[“top_issues”].append({

“type”: atype,

“count”: count,

“trend”: “increasing”, # 需要与上周对比

})

# 生成优化建议

optimizer = CostOptimizer(self.obs.token_tracker)

report[“recommendations”] = optimizer.analyze_and_recommend()

return report

def apply_optimization(self, recommendation: dict) -> bool:

“””应用一条优化建议”””

opt_type = recommendation[“type”]

if opt_type == “model_downgrade”:

# 更新路由规则:简单任务使用更便宜的模型

print(f”[OPTIMIZE] 对简单任务切换到更便宜的模型”)

self.optimization_history.append({

“recommendation”: recommendation,

“applied_at”: time.time(),

“status”: “applied”,

})

return True

elif opt_type == “cache_optimization”:

print(f”[OPTIMIZE] 启用Prompt Caching”)

self.optimization_history.append({

“recommendation”: recommendation,

“applied_at”: time.time(),

“status”: “applied”,

})

return True

return False

`


九、总结

9.1 核心要点回顾

回顾本文的关键设计原则:

1. 可观测性三大支柱在Agent场景下的演化

  • 日志从”请求/响应”变成”推理链/决策过程”
  • 指标从”QPS/延迟”变成”Token/循环次数/任务完成率”
  • 追踪从”跨服务调用”变成”跨步骤决策链”

2. Agent特有指标的四大类别

  • 成本指标:Token消耗、API成本、缓存命中率
  • 质量指标:任务完成率、幻觉分数、用户评分
  • 效率指标:循环次数、工具调用效率、首token延迟
  • 安全指标:护栏触发、PII检测、权限越级

3. 分布式追踪方案选择

  • Langfuse:开源首选,可自部署
  • LangSmith:LangChain生态最优解
  • Helicone:零侵入快速接入

4. 异常检测的两层架构

  • 规则层:确定性异常(超时、token超限、循环检测)
  • 统计层:概率性异常(Z-score、趋势分析)

5. 成本监控是Agent的”生命线”

  • 预算制:日/月/会话/请求四级预算
  • 实时告警:70%警告,90%紧急
  • 优化闭环:采集→分析→建议→应用→验证

6. 五大常见陷阱

  • 过度日志:分层记录,正常摘要/异常详情
  • 性能开销:采集同步/写入异步,绝不阻塞主流程
  • 忽视工具:工具追踪粒度=LLM追踪粒度
  • 忽视反馈:用户满意度是最核心的质量指标
  • 没有闭环:数据必须驱动优化,否则等于没采集

9.2 设计检查清单

在为你的Agent系统设计可观测性时,用这个清单自查:

`

□ 日志:是否实现了分层记录(TRACE/DETAIL/SUMMARY)?

□ 日志:是否在异常时自动提升日志粒度?

□ 指标:是否覆盖了Token、工具、循环、任务完成率四大类?

□ 追踪:是否能可视化一次Agent执行的完整决策链?

□ 成本:是否设置了多级预算和实时告警?

□ 异常:是否同时有规则检测和统计检测?

□ 告警:是否实现了告警抑制(避免告警风暴)?

□ 反馈:是否采集了用户满意度数据并关联到追踪?

□ 性能:观测写入是否异步化,不阻塞Agent主流程?

□ 闭环:是否有定期的数据回顾和优化机制?

`

9.3 推荐技术栈

根据团队规模和阶段,推荐不同的技术栈:

个人开发者/小团队(起步阶段):

  • 追踪:Langfuse(免费自部署)
  • 指标:Prometheus + Grafana
  • 日志:结构化JSON日志 + Loki
  • 成本:手动Token追踪 + 简单脚本

中型团队(生产阶段):

  • 追踪:Langfuse / LangSmith
  • 指标:Datadog / Grafana Cloud
  • 日志:ELK Stack / ClickHouse
  • 成本:自建Dashboard + 自动告警

大型企业(规模化阶段):

  • 追踪:自建 + 商业方案混合
  • 指标:自建Prometheus集群 + 商业APM
  • 日志:ClickHouse / BigQuery
  • 成本:FinOps平台集成 + 自动优化

9.4 最后的话

可观测性不是锦上添花,而是Agent系统的生存基础设施。没有可观测性的Agent就像没有仪表盘的飞机——可能飞得很好,但你永远不知道它什么时候会出问题,出了问题也找不到原因。

投入10%的工程时间在可观测性上,能节省90%的排障时间。这不是夸张——在Agent这种非确定性系统中,这是被反复验证的真理。

本篇是Harness Engineering系列的完结篇。至此,我们已经完整地拆解了Harness六层架构的每一层设计。希望这个系列能帮助你系统性地理解和实践AI Agent工程。


系列导航

>

第一篇:什么是Harness Engineering

第二篇:六层架构详解

第三篇:Prompt层设计

第四篇:Context层设计

第五篇:Tool层设计

第六篇:Memory层设计

第七篇:Eval层设计

第八篇:Observability层设计(本文)

评论

发表回复

您的邮箱地址不会被公开。 必填项已用 * 标注