Agent 系统的行为往往是非确定性的、多步骤的、涉及复杂工具调用的,这让传统的调试方法变得力不从心。本文探讨如何构建 Agent 系统的可观测性体系,让"不可预测"变得可控可追踪。
一、Agent 系统的可观测性挑战
1.1 为什么 Agent 难以调试?
传统软件的行为是确定性的:相同的输入经过相同的处理逻辑,必然产生相同的输出。但 Agent 系统打破了这一确定性神话。
┌─────────────────────────────────────────────────────────────────────────────┐
│ Agent 系统的不确定性来源 │
├─────────────────────────────────────────────────────────────────────────────┤
│ │
│ ┌─────────────────────────────────────────────────────────────────────┐ │
│ │ 1. LLM 的固有不确定性 │ │
│ │ │ │
│ │ 同一 Prompt → 不同模型可能产生不同回复 │ │
│ │ 相同输入 → 即使同模型也可能有随机性(temperature > 0) │ │
│ │ 上下文变化 → 微小变化可能导致截然不同的决策 │ │
│ └─────────────────────────────────────────────────────────────────────┘ │
│ ↓ │
│ ┌─────────────────────────────────────────────────────────────────────┐ │
│ │ 2. 多步推理链路 │ │
│ │ │ │
│ │ Agent 执行往往需要多轮交互: │ │
│ │ 用户请求 → 意图识别 → 工具选择 → 参数生成 → 工具执行 → 结果处理 │ │
│ │ → 可能触发新的工具调用 → 最终回复 │ │
│ │ │ │
│ │ 任何一步出错都可能导致最终结果偏离预期 │ │
│ └─────────────────────────────────────────────────────────────────────┘ │
│ ↓ │
│ ┌─────────────────────────────────────────────────────────────────────┐ │
│ │ 3. 外部工具的不确定性 │ │
│ │ │ │
│ │ API 响应可能超时、返回错误数据 │ │
│ │ 网络波动导致间歇性失败 │ │
│ │ 第三方服务行为可能变化 │ │
│ │ 工具返回格式可能与预期不符 │ │
│ └─────────────────────────────────────────────────────────────────────┘ │
│ ↓ │
│ ┌─────────────────────────────────────────────────────────────────────┐ │
│ │ 4. 组合爆炸 │ │
│ │ │ │
│ │ 10 个工具 × 10 种参数组合 × 5 种状态 = 500 种可能路径 │ │
│ │ 难以穷尽测试所有场景 │ │
│ └─────────────────────────────────────────────────────────────────────┘ │
│ │
└─────────────────────────────────────────────────────────────────────────────┘1.2 可观测性的三大支柱
可观测性(Observability)≠ 监控(Monitoring)。传统监控告诉你"系统坏了",可观测性告诉你"为什么坏了"。
┌─────────────────────────────────────────────────────────────────────────────┐
│ 可观测性三大支柱 │
├─────────────────────────────────────────────────────────────────────────────┤
│ │
│ ┌───────────────┐ │
│ │ Traces │ │
│ │ 请求链路追踪 │ │
│ └───────────────┘ │
│ │ │
│ │ │
│ ┌───────────────┐ ┌───────────────┐ ┌───────────────┐ │
│ │ Metrics │ ←→ │ Traces │ ←→ │ Logs │ │
│ │ 可聚合指标 │ │ 完整链路 │ │ 离散事件 │ │
│ └───────────────┘ └───────────────┘ └───────────────┘ │
│ ↑ ↑ ↑ │
│ │ │ │ │
│ ┌────────────┐ ┌────────────┐ ┌────────────┐ │
│ │ Prometheus │ │ Jaeger │ │ ELK │ │
│ │ Grafana │ │ Zipkin │ │ Loki │ │
│ └────────────┘ └────────────┘ └────────────┘ │
│ │
└─────────────────────────────────────────────────────────────────────────────┘二、LangSmith 与 LangFuse 集成
2.1 LangSmith 集成
LangSmith 是 LangChain 官方推出的可观测性平台,深度集成 LangChain。
import os
from langchain_openai import ChatOpenAI
from langchain.prompts import ChatPromptTemplate
from langchain.schema import StrOutputParser
from langchain.callbacks.tracers.langsmith import LangSmithTracer
from langchain.callbacks.manager import CallbackManager
from langchain.callbacks.base import BaseCallbackHandler
# ============== LangSmith 配置 ==============
# 设置环境变量
os.environ["LANGCHAIN_TRACING_V2"] = "true"
os.environ["LANGCHAIN_API_KEY"] = "your-langsmith-api-key"
os.environ["LANGCHAIN_PROJECT"] = "agent-debugging" # 项目名称
# 可选:添加数据集用于评估
os.environ["LANGCHAIN_DATASET"] = "agent-evaluation-dataset"
class AgentCallbackHandler(BaseCallbackHandler):
"""
自定义回调处理器:记录 Agent 特定事件
可以在这里:
1. 记录自定义指标
2. 添加业务相关的上下文
3. 实现自定义告警逻辑
"""
def __init__(self):
self.agent_decisions = []
self.tool_calls = []
self.token_usage = {"prompt": 0, "completion": 0, "total": 0}
def on_llm_start(self, serialized, prompts, **kwargs):
"""LLM 开始调用"""
print(f"[Callback] LLM 开始处理,Prompt tokens: {len(prompts)}")
def on_llm_end(self, response, **kwargs):
"""LLM 调用结束"""
if hasattr(response, "usage_metadata"):
usage = response.usage_metadata
self.token_usage["prompt"] += usage.get("input_tokens", 0)
self.token_usage["completion"] += usage.get("output_tokens", 0)
self.token_usage["total"] += usage.get("total_tokens", 0)
def on_chain_start(self, serialized, inputs, **kwargs):
"""Chain 开始执行"""
chain_name = serialized.get("name", "unknown")
print(f"[Callback] Chain 开始: {chain_name}")
def on_chain_end(self, outputs, **kwargs):
"""Chain 执行结束"""
print(f"[Callback] Chain 结束,输出: {str(outputs)[:200]}")
def on_tool_start(self, serialized, input_str, **kwargs):
"""工具开始执行"""
tool_name = serialized.get("name", "unknown")
print(f"[Callback] 工具开始: {tool_name}, 输入: {input_str[:100]}")
self.tool_calls.append({
"tool": tool_name,
"input": input_str,
"status": "started"
})
def on_tool_end(self, output, **kwargs):
"""工具执行结束"""
self.tool_calls[-1]["status"] = "completed"
self.tool_calls[-1]["output"] = str(output)[:200]
print(f"[Callback] 工具结束: {self.tool_calls[-1]['tool']}, 输出: {str(output)[:100]}")
def on_tool_error(self, error, **kwargs):
"""工具执行出错"""
self.tool_calls[-1]["status"] = "error"
self.tool_calls[-1]["error"] = str(error)
print(f"[Callback] 工具错误: {error}")
def create_langsmith_callback_manager(project_name: str = "agent-debugging"):
"""创建 LangSmith 回调管理器"""
# 创建 LangSmith tracer
tracer = LangSmithTracer(
project_name=project_name,
# 可选:添加 tags 用于过滤
tags=["production", "v1.0"],
# 可选:添加 metadata
metadata={
"environment": os.getenv("ENV", "development"),
"version": os.getenv("APP_VERSION", "1.0.0")
}
)
# 创建自定义 handler
custom_handler = AgentCallbackHandler()
# 创建 callback manager
callback_manager = CallbackManager(
handlers=[tracer, custom_handler]
)
return callback_manager, custom_handler
# 使用示例
async def demo_langsmith():
"""LangSmith 集成示例"""
# 创建带回调的 LLM
callback_manager, handler = create_langsmith_callback_manager()
llm = ChatOpenAI(
model="gpt-4o",
temperature=0,
callback_manager=callback_manager
)
# 创建简单的 chain
prompt = ChatPromptTemplate.from_template(
"用一句话解释 {topic},用中文回答"
)
chain = prompt | llm | StrOutputParser()
# 执行
result = await chain.ainvoke({"topic": "量子计算"})
print(f"\n最终结果: {result}")
print(f"Token 消耗: {handler.token_usage}")
print(f"工具调用: {handler.tool_calls}")2.2 LangFuse 集成
LangFuse 是开源的可观测性方案,支持自部署,数据完全私有。
"""
LangFuse 集成:开源可观测性方案
"""
from langfuse import Langfuse
from langfuse.callback import CallbackHandler
from langchain_openai import ChatOpenAI
from langchain.schema import HumanMessage
import os
# ============== LangFuse 配置 ==============
# 方式1:使用云服务
# os.environ["LANGFUSE_PUBLIC_KEY"] = "pk-xxx"
# os.environ["LANGFUSE_SECRET_KEY"] = "sk-xxx"
# langfuse = Langfuse()
# 方式2:自部署
langfuse = Langfuse(
public_key=os.getenv("LANGFUSE_PUBLIC_KEY", ""),
secret_key=os.getenv("LANGFUSE_SECRET_KEY", ""),
host=os.getenv("LANGFUSE_HOST", "http://localhost:3000") # 自部署地址
)
def create_langfuse_callback():
"""创建 LangFuse 回调处理器"""
return CallbackHandler(
langfuse=langfuse,
# 会话相关
session_id="session-123", # 用于追踪同一会话
user_id="user-456", # 用于追踪用户
# 追踪元数据
tags=["agent", "production"],
metadata={
"agent_version": "1.0.0",
"environment": "production"
}
)
# ============== 追踪 Agent 执行 ==============
class AgentTracer:
"""
Agent 追踪器:封装 LangFuse 功能
功能:
1. 追踪每个 Agent 决策
2. 记录工具调用链路
3. 收集性能指标
"""
def __init__(self, agent_name: str):
self.agent_name = agent_name
self.langfuse = langfuse
self.trace = None
self.spans = []
def start_trace(self, name: str, input_data: dict) -> str:
"""
开始追踪
Returns:
trace_id: 用于关联后续操作
"""
self.trace = self.langfuse.trace(
name=f"{self.agent_name}.{name}",
input=input_data,
metadata={
"agent": self.agent_name,
"type": "agent_execution"
}
)
return self.trace.id
def create_span(self, name: str, metadata: dict = None):
"""创建子跨度"""
class SpanContext:
def __init__(self, tracer, name, metadata):
self.tracer = tracer
self.name = name
self.metadata = metadata or {}
self.start_time = None
self.end_time = None
self.observations = []
def __enter__(self):
import time
self.start_time = time.time()
return self
def __exit__(self, exc_type, exc_val, exc_tb):
import time
self.end_time = time.time()
duration = self.end_time - self.start_time
# 记录到 LangFuse
if self.tracer.trace:
self.tracer.trace.span(
name=self.name,
input=self.metadata.get("input"),
output=self.metadata.get("output"),
metadata={
**self.metadata,
"duration_seconds": duration,
"status": "success" if exc_type is None else "error"
}
)
return False
return SpanContext(self, name, metadata)
def record_decision(
self,
step: str,
reasoning: str,
action: str,
confidence: float
) -> None:
"""记录 Agent 决策点"""
if not self.trace:
return
self.trace.generation(
name=f"decision.{step}",
input={"step": step, "reasoning": reasoning},
output={
"action": action,
"confidence": confidence
},
metadata={
"type": "agent_decision",
"step": step
}
)
def record_tool_call(
self,
tool_name: str,
input_args: dict,
output: Any,
error: str = None
) -> None:
"""记录工具调用"""
if not self.trace:
return
self.trace.span(
name=f"tool.{tool_name}",
input=input_args,
output=output if not error else None,
metadata={
"type": "tool_call",
"tool_name": tool_name,
"error": error
}
)
def end_trace(self, output: Any, metadata: dict = None) -> None:
"""结束追踪"""
if not self.trace:
return
self.trace.update(
output=output,
metadata=metadata or {}
)
self.trace = None
# 使用示例
async def demo_langfuse():
"""LangFuse 使用示例"""
tracer = AgentTracer("weather_assistant")
# 开始追踪
trace_id = tracer.start_trace(
"get_weather_and_remind",
{"user_input": "明天北京天气怎么样?"}
)
# 记录决策
tracer.record_decision(
step="intent_detection",
reasoning="用户询问天气,需要调用天气查询工具",
action="call_get_weather",
confidence=0.95
)
# 记录工具调用
with tracer.create_span("get_weather", {"city": "北京", "date": "明天"}):
weather_result = {"temp": 15, "condition": "晴"}
tracer.record_tool_call(
"get_weather",
{"city": "北京", "date": "明天"},
weather_result
)
# 结束追踪
tracer.end_trace(
output={"answer": "北京明天天气晴朗,温度15°C"},
metadata={"tokens_used": 500}
)三、Trace 设计:每步决策可追踪
3.1 完整 Trace 架构
┌─────────────────────────────────────────────────────────────────────────────┐
│ Agent Trace 完整架构 │
├─────────────────────────────────────────────────────────────────────────────┤
│ │
│ 用户输入: "帮我查下张三的订单,如果已发货就设置个提醒" │
│ │
│ ┌─────────────────────────────────────────────────────────────────────┐ │
│ │ Trace Root │ │
│ │ trace_id: abc123 │ │
│ │ start_time: 2024-01-15 10:00:00 │ │
│ │ duration: 2.5s │ │
│ └─────────────────────────────────────────────────────────────────────┘ │
│ │ │
│ ┌──────────────────────────┼──────────────────────────┐ │
│ ↓ ↓ ↓ │
│ ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ │
│ │ Span 1 │ │ Span 2 │ │ Span 3 │ │
│ │ 意图识别 │ │ 工具调用 │ │ 工具调用 │ │
│ │ intent.detect│ │ query.order │ │ set.reminder│ │
│ └─────────────┘ └─────────────┘ └─────────────┘ │
│ │ │ │ │
│ ↓ ↓ ↓ │
│ ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ │
│ │ Generation │ │ Event │ │ Event │ │
│ │ 决策理由 │ │ 工具输入 │ │ 工具输出 │ │
│ │ 置信度0.92 │ │ order_id │ │ 已发货 │ │
│ └─────────────┘ └─────────────┘ └─────────────┘ │
│ │
│ 嵌套关系: │
│ Trace │
│ └── Span: 意图识别 │
│ └── Generation: 推理过程 │
│ └── Span: 订单查询 │
│ └── Event: API调用详情 │
│ └── Span: 设置提醒 │
│ └── Event: 提醒创建成功 │
│ │
└─────────────────────────────────────────────────────────────────────────────┘3.2 追踪实现代码
"""
Agent 追踪系统完整实现
支持:LangSmith、LangFuse、自建追踪
"""
import json
import time
import uuid
from typing import Any, Dict, List, Optional, Callable
from dataclasses import dataclass, field
from datetime import datetime
from enum import Enum
from contextvars import ContextVar
from collections import defaultdict
# 追踪上下文
_current_trace: ContextVar[Optional["TraceContext"]] = ContextVar(
"current_trace",
default=None
)
class SpanKind(Enum):
"""跨度类型"""
INTERNAL = "internal" # 内部操作
LLM = "llm" # LLM 调用
TOOL = "tool" # 工具调用
CHAIN = "chain" # Chain 执行
RETRIEVER = "retriever" # 检索操作
class SpanStatus(Enum):
"""跨度状态"""
OK = "ok"
ERROR = "error"
UNSET = "unset"
@dataclass
class Span:
"""追踪跨度"""
span_id: str
name: str
kind: SpanKind
trace_id: str
parent_span_id: Optional[str]
start_time: datetime = field(default_factory=datetime.now)
end_time: Optional[datetime] = None
input_data: Any = None
output_data: Any = None
error: Optional[str] = None
metadata: Dict[str, Any] = field(default_factory=dict)
attributes: Dict[str, Any] = field(default_factory=dict)
# 关联数据
tags: List[str] = field(default_factory=list)
@property
def duration_ms(self) -> float:
"""计算持续时间(毫秒)"""
if self.end_time:
return (self.end_time - self.start_time).total_seconds() * 1000
return 0.0
def to_dict(self) -> Dict[str, Any]:
"""转换为字典"""
return {
"span_id": self.span_id,
"name": self.name,
"kind": self.kind.value,
"trace_id": self.trace_id,
"parent_span_id": self.parent_span_id,
"start_time": self.start_time.isoformat(),
"end_time": self.end_time.isoformat() if self.end_time else None,
"duration_ms": self.duration_ms,
"input": self._truncate(self.input_data),
"output": self._truncate(self.output_data),
"error": self.error,
"metadata": self.metadata,
"attributes": self.attributes,
"tags": self.tags
}
def _truncate(self, data: Any, max_len: int = 1000) -> Any:
"""截断数据"""
if data is None:
return None
s = str(data)
if len(s) > max_len:
return s[:max_len] + "..."
return data
@dataclass
class TraceContext:
"""追踪上下文"""
trace_id: str
name: str
start_time: datetime = field(default_factory=datetime.now)
end_time: Optional[datetime] = None
spans: List[Span] = field(default_factory=list)
_span_stack: List[str] = field(default_factory=list) # 用于嵌套追踪
metadata: Dict[str, Any] = field(default_factory=dict)
tags: List[str] = field(default_factory=list)
# 指标收集
metrics: Dict[str, Any] = field(default_factory=lambda: defaultdict(int))
@property
def duration_ms(self) -> float:
"""总持续时间"""
if self.end_time:
return (self.end_time - self.start_time).total_seconds() * 1000
return (datetime.now() - self.start_time).total_seconds() * 1000
def add_span(self, span: Span) -> None:
"""添加跨度"""
self.spans.append(span)
def get_current_span(self) -> Optional[Span]:
"""获取当前活跃跨度"""
if not self._span_stack:
return None
span_id = self._span_stack[-1]
for span in reversed(self.spans):
if span.span_id == span_id:
return span
return None
def to_dict(self) -> Dict[str, Any]:
"""转换为字典"""
return {
"trace_id": self.trace_id,
"name": self.name,
"start_time": self.start_time.isoformat(),
"end_time": self.end_time.isoformat() if self.end_time else None,
"duration_ms": self.duration_ms,
"spans": [s.to_dict() for s in self.spans],
"metadata": self.metadata,
"tags": self.tags,
"metrics": dict(self.metrics)
}
class AgentTracer:
"""
Agent 追踪器:完整的追踪系统实现
功能:
1. 自动追踪 LLM 调用
2. 自动追踪工具执行
3. 支持嵌套跨度
4. 收集关键指标
5. 支持多后端导出(LangSmith、LangFuse、自建)
"""
def __init__(
self,
service_name: str,
export_backends: Optional[List[Callable]] = None
):
self.service_name = service_name
self.export_backends = export_backends or []
def start_trace(
self,
name: str,
input_data: Any = None,
metadata: Optional[Dict] = None,
tags: Optional[List[str]] = None
) -> TraceContext:
"""
开始追踪会话
Args:
name: 追踪名称
input_data: 输入数据
metadata: 元数据
tags: 标签
Returns:
TraceContext: 追踪上下文
"""
trace = TraceContext(
trace_id=str(uuid.uuid4())[:16],
name=name,
metadata=metadata or {},
tags=tags or []
)
# 存入上下文
_current_trace.set(trace)
# 创建根跨度
root_span = Span(
span_id=str(uuid.uuid4())[:16],
name=name,
kind=SpanKind.INTERNAL,
trace_id=trace.trace_id,
parent_span_id=None,
input_data=input_data
)
trace.add_span(root_span)
trace._span_stack.append(root_span.span_id)
return trace
def end_trace(
self,
output_data: Any = None,
error: Optional[str] = None,
metadata: Optional[Dict] = None
) -> Optional[TraceContext]:
"""结束追踪"""
trace = _current_trace.get()
if not trace:
return None
trace.end_time = datetime.now()
# 更新根跨度
root_span = trace.spans[0]
root_span.end_time = trace.end_time
root_span.output_data = output_data
root_span.error = error
# 更新元数据
if metadata:
trace.metadata.update(metadata)
# 导出到后端
for backend in self.export_backends:
try:
backend(trace)
except Exception as e:
print(f"导出追踪数据失败: {e}")
# 清理上下文
_current_trace.set(None)
return trace
def create_span(
self,
name: str,
kind: SpanKind,
input_data: Any = None,
metadata: Optional[Dict] = None
) -> Span:
"""创建新的跨度"""
trace = _current_trace.get()
if not trace:
# 没有活跃追踪,创建一个匿名的
trace = self.start_trace(f"implicit.{name}")
# 获取父跨度
parent_span = trace.get_current_span()
parent_span_id = parent_span.span_id if parent_span else None
# 创建新跨度
span = Span(
span_id=str(uuid.uuid4())[:16],
name=name,
kind=kind,
trace_id=trace.trace_id,
parent_span_id=parent_span_id,
input_data=input_data,
metadata=metadata or {}
)
trace.add_span(span)
trace._span_stack.append(span.span_id)
return span
def end_span(
self,
span: Span,
output_data: Any = None,
error: Optional[str] = None,
attributes: Optional[Dict] = None
) -> None:
"""结束跨度"""
span.end_time = datetime.now()
span.output_data = output_data
span.error = error
if attributes:
span.attributes.update(attributes)
# 更新追踪上下文
trace = _current_trace.get()
if trace and span.span_id in trace._span_stack:
trace._span_stack.remove(span.span_id)
# 更新指标
if span.kind == SpanKind.LLM and output_data:
trace.metrics["llm_calls"] += 1
if hasattr(output_data, "usage_metadata"):
usage = output_data.usage_metadata
trace.metrics["prompt_tokens"] += usage.get("input_tokens", 0)
trace.metrics["completion_tokens"] += usage.get("output_tokens", 0)
if span.kind == SpanKind.TOOL:
trace.metrics["tool_calls"] += 1
def span(self, name: str, kind: SpanKind = SpanKind.INTERNAL):
"""
上下文管理器:自动管理跨度生命周期
用法:
with tracer.span("query_order", SpanKind.TOOL) as span:
result = query_order(order_id)
span.output_data = result
"""
class SpanContext:
def __init__(ctx, tracer, name, kind):
ctx.tracer = tracer
ctx.name = name
ctx.kind = kind
ctx.span = None
def __enter__(ctx):
ctx.span = ctx.tracer.create_span(
name=ctx.name,
kind=ctx.kind
)
return ctx.span
def __exit__(ctx, exc_type, exc_val, exc_tb):
if exc_type:
ctx.span.error = str(exc_val)
ctx.tracer.end_span(ctx.span)
return False
return SpanContext(self, name, kind)
def trace_llm(
self,
model_name: str,
prompt: Any,
metadata: Optional[Dict] = None
):
"""
LLM 调用追踪装饰器
用法:
@tracer.trace_llm("gpt-4", {"temperature": 0})
async def call_llm(messages):
return await llm.agenerate(messages)
"""
def decorator(func: Callable) -> Callable:
async def async_wrapper(*args, **kwargs):
with self.span(f"llm.{model_name}", SpanKind.LLM) as span:
span.attributes["model"] = model_name
span.attributes["prompt_length"] = len(str(prompt))
result = await func(*args, **kwargs)
if hasattr(result, "usage_metadata"):
span.attributes["usage"] = result.usage_metadata
return result
def sync_wrapper(*args, **kwargs):
with self.span(f"llm.{model_name}", SpanKind.LLM) as span:
span.attributes["model"] = model_name
result = func(*args, **kwargs)
return result
import asyncio
if asyncio.iscoroutinefunction(func):
return async_wrapper
return sync_wrapper
return decorator
def trace_tool(self, tool_name: str):
"""
工具调用追踪装饰器
用法:
@tracer.trace_tool("get_weather")
def get_weather(city, date):
return weather_api.query(city, date)
"""
def decorator(func: Callable) -> Callable:
async def async_wrapper(*args, **kwargs):
with self.span(f"tool.{tool_name}", SpanKind.TOOL) as span:
span.attributes["tool_name"] = tool_name
span.attributes["input"] = kwargs
try:
result = await func(*args, **kwargs)
span.attributes["success"] = True
return result
except Exception as e:
span.attributes["success"] = False
span.attributes["error"] = str(e)
raise
def sync_wrapper(*args, **kwargs):
with self.span(f"tool.{tool_name}", SpanKind.TOOL) as span:
span.attributes["tool_name"] = tool_name
try:
result = func(*args, **kwargs)
span.attributes["success"] = True
return result
except Exception as e:
span.attributes["success"] = False
span.attributes["error"] = str(e)
raise
import asyncio
if asyncio.iscoroutinefunction(func):
return async_wrapper
return sync_wrapper
return decorator
# ============== 导出后端示例 ==============
def export_to_langsmith(trace: TraceContext) -> None:
"""导出到 LangSmith"""
# 实现 LangSmith 导出逻辑
pass
def export_to_langfuse(trace: TraceContext) -> None:
"""导出到 LangFuse"""
# 实现 LangFuse 导出逻辑
pass
def export_to_file(trace: TraceContext, filepath: str = "traces.jsonl") -> None:
"""导出到文件"""
import json
with open(filepath, "a") as f:
f.write(json.dumps(trace.to_dict(), ensure_ascii=False) + "\n")
def export_to_stdout(trace: TraceContext) -> None:
"""导出到标准输出(调试用)"""
import json
print(json.dumps(trace.to_dict(), indent=2, ensure_ascii=False))四、关键指标体系
4.1 Agent 核心指标
# Agent 可观测性指标配置
metrics_config:
# 核心业务指标
business_metrics:
- name: "agent.task.success_rate"
description: "任务成功率"
type: "gauge"
aggregation: "ratio"
- name: "agent.task.avg_duration_ms"
description: "平均任务耗时"
type: "histogram"
buckets: [100, 500, 1000, 2000, 5000, 10000, 30000]
- name: "agent.task.user_satisfaction"
description: "用户满意度评分"
type: "gauge"
aggregation: "avg"
# LLM 调用指标
llm_metrics:
- name: "llm.calls.total"
description: "LLM 调用总次数"
type: "counter"
- name: "llm.tokens.prompt"
description: "Prompt Token 消耗"
type: "counter"
- name: "llm.tokens.completion"
description: "Completion Token 消耗"
type: "counter"
- name: "llm.latency.p50_ms"
description: "LLM 延迟 P50"
type: "gauge"
- name: "llm.latency.p99_ms"
description: "LLM 延迟 P99"
type: "gauge"
# 工具调用指标
tool_metrics:
- name: "tool.calls.total"
description: "工具调用总次数"
type: "counter"
labels: ["tool_name"]
- name: "tool.calls.success_rate"
description: "工具调用成功率"
type: "gauge"
labels: ["tool_name"]
- name: "tool.latency.avg_ms"
description: "工具平均延迟"
type: "gauge"
labels: ["tool_name"]
- name: "tool.errors.total"
description: "工具错误总次数"
type: "counter"
labels: ["tool_name", "error_type"]
# Agent 行为指标
agent_metrics:
- name: "agent.steps.avg"
description: "平均决策步数"
type: "histogram"
buckets: [1, 2, 3, 5, 7, 10, 15, 20]
- name: "agent.loops.detected"
description: "检测到的循环次数"
type: "counter"
- name: "agent.escalations.human"
description: "转人工次数"
type: "counter"4.2 指标收集器实现
"""
Agent 指标收集器
"""
from typing import Dict, List, Optional, Any
from dataclasses import dataclass, field
from datetime import datetime
from collections import defaultdict
import time
@dataclass
class MetricPoint:
"""指标数据点"""
name: str
value: float
timestamp: datetime = field(default_factory=datetime.now)
labels: Dict[str, str] = field(default_factory=dict)
class MetricsCollector:
"""
指标收集器
功能:
1. 计数指标
2. 数值指标
3. 直方图指标
4. 聚合计算
"""
def __init__(self):
self._counters: Dict[str, float] = defaultdict(float)
self._gauges: Dict[str, float] = {}
self._histograms: Dict[str, List[float]] = defaultdict(list)
self._last_values: Dict[str, float] = {}
def increment(
self,
name: str,
value: float = 1.0,
labels: Optional[Dict[str, str]] = None
) -> None:
"""增加计数"""
key = self._make_key(name, labels)
self._counters[key] += value
def gauge(
self,
name: str,
value: float,
labels: Optional[Dict[str, str]] = None
) -> None:
"""设置仪表值"""
key = self._make_key(name, labels)
self._gauges[key] = value
self._last_values[key] = value
def histogram(
self,
name: str,
value: float,
labels: Optional[Dict[str, str]] = None
) -> None:
"""记录直方图值"""
key = self._make_key(name, labels)
self._histograms[key].append(value)
def _make_key(self, name: str, labels: Optional[Dict[str, str]]) -> str:
"""生成带标签的键"""
if not labels:
return name
label_str = ",".join(f"{k}={v}" for k, v in sorted(labels.items()))
return f"{name}{{{label_str}}}"
def get_prometheus_metrics(self) -> str:
"""生成 Prometheus 格式的指标"""
lines = []
# 计数器
for key, value in self._counters.items():
name = key.split("{")[0]
lines.append(f"# TYPE {name} counter")
lines.append(f"{key} {value}")
# 仪表
for key, value in self._gauges.items():
name = key.split("{")[0]
lines.append(f"# TYPE {name} gauge")
lines.append(f"{key} {value}")
# 直方图
for key, values in self._histograms.items():
name = key.split("{")[0]
lines.append(f"# TYPE {name} histogram")
# 计算分位数
sorted_values = sorted(values)
total = len(sorted_values)
for quantile in [0.5, 0.9, 0.95, 0.99]:
idx = int(total * quantile)
value = sorted_values[idx] if idx < total else sorted_values[-1]
labels = key.split("{")[1].rstrip("}") if "{" in key else ""
q_key = f'{name}_quantile{{{labels},quantile="{quantile}"}}'
lines.append(f"{q_key} {value}")
# 总和与计数
sum_key = f'{name}_sum{key[key.index("{"):] if "{" in key else ""}'
count_key = f'{name}_count{key[key.index("{"):] if "{" in key else ""}'
lines.append(f"{sum_key} {sum(values)}")
lines.append(f"{count_key} {total}")
return "\n".join(lines)
def get_summary(self) -> Dict[str, Any]:
"""获取指标摘要"""
return {
"counters": dict(self._counters),
"gauges": dict(self._gauges),
"histograms": {
k: {
"count": len(v),
"mean": sum(v) / len(v),
"min": min(v),
"max": max(v),
"p50": self._percentile(v, 0.5),
"p95": self._percentile(v, 0.95),
"p99": self._percentile(v, 0.99)
}
for k, v in self._histograms.items()
}
}
def _percentile(self, values: List[float], p: float) -> float:
"""计算百分位数"""
if not values:
return 0.0
sorted_values = sorted(values)
idx = int(len(sorted_values) * p)
return sorted_values[min(idx, len(sorted_values) - 1)]
# 全局指标收集器
_global_metrics = MetricsCollector()
def get_metrics_collector() -> MetricsCollector:
"""获取全局指标收集器"""
return _global_metrics
# ============== Agent 指标追踪装饰器 ==============
def track_agent_metrics(operation: str):
"""追踪 Agent 操作指标的装饰器"""
def decorator(func):
async def async_wrapper(*args, **kwargs):
start_time = time.time()
metrics = get_metrics_collector()
try:
result = await func(*args, **kwargs)
duration_ms = (time.time() - start_time) * 1000
# 记录指标
metrics.increment("agent.operations.total", labels={"operation": operation})
metrics.histogram("agent.operations.duration_ms", duration_ms, {"operation": operation})
return result
except Exception as e:
metrics.increment("agent.operations.errors", labels={
"operation": operation,
"error_type": type(e).__name__
})
raise
def sync_wrapper(*args, **kwargs):
start_time = time.time()
metrics = get_metrics_collector()
try:
result = func(*args, **kwargs)
duration_ms = (time.time() - start_time) * 1000
metrics.increment("agent.operations.total", labels={"operation": operation})
metrics.histogram("agent.operations.duration_ms", duration_ms, {"operation": operation})
return result
except Exception as e:
metrics.increment("agent.operations.errors", labels={
"operation": operation,
"error_type": type(e).__name__
})
raise
import asyncio
if asyncio.iscoroutinefunction(func):
return async_wrapper
return sync_wrapper
return decorator五、调试技巧:回放与断点
5.1 回放系统
"""
Agent 执行回放系统
支持:记录、重放、调试
"""
import json
import pickle
from typing import Any, Dict, List, Optional
from dataclasses import dataclass, field
from datetime import datetime
from pathlib import Path
from .tracing import TraceContext
@dataclass
class ReplayEvent:
"""回放事件"""
timestamp: datetime
event_type: str # "llm_call", "tool_call", "decision", "error"
data: Dict[str, Any]
@dataclass
class ReplaySession:
"""回放会话"""
session_id: str
trace_id: str
start_time: datetime
end_time: Optional[datetime] = None
events: List[ReplayEvent] = field(default_factory=list)
# 原始输入输出
input_data: Any = None
output_data: Any = None
# 评估结果
evaluation: Dict[str, Any] = field(default_factory=dict)
class ReplayRecorder:
"""
回放记录器:记录 Agent 执行过程
用于:
1. 生产问题排查
2. 测试用例生成
3. 模型效果评估
"""
def __init__(self, storage_path: str = "./replays"):
self.storage_path = Path(storage_path)
self.storage_path.mkdir(parents=True, exist_ok=True)
self._current_session: Optional[ReplaySession] = None
self._sessions: Dict[str, ReplaySession] = {}
def start_recording(
self,
session_id: str,
trace_id: str,
input_data: Any
) -> ReplaySession:
"""开始记录"""
session = ReplaySession(
session_id=session_id,
trace_id=trace_id,
start_time=datetime.now(),
input_data=input_data
)
self._current_session = session
self._sessions[session_id] = session
return session
def record_event(
self,
event_type: str,
data: Dict[str, Any]
) -> None:
"""记录事件"""
if not self._current_session:
return
event = ReplayEvent(
timestamp=datetime.now(),
event_type=event_type,
data=data
)
self._current_session.events.append(event)
def end_recording(
self,
output_data: Any,
evaluation: Optional[Dict[str, Any]] = None
) -> ReplaySession:
"""结束记录"""
if not self._current_session:
raise RuntimeError("没有活跃的记录会话")
self._current_session.end_time = datetime.now()
self._current_session.output_data = output_data
if evaluation:
self._current_session.evaluation = evaluation
# 持久化
self._persist_session(self._current_session)
session = self._current_session
self._current_session = None
return session
def _persist_session(self, session: ReplaySession) -> None:
"""持久化会话"""
filepath = self.storage_path / f"{session.session_id}.json"
data = {
"session_id": session.session_id,
"trace_id": session.trace_id,
"start_time": session.start_time.isoformat(),
"end_time": session.end_time.isoformat() if session.end_time else None,
"input": self._serialize(session.input_data),
"output": self._serialize(session.output_data),
"evaluation": session.evaluation,
"events": [
{
"timestamp": e.timestamp.isoformat(),
"event_type": e.event_type,
"data": e.data
}
for e in session.events
]
}
with open(filepath, "w", encoding="utf-8") as f:
json.dump(data, f, ensure_ascii=False, indent=2)
def _serialize(self, data: Any) -> Any:
"""序列化数据"""
try:
return json.dumps(data, ensure_ascii=False, default=str)
except:
return str(data)
def load_session(self, session_id: str) -> Optional[ReplaySession]:
"""加载会话"""
filepath = self.storage_path / f"{session_id}.json"
if not filepath.exists():
return None
with open(filepath, "r", encoding="utf-8") as f:
data = json.load(f)
session = ReplaySession(
session_id=data["session_id"],
trace_id=data["trace_id"],
start_time=datetime.fromisoformat(data["start_time"]),
end_time=datetime.fromisoformat(data["end_time"]) if data.get("end_time") else None,
input_data=data.get("input"),
output_data=data.get("output"),
evaluation=data.get("evaluation", {})
)
session.events = [
ReplayEvent(
timestamp=datetime.fromisoformat(e["timestamp"]),
event_type=e["event_type"],
data=e["data"]
)
for e in data.get("events", [])
]
return session
def list_sessions(
self,
limit: int = 100,
offset: int = 0
) -> List[Dict[str, Any]]:
"""列出会话"""
sessions = []
for filepath in sorted(self.storage_path.glob("*.json"), reverse=True):
try:
with open(filepath, "r", encoding="utf-8") as f:
data = json.load(f)
sessions.append({
"session_id": data["session_id"],
"trace_id": data["trace_id"],
"start_time": data["start_time"],
"event_count": len(data.get("events", []))
})
except:
continue
return sessions[offset:offset + limit]
class ReplayDebugger:
"""
回放调试器:重放历史执行,支持断点
"""
def __init__(self, recorder: ReplayRecorder):
self.recorder = recorder
self._breakpoints: Dict[str, bool] = {}
def set_breakpoint(self, event_type: str) -> None:
"""设置断点"""
self._breakpoints[event_type] = True
def clear_breakpoint(self, event_type: str) -> None:
"""清除断点"""
self._breakpoints.pop(event_type, None)
def replay(
self,
session_id: str,
callback: Optional[callable] = None
) -> None:
"""
重放会话
Args:
session_id: 会话 ID
callback: 每个事件的回调函数,可用于调试
"""
session = self.recorder.load_session(session_id)
if not session:
raise ValueError(f"会话不存在: {session_id}")
print(f"开始重放会话: {session_id}")
print(f"开始时间: {session.start_time}")
print(f"事件数: {len(session.events)}")
print(f"输入: {session.input_data}")
print("-" * 80)
for i, event in enumerate(session.events):
# 检查断点
if self._breakpoints.get(event.event_type):
print(f"\n[断点] 事件类型: {event.event_type}")
input("按 Enter 继续...")
print(f"\n[{i+1}] {event.timestamp} - {event.event_type}")
print(json.dumps(event.data, indent=2, ensure_ascii=False, default=str))
# 执行回调
if callback:
callback(event)
print("\n" + "-" * 80)
print(f"重放结束")
print(f"输出: {session.output_data}")
print(f"评估: {session.evaluation}")5.2 调试面板
"""
Agent 调试面板:实时查看 Agent 执行状态
"""
import asyncio
from typing import Dict, List, Any, Optional
from dataclasses import dataclass
from datetime import datetime
import streamlit as st
from tracelite import TraceLite
@dataclass
class DebugState:
"""调试状态"""
current_step: int = 0
total_steps: int = 0
current_span: Optional[str] = None
decision_reasoning: Optional[str] = None
last_tool_call: Optional[Dict] = None
error: Optional[str] = None
class AgentDebugPanel:
"""
Agent 调试面板
功能:
1. 实时显示执行状态
2. 查看历史追踪
3. 设置条件断点
4. 分析性能瓶颈
"""
def __init__(self):
self.tracer = TraceLite()
self.state = DebugState()
def render(self):
"""渲染调试面板"""
st.set_page_config(page_title="Agent Debug Panel")
# 侧边栏:会话列表
with st.sidebar:
st.title("调试会话")
# 会话选择
sessions = self.tracer.list_sessions()
session_options = {s["session_id"]: s for s in sessions}
selected = st.selectbox(
"选择会话",
options=list(session_options.keys()),
format_func=lambda x: f"{x} ({session_options[x]['event_count']} 事件)"
)
if selected:
self._render_session_detail(session_options[selected])
# 主区域:执行可视化
st.title("Agent 执行追踪")
# 追踪树
self._render_trace_tree()
# 指标面板
self._render_metrics_panel()
# 日志流
self._render_log_stream()
def _render_session_detail(self, session: Dict):
"""渲染会话详情"""
st.subheader("会话信息")
col1, col2 = st.columns(2)
with col1:
st.metric("事件数", session.get("event_count", 0))
with col2:
st.metric("开始时间", session.get("start_time", "N/A"))
# 加载会话
full_session = self.tracer.load_session(session["session_id"])
if full_session:
st.json({
"input": full_session.input_data,
"output": full_session.output_data,
"evaluation": full_session.evaluation
})
def _render_trace_tree(self):
"""渲染追踪树"""
st.subheader("执行追踪")
# 模拟追踪数据
trace_data = {
"name": "agent.query_order",
"duration_ms": 2500,
"spans": [
{
"id": "span1",
"name": "intent_detection",
"kind": "llm",
"duration_ms": 500,
"status": "ok"
},
{
"id": "span2",
"name": "query_order",
"kind": "tool",
"duration_ms": 1200,
"status": "ok",
"children": [
{
"id": "span2.1",
"name": "api_call",
"kind": "internal",
"duration_ms": 1000,
"status": "ok"
}
]
},
{
"id": "span3",
"name": "format_response",
"kind": "llm",
"duration_ms": 800,
"status": "ok"
}
]
}
# 使用树形图展示
self._render_span_tree(trace_data["spans"], level=0)
def _render_span_tree(self, spans: List[Dict], level: int):
"""递归渲染跨度树"""
for span in spans:
indent = " " * level
status_emoji = "✅" if span["status"] == "ok" else "❌"
# 展开/折叠
with st.expander(f"{indent}{status_emoji} {span['name']} ({span['duration_ms']}ms)"):
st.json(span)
# 递归渲染子跨度
if "children" in span:
self._render_span_tree(span["children"], level + 1)
def _render_metrics_panel(self):
"""渲染指标面板"""
st.subheader("关键指标")
col1, col2, col3, col4 = st.columns(4)
with col1:
st.metric("LLM 调用", 3, delta="-1 vs 上次")
with col2:
st.metric("Token 消耗", "1.2K", delta="-15%")
with col3:
st.metric("工具调用", 1, delta="0 vs 上次")
with col4:
st.metric("执行时间", "2.5s", delta="-0.3s")
def _render_log_stream(self):
"""渲染日志流"""
st.subheader("实时日志")
# 模拟日志
logs = [
{"time": "10:00:01", "level": "INFO", "message": "开始处理请求"},
{"time": "10:00:01", "level": "DEBUG", "message": "意图识别: query_order"},
{"time": "10:00:02", "level": "INFO", "message": "调用工具: query_order"},
{"time": "10:00:03", "level": "INFO", "message": "工具执行成功"},
{"time": "10:00:03", "level": "DEBUG", "message": "生成响应..."},
{"time": "10:00:04", "level": "INFO", "message": "请求处理完成"},
]
for log in logs:
color = {
"INFO": "white",
"DEBUG": "gray",
"WARNING": "yellow",
"ERROR": "red"
}.get(log["level"], "white")
st.text(f"[{log['time']}] {log['level']}: {log['message']}")六、告警规则与自动告警
6.1 告警规则配置
# Agent 告警规则配置
alerting_config:
# 告警通道
channels:
- name: "slack"
type: "webhook"
url: "https://hooks.slack.com/services/xxx"
enabled: true
- name: "email"
type: "smtp"
recipients: ["oncall@company.com"]
enabled: false
- name: "pagerduty"
type: "pagerduty"
integration_key: "xxx"
enabled: true
# 告警规则
rules:
# 无限循环检测
- name: "infinite_loop_detection"
description: "检测 Agent 进入无限循环"
condition: >
agent.loops.detected > 0
and agent.steps.current > 20
severity: "critical"
cooldown: 5m
channels: ["slack", "pagerduty"]
# 工具调用异常
- name: "tool_failure_rate_high"
description: "工具调用失败率过高"
condition: >
rate(tool.calls.errors) / rate(tool.calls.total) > 0.1
severity: "warning"
cooldown: 10m
channels: ["slack"]
# 延迟过高
- name: "latency_high"
description: "LLM 响应延迟过高"
condition: >
llm.latency.p99 > 10000 # 10秒
severity: "warning"
cooldown: 5m
channels: ["slack"]
# Token 消耗异常
- name: "token_consumption_anomaly"
description: "Token 消耗异常"
condition: >
agent.tokens.current > agent.tokens.avg * 3
severity: "warning"
cooldown: 15m
channels: ["slack"]
# 错误率过高
- name: "error_rate_critical"
description: "Agent 错误率超过阈值"
condition: >
rate(agent.operations.errors) / rate(agent.operations.total) > 0.05
severity: "critical"
cooldown: 5m
channels: ["slack", "pagerduty"]
# 转人工率过高
- name: "human_escalation_high"
description: "转人工频率过高"
condition: >
rate(agent.escalations.human) / rate(agent.operations.total) > 0.1
severity: "warning"
cooldown: 30m
channels: ["slack"]
# 成本超限
- name: "cost_over_budget"
description: "API 成本超过预算"
condition: >
agent.cost.daily > 1000 # 美元
severity: "warning"
cooldown: 1h
channels: ["slack", "email"]6.2 告警系统实现
"""
Agent 告警系统
"""
import asyncio
from typing import Dict, List, Optional, Callable, Any
from dataclasses import dataclass, field
from datetime import datetime, timedelta
from enum import Enum
import logging
logger = logging.getLogger(__name__)
class AlertSeverity(Enum):
"""告警级别"""
INFO = "info"
WARNING = "warning"
CRITICAL = "critical"
@dataclass
class Alert:
"""告警"""
name: str
message: str
severity: AlertSeverity
timestamp: datetime = field(default_factory=datetime.now)
labels: Dict[str, str] = field(default_factory=dict)
# 关联数据
trace_id: Optional[str] = None
session_id: Optional[str] = None
# 上下文
context: Dict[str, Any] = field(default_factory=dict)
@dataclass
class AlertRule:
"""告警规则"""
name: str
description: str
condition: Callable[[Dict], bool] # 返回 True 触发告警
severity: AlertSeverity
# 抑制配置
cooldown: timedelta = field(default_factory=lambda: timedelta(minutes=5))
last_triggered: Optional[datetime] = None
# 通道
channels: List[str] = field(default_factory=list)
# 启用状态
enabled: bool = True
class AlertManager:
"""
告警管理器
功能:
1. 定义和评估告警规则
2. 告警聚合和抑制
3. 多通道通知
4. 告警历史记录
"""
def __init__(self):
self.rules: Dict[str, AlertRule] = {}
self.channels: Dict[str, Callable] = {}
self.alert_history: List[Alert] = []
# 初始化默认规则
self._init_default_rules()
def _init_default_rules(self):
"""初始化默认告警规则"""
# 无限循环检测
self.add_rule(AlertRule(
name="infinite_loop_detection",
description="检测 Agent 进入无限循环",
condition=lambda m: m.get("agent.loops.detected", 0) > 0 and
m.get("agent.steps.current", 0) > 20,
severity=AlertSeverity.CRITICAL,
channels=["slack", "pagerduty"]
))
# 工具失败率
self.add_rule(AlertRule(
name="tool_failure_rate_high",
description="工具调用失败率超过 10%",
condition=lambda m: (
m.get("tool.calls.total", 0) > 0 and
m.get("tool.calls.errors", 0) / m.get("tool.calls.total", 1) > 0.1
),
severity=AlertSeverity.WARNING,
channels=["slack"]
))
# 延迟过高
self.add_rule(AlertRule(
name="latency_high",
description="LLM P99 延迟超过 10 秒",
condition=lambda m: m.get("llm.latency.p99", 0) > 10000,
severity=AlertSeverity.WARNING,
channels=["slack"]
))
def add_rule(self, rule: AlertRule) -> None:
"""添加告警规则"""
self.rules[rule.name] = rule
def add_channel(self, name: str, sender: Callable) -> None:
"""添加告警通道"""
self.channels[name] = sender
def evaluate(self, metrics: Dict[str, Any]) -> List[Alert]:
"""
评估所有规则
Args:
metrics: 当前指标数据
Returns:
触发的告警列表
"""
triggered_alerts = []
for rule in self.rules.values():
if not rule.enabled:
continue
# 检查冷却期
if rule.last_triggered:
if datetime.now() - rule.last_triggered < rule.cooldown:
continue
# 评估条件
try:
if rule.condition(metrics):
alert = Alert(
name=rule.name,
message=rule.description,
severity=rule.severity,
context={"metrics": metrics}
)
triggered_alerts.append(alert)
rule.last_triggered = datetime.now()
except Exception as e:
logger.error(f"规则 {rule.name} 评估失败: {e}")
return triggered_alerts
async def send_alerts(self, alerts: List[Alert]) -> None:
"""发送告警"""
for alert in alerts:
# 记录历史
self.alert_history.append(alert)
# 发送到各通道
for rule in self.rules.values():
if rule.name == alert.name:
for channel_name in rule.channels:
if channel_name in self.channels:
try:
await self.channels[channel_name](alert)
except Exception as e:
logger.error(f"发送告警到 {channel_name} 失败: {e}")
def get_active_alerts(self) -> List[Alert]:
"""获取活跃告警(最近 1 小时内)"""
cutoff = datetime.now() - timedelta(hours=1)
return [
a for a in self.alert_history
if a.timestamp > cutoff
]
def get_alert_history(
self,
limit: int = 100,
severity: Optional[AlertSeverity] = None
) -> List[Alert]:
"""获取告警历史"""
alerts = self.alert_history[-limit:]
if severity:
alerts = [a for a in alerts if a.severity == severity]
return alerts
# ============== 告警通道实现 ==============
async def slack_sender(webhook_url: str, alert: Alert) -> None:
"""发送 Slack 告警"""
import aiohttp
severity_emoji = {
AlertSeverity.INFO: "ℹ️",
AlertSeverity.WARNING: "⚠️",
AlertSeverity.CRITICAL: "🚨"
}
payload = {
"text": f"{severity_emoji[alert.severity]} *{alert.name}*",
"blocks": [
{
"type": "header",
"text": {
"type": "plain_text",
"text": f"Agent 告警: {alert.name}"
}
},
{
"type": "section",
"fields": [
{"type": "mrkdwn", "text": f"*级别:*\n{alert.severity.value}"},
{"type": "mrkdwn", "text": f"*时间:*\n{alert.timestamp.isoformat()}"}
]
},
{
"type": "section",
"text": {
"type": "mrkdwn",
"text": f"*描述:*\n{alert.message}"
}
}
]
}
async with aiohttp.ClientSession() as session:
await session.post(webhook_url, json=payload)
# ============== 与 Agent 集成 ==============
class AlertingMiddleware:
"""
告警中间件:集成到 Agent 执行链路
"""
def __init__(self, alert_manager: AlertManager):
self.alert_manager = alert_manager
self.metrics_collector = get_metrics_collector()
async def check_alerts(self) -> None:
"""检查告警规则"""
metrics = self.metrics_collector.get_summary()
# 评估规则
alerts = self.alert_manager.evaluate(metrics)
# 发送告警
if alerts:
await self.alert_manager.send_alerts(alerts)七、完整可观测性集成代码
"""
Agent 可观测性完整集成示例
使用 LangChain + LangSmith/LangFuse + 自建追踪系统
"""
import os
from typing import Dict, Any, List, Optional
from langchain_openai import ChatOpenAI
from langchain.prompts import ChatPromptTemplate
from langchain.schema import StrOutputParser
from langchain.tools import tool
from langchain.agents import AgentExecutor, create_openai_functions_agent
from langchain.callbacks.manager import CallbackManager
from agent_tracing import AgentTracer, SpanKind
from agent_metrics import get_metrics_collector, track_agent_metrics
from agent_alerting import AlertManager, AlertingMiddleware
class ObservableAgent:
"""
可观测 Agent:集成完整的可观测性能力
功能:
1. 全链路追踪
2. 指标收集
3. 自动告警
4. 回放支持
"""
def __init__(
self,
llm: ChatOpenAI,
tools: List[Any],
project_name: str = "agent-debugging"
):
self.llm = llm
self.tools = tools
self.project_name = project_name
# 初始化组件
self._init_tracer()
self._init_metrics()
self._init_alerting()
self._init_replay()
def _init_tracer(self):
"""初始化追踪器"""
# 导出后端列表
export_backends = []
# 如果配置了 LangSmith
if os.getenv("LANGCHAIN_TRACING_V2"):
from langchain.callbacks.tracers.langsmith import LangSmithTracer
export_backends.append(
lambda trace: LangSmithTracer(project_name=self.project_name).on_chain_end(
{"trace": trace.to_dict()}
)
)
self.tracer = AgentTracer(
service_name="observable_agent",
export_backends=export_backends
)
def _init_metrics(self):
"""初始化指标收集"""
self.metrics = get_metrics_collector()
def _init_alerting(self):
"""初始化告警"""
self.alert_manager = AlertManager()
# 配置告警通道
if os.getenv("SLACK_WEBHOOK_URL"):
from agent_alerting import slack_sender
self.alert_manager.add_channel(
"slack",
lambda alert: slack_sender(os.getenv("SLACK_WEBHOOK_URL"), alert)
)
self.alerting_middleware = AlertingMiddleware(self.alert_manager)
def _init_replay(self):
"""初始化回放系统"""
from agent_replay import ReplayRecorder, ReplayDebugger
self.replay_recorder = ReplayRecorder()
self.replay_debugger = ReplayDebugger(self.replay_recorder)
@track_agent_metrics("agent.run")
async def run(self, user_input: str, session_id: Optional[str] = None) -> Dict[str, Any]:
"""
运行 Agent
Args:
user_input: 用户输入
session_id: 会话 ID(用于回放)
Returns:
{
"response": str, # Agent 回复
"trace_id": str, # 追踪 ID
"metrics": Dict, # 指标快照
"success": bool # 是否成功
}
"""
import uuid
trace_id = str(uuid.uuid4())[:16]
session_id = session_id or trace_id
# 开始记录
self.replay_recorder.start_recording(
session_id=session_id,
trace_id=trace_id,
input_data=user_input
)
# 开始追踪
self.tracer.start_trace(
name=f"agent.{session_id}",
input_data={"user_input": user_input},
metadata={"session_id": session_id}
)
try:
# 构建 Agent
prompt = ChatPromptTemplate.from_messages([
("system", "你是一个智能助手,使用工具来回答用户问题。"),
("human", "{input}"),
("ai", "{agent_scratchpad}")
])
agent = create_openai_functions_agent(
llm=self.llm,
tools=self.tools,
prompt=prompt
)
agent_executor = AgentExecutor.from_agent_and_tools(
agent=agent,
tools=self.tools,
verbose=True,
handle_parsing_errors=True
)
# 执行
with self.tracer.span("agent_execution", SpanKind.CHAIN):
result = await agent_executor.ainvoke({"input": user_input})
# 结束追踪
output = result.get("output", "")
self.tracer.end_trace(output_data=output)
# 结束记录
self.replay_recorder.end_recording(
output_data=output,
evaluation={"success": True}
)
# 检查告警
await self.alerting_middleware.check_alerts()
return {
"response": output,
"trace_id": trace_id,
"session_id": session_id,
"metrics": self.metrics.get_summary(),
"success": True
}
except Exception as e:
# 错误处理
self.tracer.end_trace(error=str(e))
self.replay_recorder.end_recording(
output_data=None,
evaluation={"success": False, "error": str(e)}
)
return {
"response": f"抱歉,处理您的请求时遇到错误: {str(e)}",
"trace_id": trace_id,
"session_id": session_id,
"error": str(e),
"success": False
}
def debug_session(self, session_id: str):
"""调试历史会话"""
self.replay_debugger.replay(session_id)
# ============== 使用示例 ==============
async def demo_observable_agent():
"""可观测 Agent 使用示例"""
# 初始化 LLM
llm = ChatOpenAI(model="gpt-4o", temperature=0)
# 定义工具
@tool
def get_weather(city: str) -> str:
"""获取城市天气"""
weathers = {"北京": "晴 15°C", "上海": "多云 22°C"}
return weathers.get(city, "未知")
tools = [get_weather]
# 创建可观测 Agent
agent = ObservableAgent(llm=llm, tools=tools)
# 运行
result = await agent.run("北京天气怎么样?")
print(f"回复: {result['response']}")
print(f"追踪 ID: {result['trace_id']}")
# 调试历史会话
# agent.debug_session(result['session_id'])
if __name__ == "__main__":
import asyncio
asyncio.run(demo_observable_agent())八、总结
Agent 系统的可观测性是构建可靠生产系统的基石。通过追踪、指标、日志三大支柱,我们可以让"不可预测"的 Agent 变得可控可追踪。
核心要点回顾:
追踪是基础:完整记录每个决策点、工具调用的输入输出
指标要全面:从业务、LLM、工具、Agent 行为多个维度收集
告警要精准:避免告警疲劳,聚焦真正需要人工介入的问题
回放是调试利器:生产问题的完美复现离不开完整的执行记录
技术栈推荐:
可观测性不是事后补救,而是从设计之初就需要考虑的架构问题。投入建设的可观测性基础设施,终将在生产问题排查时得到丰厚回报。
评论区