一、引言:可观测性是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工程。
系列导航
>
– 第八篇:Observability层设计(本文)
发表回复