Gemini模型+ADK避坑指南:手把手构建带记忆的天气智能体团队
构建带记忆的智能体团队:从状态管理到多智能体协作的实战进阶
在当今AI应用开发领域,构建一个能够理解上下文、记住用户偏好并协同工作的智能体系统,正从“锦上添花”演变为“不可或缺”的核心能力。对于中高级开发者而言,单纯调用大模型API返回一次性答案的时代已经过去。真正的挑战在于如何设计一个具备状态感知和团队协作能力的智能系统,使其能够像人类助手一样,在连续对话中提供个性化、连贯的服务。
想象这样一个场景:用户第一次询问“纽约天气如何?”时,你构建的天气助手不仅返回了当前气温,还记住了用户偏好使用摄氏度。当用户几轮对话后再次询问“伦敦呢?”,助手能自动以摄氏度单位回应,无需用户重复声明。更进一步,当用户说“你好”时,由专门的问候智能体处理;询问天气时,由专业天气智能体响应;告别时,则由告别智能体优雅收尾——这就是智能体团队的魅力所在。
本文将深入探讨如何利用Google的Agent Development Kit(ADK)和Gemini模型,构建一个具备记忆功能和团队协作能力的天气服务智能体系统。我们将超越基础教程,聚焦于会话状态存储、多智能体路由机制以及用户偏好记忆等进阶特性,并通过“华氏/摄氏自动转换”等实战案例,手把手带你避开开发中的常见陷阱。
1. 智能体状态管理:让AI记住用户偏好
在传统对话系统中,每次用户查询都被视为独立事件,系统无法记住之前的交互内容。这种“健忘症”严重限制了AI助手的实用性。智能体状态管理正是为了解决这一问题而生——它允许系统在会话过程中持久化存储和检索关键信息。
1.1 会话状态的核心价值
会话状态不仅仅是存储用户ID或对话历史那么简单。一个设计良好的状态管理系统应该能够:
- 记住用户偏好:温度单位、语言风格、信息详细程度等个性化设置
- 维护对话上下文:理解当前对话在更广泛任务中的位置
- 跟踪任务进度:对于多步骤操作,记录已完成和待完成的步骤
- 存储临时数据:缓存中间计算结果,避免重复计算
在ADK框架中,状态管理通过Session对象实现。每个会话都有唯一的标识符,并可以存储任意键值对形式的状态数据。下面是一个基础的状态管理示例:
from google.adk.sessions import InMemorySessionService
# 创建会话服务
session_service = InMemorySessionService()
# 创建新会话并设置初始状态
session = session_service.create_session(
app_name="weather_app",
user_id="user_123",
session_id="session_001",
state={
"user_preference_temperature_unit": "Celsius",
"last_queried_city": None,
"interaction_count": 0
}
)
# 在后续交互中更新状态
def update_user_preference(session, unit):
"""更新用户的温度单位偏好"""
if session and session.state is not None:
session.state["user_preference_temperature_unit"] = unit
session.state["preference_updated_at"] = datetime.now().isoformat()
return True
return False
注意:在实际生产环境中,
InMemorySessionService仅适用于开发测试。对于需要持久化和水平扩展的场景,需要实现基于数据库(如Redis、PostgreSQL)的自定义会话服务。
1.2 状态感知工具的设计模式
要让工具(Tool)能够访问和修改会话状态,ADK提供了ToolContext机制。通过将状态感知逻辑嵌入工具函数,我们可以创建真正个性化的用户体验。
以下是一个状态感知的天气查询工具实现:
from typing import Dict, Any, Optional
from google.adk.tools import FunctionTool, ToolContext
def get_weather_with_preference(
city: str,
tool_context: Optional[ToolContext] = None
) -> Dict[str, Any]:
"""
获取指定城市的天气,并考虑用户的温度单位偏好
参数:
city: 城市名称
tool_context: 可选的工具上下文,包含会话状态
返回:
包含天气信息的字典,温度以用户偏好单位显示
"""
# 从模拟数据库获取天气数据
weather_data = mock_weather_db.get(city)
if not weather_data:
return {"status": "error", "message": f"未找到{city}的天气数据"}
# 获取用户温度偏好(默认为摄氏度)
temperature_unit = "Celsius"
if tool_context and tool_context.state:
temperature_unit = tool_context.state.get(
"user_preference_temperature_unit",
"Celsius"
)
# 原始数据为华氏度
temp_f = weather_data["temperature"]
temp_c = round((temp_f - 32) * 5/9, 1)
# 根据用户偏好选择显示单位
display_temp = temp_c if temperature_unit == "Celsius" else temp_f
# 更新会话状态中的最后查询记录
if tool_context and tool_context.state is not None:
tool_context.state["last_queried_city"] = city
tool_context.state["last_query_time"] = datetime.now().isoformat()
return {
"status": "success",
"city": city,
"temperature": {
"value": display_temp,
"unit": temperature_unit,
"original_fahrenheit": temp_f,
"original_celsius": temp_c
},
"condition": weather_data["condition"],
"humidity": weather_data["humidity"],
"timestamp": datetime.now().isoformat()
}
# 将函数转换为ADK工具
weather_tool_stateful = FunctionTool(get_weather_with_preference)
这个工具的设计有几个关键点:
- 向后兼容:
tool_context参数是可选的,确保工具在无状态环境下也能工作 - 默认值处理:当状态中不存在偏好设置时,使用合理的默认值(摄氏度)
- 状态更新:不仅读取状态,还在适当时候更新状态,记录用户行为
- 数据完整性:同时保存原始数据和转换后的数据,便于后续处理
1.3 用户偏好检测与更新机制
智能体不仅需要记住用户偏好,还需要能够从对话中检测和更新这些偏好。这通常通过自然语言理解(NLU)和指令(Instruction)设计来实现。
在智能体的指令中,我们可以明确指定如何检测和更新用户偏好:
stateful_agent = Agent(
name="StatefulWeatherAgent",
description="我是一个能够记住用户偏好的天气智能体",
model="gemini-2.0-flash",
tools=[weather_tool_stateful],
instruction="""
你是一个个性化的天气信息助手。你的核心任务是:
1. 当用户查询城市天气时,使用get_weather_with_preference工具获取信息
2. 自动检测并更新用户的温度单位偏好
3. 在连续对话中保持上下文一致性
关于温度单位偏好的处理规则:
* 当用户提到"摄氏度"、"摄氏"、"°C"、"C"时,立即更新偏好为"Celsius"
* 当用户提到"华氏度"、"华氏"、"°F"、"F"时,立即更新偏好为"Fahrenheit"
* 更新偏好后,需要向用户确认:"已为您设置为使用{单位}显示温度"
重要行为准则:
- 永远不要询问用户"您想用哪种温度单位?",而是从对话中推断
- 如果用户没有明确指定,使用当前会话中存储的偏好设置
- 每次查询天气时,都使用用户偏好的温度单位显示结果
- 在响应中明确显示使用的温度单位,例如:"当前温度25°C(77°F)"
示例交互:
用户:"我喜欢用华氏度"
你:[更新状态为Fahrenheit] "好的,已为您设置为使用华氏度显示温度"
用户:"纽约天气怎么样?"
你:[调用工具查询纽约天气,使用华氏度显示] "纽约当前天气晴朗,温度72°F(22.2°C),湿度50%"
用户:"切换到摄氏度"
你:[更新状态为Celsius] "已切换到摄氏度显示"
用户:"伦敦呢?"
你:[调用工具查询伦敦天气,使用摄氏度显示] "伦敦当前雨天,温度15.6°C(60°F),湿度80%"
""",
output_key="agent_response"
)
这种设计模式的优势在于:
- 无感学习:用户无需显式设置,系统从对话中自动学习
- 即时生效:偏好更新后立即应用于后续查询
- 透明反馈:系统明确告知用户当前的设置状态
- 容错处理:当无法确定偏好时,使用合理的默认值
1.4 状态持久化策略
对于生产环境应用,状态管理需要考虑持久化策略。ADK的会话服务接口允许我们实现自定义的存储后端。
以下是一个基于SQLite的简单持久化示例:
import sqlite3
from datetime import datetime
from typing import Optional, Dict, Any
from google.adk.sessions import Session, SessionService
class SQLiteSessionService(SessionService):
"""基于SQLite的持久化会话服务"""
def __init__(self, db_path: str = "sessions.db"):
self.db_path = db_path
self._init_database()
def _init_database(self):
"""初始化数据库表结构"""
conn = sqlite3.connect(self.db_path)
cursor = conn.cursor()
cursor.execute("""
CREATE TABLE IF NOT EXISTS sessions (
app_name TEXT NOT NULL,
user_id TEXT NOT NULL,
session_id TEXT NOT NULL,
state_json TEXT NOT NULL,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (app_name, user_id, session_id)
)
""")
conn.commit()
conn.close()
def create_session(
self,
app_name: str,
user_id: str,
session_id: str,
state: Optional[Dict[str, Any]] = None
) -> Session:
"""创建新会话"""
if state is None:
state = {}
session = Session(
app_name=app_name,
user_id=user_id,
session_id=session_id,
state=state
)
# 保存到数据库
self._save_session(session)
return session
def get_session(
self,
app_name: str,
user_id: str,
session_id: str
) -> Optional[Session]:
"""获取现有会话"""
conn = sqlite3.connect(self.db_path)
cursor = conn.cursor()
cursor.execute("""
SELECT state_json FROM sessions
WHERE app_name = ? AND user_id = ? AND session_id = ?
""", (app_name, user_id, session_id))
row = cursor.fetchone()
conn.close()
if row:
import json
state = json.loads(row[0])
return Session(
app_name=app_name,
user_id=user_id,
session_id=session_id,
state=state
)
return None
def _save_session(self, session: Session):
"""保存会话到数据库"""
import json
conn = sqlite3.connect(self.db_path)
cursor = conn.cursor()
state_json = json.dumps(session.state)
cursor.execute("""
INSERT OR REPLACE INTO sessions
(app_name, user_id, session_id, state_json, updated_at)
VALUES (?, ?, ?, ?, CURRENT_TIMESTAMP)
""", (session.app_name, session.user_id,
session.session_id, state_json))
conn.commit()
conn.close()
这个自定义会话服务提供了:
- 数据持久化:会话状态保存在SQLite数据库中
- 会话恢复:应用重启后可以恢复之前的会话状态
- 状态版本管理:通过时间戳跟踪状态更新时间
- 易于扩展:可以轻松迁移到PostgreSQL、MySQL等生产级数据库
2. 多智能体协作架构设计
单个智能体能力有限,真正的强大之处在于多个智能体协同工作。ADK的智能体团队(Agent Team)功能允许我们构建一个由多个专门化智能体组成的系统,每个智能体负责特定任务,通过路由机制协同工作。
2.1 智能体团队的核心概念
智能体团队不是简单地将多个智能体堆砌在一起,而是需要精心设计的协作架构:
| 组件类型 | 职责 | 示例 |
|---|---|---|
| 路由智能体 | 分析用户输入,决定由哪个子智能体处理 | 主调度器,根据意图路由请求 |
| 专业智能体 | 处理特定领域的任务 | 天气查询、问候、告别等专业智能体 |
| 工具智能体 | 提供特定功能调用 | 数据查询、计算、外部API调用 |
| 协调智能体 | 管理多个智能体间的协作 | 处理需要多步骤、多智能体参与的任务 |
在天气服务示例中,我们可以设计这样的团队结构:
from google.adk.agents import Agent
# 1. 专业问候智能体
greeting_agent = Agent(
name="GreetingAgent",
description="我负责处理用户的问候和欢迎消息",
model="gemini-2.0-flash",
tools=[greeting_tool],
instruction="""
你是一个友好的问候助手。当用户说"你好"、"早上好"、"嗨"等问候语时,
你负责回应热情的欢迎消息。
你的回应应该:
1. 热情但不过度
2. 个性化(如果知道用户名字)
3. 简要介绍可提供的服务
4. 保持对话流畅自然
不要处理非问候相关的查询,这些应该由其他智能体处理。
"""
)
# 2. 专业天气智能体
weather_agent = Agent(
name="WeatherAgent",
description="我专门提供城市天气信息查询",
model="gemini-2.0-flash",
tools=[weather_tool_stateful],
instruction="""
你是一个专业的天气信息提供者。你的唯一任务是:
1. 当用户提到城市名称时,立即使用get_weather_with_preference工具查询天气
2. 以清晰、友好的格式呈现天气信息
3. 包含温度、天气状况、湿度等关键信息
4. 使用用户偏好的温度单位显示
重要规则:
- 不要询问用户想查询哪个城市,直接从消息中提取
- 不要进行无关的闲聊或解释
- 如果工具返回错误,礼貌地告知用户
示例:
用户:"纽约天气怎么样?"
你:[调用工具] "纽约当前天气晴朗,温度22°C(72°F),湿度50%"
用户:"伦敦和东京的天气呢?"
你:[分别调用工具] "伦敦:雨天,15°C(59°F),湿度80%。东京:多云,26°C(79°F),湿度65%"
"""
)
# 3. 专业告别智能体
farewell_agent = Agent(
name="FarewellAgent",
description="我负责处理用户的告别和结束对话",
model="gemini-2.0-flash",
tools=[farewell_tool],
instruction="""
你是一个礼貌的告别助手。当用户表示要结束对话时,
你负责提供友好的告别消息。
常见告别信号包括:
- "再见"、"拜拜"、"下次聊"
- "谢谢"、"感谢帮助"
- "我要走了"、"结束对话"
你的回应应该:
1. 表达感谢
2. 提供适当的结束语
3. 保持友好和专业
4. 不要继续对话话题
示例:
用户:"谢谢,再见!"
你:"不客气!很高兴为您服务。再见!"
"""
)
# 4. 路由智能体(团队协调者)
weather_team_router = Agent(
name="WeatherTeamRouter",
description="我是天气服务团队的总调度员,负责将请求路由给合适的专家",
model="gemini-2.0-flash",
tools=[], # 路由智能体通常不直接使用工具
sub_agents=[greeting_agent, weather_agent, farewell_agent],
instruction="""
你是天气服务智能体团队的调度中心。你的职责是分析用户输入,
并将其路由给最合适的专业智能体处理。
路由规则:
1. 问候检测:
- 如果用户输入包含问候语(如"你好"、"嗨"、"早上好"等)
- 或者输入是简单的打招呼
- 路由给GreetingAgent
2. 天气查询检测:
- 如果用户输入包含城市名称(如"纽约"、"伦敦"、"东京"等)
- 或者输入明显是询问天气(如"天气如何"、"温度多少")
- 路由给WeatherAgent
3. 告别检测:
- 如果用户输入包含告别语(如"再见"、"谢谢"、"结束"等)
- 或者用户明确表示要结束对话
- 路由给FarewellAgent
4. 模糊或复杂查询:
- 如果输入同时包含多个意图
- 或者意图不明确
- 先尝试澄清,然后根据澄清结果路由
重要原则:
- 不要自己处理专业任务,总是委托给专业智能体
- 保持路由决策快速准确
- 如果无法确定路由目标,礼貌地请求用户澄清
示例路由:
用户:"你好!"
你:[路由给GreetingAgent]
用户:"纽约今天天气怎么样?"
你:[路由给WeatherAgent]
用户:"谢谢你的帮助,再见!"
你:[路由给FarewellAgent]
用户:"先问个好,然后告诉我伦敦天气"
你:[先路由给GreetingAgent,然后路由给WeatherAgent]
""",
output_key="routed_response"
)
2.2 智能体间的通信与协调
智能体团队的核心挑战在于如何让智能体之间有效通信和协调。ADK提供了几种机制:
1. 显式路由 通过主智能体的指令明确指定路由逻辑,如上例所示。这种方式控制性强,但需要手动维护路由规则。
2. 基于LLM的自动路由 利用大语言模型的意图识别能力自动决定路由:
def auto_route_agent(user_input: str, available_agents: List[Agent]) -> Optional[Agent]:
"""
基于LLM的自动路由函数
参数:
user_input: 用户输入文本
available_agents: 可用智能体列表
返回:
最适合处理该输入的智能体,或None
"""
# 构建智能体描述供LLM参考
agent_descriptions = []
for agent in available_agents:
agent_descriptions.append(f"- {agent.name}: {agent.description}")
prompt = f"""
根据用户输入,选择最合适的智能体进行处理。
可用智能体:
{chr(10).join(agent_descriptions)}
用户输入:{user_input}
请只返回智能体名称,不要添加任何解释。
如果没有合适的智能体,返回"None"。
"""
# 调用LLM进行路由决策
response = llm_client.complete(prompt)
selected_agent_name = response.strip()
# 查找对应的智能体
for agent in available_agents:
if agent.name == selected_agent_name:
return agent
return None
3. 智能体链式调用 一个智能体的输出可以作为另一个智能体的输入,形成处理流水线:
class ChainedAgentSystem:
"""链式智能体系统"""
def __init__(self, agents: List[Agent]):
self.agents = agents
self.execution_history = []
async def process_query(self, user_input: str, context: Dict = None) -> str:
"""链式处理用户查询"""
current_input = user_input
final_response = ""
for i, agent in enumerate(self.agents):
# 记录执行历史
self.execution_history.append({
"step": i,
"agent": agent.name,
"input": current_input
})
# 执行当前智能体
response = await self._run_agent(agent, current_input, context)
# 更新输入供下一个智能体使用
if i < len(self.agents) - 1:
# 对于中间智能体,将其输出作为下一个智能体的输入
current_input = response
else:
# 最后一个智能体的输出作为最终响应
final_response = response
return final_response
async def _run_agent(self, agent: Agent, input_text: str, context: Dict) -> str:
"""运行单个智能体"""
# 这里简化了实际实现
# 实际中需要调用ADK的Runner
return f"Agent {agent.name} processed: {input_text}"
2.3 团队协作的实战模式
在实际应用中,智能体团队可以采取多种协作模式:
模式一:并行处理 多个智能体同时处理同一输入的不同方面,然后合并结果。
import asyncio
from typing import List, Dict, Any
async def parallel_agent_processing(
user_input: str,
agents: List[Agent],
context: Dict[str, Any]
) -> Dict[str, Any]:
"""
并行处理:多个智能体同时处理同一输入
适用于:
- 需要多角度分析的问题
- 事实核查
- 综合评估
"""
# 创建所有智能体的任务
tasks = []
for agent in agents:
task = asyncio.create_task(
run_agent_async(agent, user_input, context)
)
tasks.append((agent.name, task))
# 等待所有任务完成
results = {}
for agent_name, task in tasks:
try:
result = await task
results[agent_name] = result
except Exception as e:
results[agent_name] = {"error": str(e)}
# 整合结果
return integrate_results(results)
def integrate_results(agent_results: Dict[str, Any]) -> Dict[str, Any]:
"""整合多个智能体的结果"""
# 简单的多数投票或加权平均
# 实际中可能需要更复杂的整合逻辑
integrated = {
"raw_results": agent_results,
"consensus": None,
"confidence": 0.0
}
# 这里可以添加具体的整合逻辑
# 例如:对于天气查询,取所有智能体的平均温度
return integrated
模式二:顺序流水线 智能体按特定顺序处理,每个智能体处理前一个的输出。
class ProcessingPipeline:
"""顺序处理流水线"""
def __init__(self):
self.pipeline = []
def add_stage(self, agent: Agent, condition=None):
"""向流水线添加处理阶段"""
self.pipeline.append({
"agent": agent,
"condition": condition, # 可选的条件函数
"description": agent.description
})
async def process(self, user_input: str, context: Dict) -> str:
"""按顺序执行流水线"""
current_output = user_input
execution_log = []
for stage in self.pipeline:
stage_info = {
"agent": stage["agent"].name,
"input": current_output,
"condition_met": True
}
# 检查条件(如果有)
if stage["condition"] and not stage["condition"](current_output, context):
stage_info["condition_met"] = False
stage_info["skipped"] = True
execution_log.append(stage_info)
continue
# 执行当前阶段
try:
result = await run_agent_async(
stage["agent"],
current_output,
context
)
current_output = result
stage_info["output"] = current_output
stage_info["success"] = True
except Exception as e:
stage_info["error"] = str(e)
stage_info["success"] = False
# 可以选择继续或中断流水线
execution_log.append(stage_info)
return {
"final_output": current_output,
"execution_log": execution_log
}
模式三:条件分支 根据输入内容或中间结果选择不同的处理路径。
class ConditionalRouter:
"""条件路由系统"""
def __init__(self):
self.routes = []
def add_route(self, condition_func, target_agent: Agent):
"""添加路由规则"""
self.routes.append({
"condition": condition_func,
"agent": target_agent
})
def add_default_route(self, default_agent: Agent):
"""添加默认路由"""
self.default_agent = default_agent
async def route_and_process(self, user_input: str, context: Dict) -> str:
"""根据条件路由并处理"""
# 检查所有路由条件
for route in self.routes:
if route["condition"](user_input, context):
# 找到匹配的路由,使用对应的智能体
return await run_agent_async(
route["agent"],
user_input,
context
)
# 没有匹配的路由,使用默认智能体
if hasattr(self, "default_agent"):
return await run_agent_async(
self.default_agent,
user_input,
context
)
# 没有默认路由,返回错误
return "抱歉,我无法处理这个请求。"
2.4 智能体团队的性能优化
当智能体团队规模增长时,性能优化变得至关重要:
1. 智能体懒加载 不是所有智能体都需要在启动时立即初始化:
class LazyAgentLoader:
"""智能体懒加载器"""
def __init__(self):
self.agent_registry = {}
self.loaded_agents = {}
def register_agent(self, name: str, factory_func):
"""注册智能体工厂函数"""
self.agent_registry[name] = factory_func
def get_agent(self, name: str) -> Agent:
"""获取智能体(按需加载)"""
if name not in self.loaded_agents:
if name in self.agent_registry:
# 第一次访问时加载
self.loaded_agents[name] = self.agent_registry[name]()
else:
raise ValueError(f"未注册的智能体: {name}")
return self.loaded_agents[name]
def preload_frequent_agents(self, agent_names: List[str]):
"""预加载高频使用的智能体"""
for name in agent_names:
if name in self.agent_registry:
_ = self.get_agent(name) # 触发加载
2. 结果缓存 对于相同或相似的查询,缓存智能体的响应:
import hashlib
import json
from datetime import datetime, timedelta
class AgentResponseCache:
"""智能体响应缓存"""
def __init__(self, ttl_seconds: int = 300): # 默认5分钟
self.cache = {}
self.ttl = timedelta(seconds=ttl_seconds)
def _generate_key(self, agent_name: str, input_text: str, context: Dict) -> str:
"""生成缓存键"""
# 基于智能体名称、输入和上下文生成唯一键
key_data = {
"agent": agent_name,
"input": input_text,
"context_hash": hash(json.dumps(context, sort_keys=True))
}
key_str = json.dumps(key_data, sort_keys=True)
return hashlib.md5(key_str.encode()).hexdigest()
def get(self, agent_name: str, input_text: str, context: Dict):
"""获取缓存响应"""
key = self._generate_key(agent_name, input_text, context)
if key in self.cache:
entry = self.cache[key]
# 检查是否过期
if datetime.now() - entry["timestamp"] < self.ttl:
return entry["response"]
else:
# 过期,删除
del self.cache[key]
return None
def set(self, agent_name: str, input_text: str, context: Dict, response):
"""设置缓存响应"""
key = self._generate_key(agent_name, input_text, context)
self.cache[key] = {
"response": response,
"timestamp": datetime.now()
}
# 简单的缓存清理(实际中可能需要更复杂的策略)
if len(self.cache) > 1000: # 限制缓存大小
# 删除最旧的条目
oldest_key = min(self.cache.keys(),
key=lambda k: self.cache[k]["timestamp"])
del self.cache[oldest_key]
3. 并发限制 控制同时运行的智能体数量,避免资源耗尽:
import asyncio
from asyncio import Semaphore
class ConcurrentAgentExecutor:
"""带并发限制的智能体执行器"""
def __init__(self, max_concurrent: int = 5):
self.semaphore = Semaphore(max_concurrent)
self.active_tasks = 0
self.max_concurrent = max_concurrent
async def execute_agent(self, agent: Agent, input_text: str, context: Dict) -> str:
"""执行智能体,带并发控制"""
async with self.semaphore:
self.active_tasks += 1
try:
# 实际执行智能体
result = await self._run_agent_impl(agent, input_text, context)
return result
finally:
self.active_tasks -= 1
async def _run_agent_impl(self, agent: Agent, input_text: str, context: Dict) -> str:
"""实际执行智能体的实现"""
# 这里调用ADK的Runner
# 简化实现
await asyncio.sleep(0.1) # 模拟处理时间
return f"Processed by {agent.name}: {input_text}"
def get_utilization(self) -> float:
"""获取当前利用率"""
if self.max_concurrent > 0:
return self.active_tasks / self.max_concurrent
return 0.0
3. 实战案例:华氏/摄氏自动转换系统
现在让我们将这些概念应用于一个完整的实战案例:构建一个能够自动检测、记忆和适应用户温度单位偏好的天气服务系统。
3.1 系统架构设计
我们的系统将包含以下组件:
- 偏好检测器:从用户输入中识别温度单位偏好
- 状态管理器:持久化存储用户偏好
- 单位转换器:在摄氏度和华氏度之间转换
- 智能路由:根据上下文选择正确的处理逻辑
- 响应格式化器:以用户偏好的格式呈现结果
以下是完整的实现:
import re
from enum import Enum
from dataclasses import dataclass
from typing import Optional, Dict, Any
from datetime import datetime
class TemperatureUnit(Enum):
"""温度单位枚举"""
CELSIUS = "celsius"
FAHRENHEIT = "fahrenheit"
KELVIN = "kelvin" # 扩展支持
@dataclass
class Temperature:
"""温度值封装类"""
value: float
unit: TemperatureUnit
def convert_to(self, target_unit: TemperatureUnit) -> 'Temperature':
"""转换到目标单位"""
if self.unit == target_unit:
return self
# 先将所有单位转换为摄氏度作为中间单位
if self.unit == TemperatureUnit.CELSIUS:
celsius = self.value
elif self.unit == TemperatureUnit.FAHRENHEIT:
celsius = (self.value - 32) * 5/9
elif self.unit == TemperatureUnit.KELVIN:
celsius = self.value - 273.15
else:
raise ValueError(f"不支持的源单位: {self.unit}")
# 从摄氏度转换到目标单位
if target_unit == TemperatureUnit.CELSIUS:
converted = celsius
elif target_unit == TemperatureUnit.FAHRENHEIT:
converted = celsius * 9/5 + 32
elif target_unit == TemperatureUnit.KELVIN:
converted = celsius + 273.15
else:
raise ValueError(f"不支持的目标单位: {target_unit}")
return Temperature(value=round(converted, 1), unit=target_unit)
def __str__(self) -> str:
"""格式化显示"""
unit_symbol = {
TemperatureUnit.CELSIUS: "°C",
TemperatureUnit.FAHRENHEIT: "°F",
TemperatureUnit.KELVIN: "K"
}
return f"{self.value}{unit_symbol[self.unit]}"
class PreferenceDetector:
"""用户偏好检测器"""
# 温度单位识别模式
CELSIUS_PATTERNS = [
r"摄氏度", r"摄氏", r"°C", r"℃", r"\bC\b",
r"摄氏温度", r"摄氏制", r"用摄氏", r"显示为摄氏"
]
FAHRENHEIT_PATTERNS = [
r"华氏度", r"华氏", r"°F", r"℉", r"\bF\b",
r"华氏温度", r"华氏制", r"用华氏", r"显示为华氏"
]
@classmethod
def detect_temperature_unit(cls, text: str) -> Optional[TemperatureUnit]:
"""从文本中检测温度单位偏好"""
text_lower = text.lower()
# 检查摄氏度模式
for pattern in cls.CELSIUS_PATTERNS:
if re.search(pattern, text_lower, re.IGNORECASE):
return TemperatureUnit.CELSIUS
# 检查华氏度模式
for pattern in cls.FAHRENHEIT_PATTERNS:
if re.search(pattern, text_lower, re.IGNORECASE):
return TemperatureUnit.FAHRENHEIT
# 检查数字后的单位符号
# 例如:"25C" 或 "77F"
unit_match = re.search(r'(\d+(?:\.\d+)?)\s*([CFK])', text, re.IGNORECASE)
if unit_match:
unit_char = unit_match.group(2).upper()
if unit_char == 'C':
return TemperatureUnit.CELSIUS
elif unit_char == 'F':
return TemperatureUnit.FAHRENHEIT
elif unit_char == 'K':
return TemperatureUnit.KELVIN
return None
@classmethod
def extract_temperature_values(cls, text: str) -> List[Dict[str, Any]]:
"""从文本中提取温度数值和单位"""
# 匹配模式:数字后跟可选的空格和单位符号
pattern = r'(\d+(?:\.\d+)?)\s*°?\s*([CFK])'
matches = re.finditer(pattern, text, re.IGNORECASE)
results = []
for match in matches:
value = float(match.group(1))
unit_char = match.group(2).upper()
if unit_char == 'C':
unit = TemperatureUnit.CELSIUS
elif unit_char == 'F':
unit = TemperatureUnit.FAHRENHEIT
elif unit_char == 'K':
unit = TemperatureUnit.KELVIN
else:
continue
results.append({
"value": value,
"unit": unit,
"original_text": match.group(0),
"position": match.start()
})
return results
class UserPreferenceManager:
"""用户偏好管理器"""
def __init__(self, session_service):
self.session_service = session_service
def get_temperature_unit(self, app_name: str, user_id: str, session_id: str) -> TemperatureUnit:
"""获取用户的温度单位偏好"""
session = self.session_service.get_session(app_name, user_id, session_id)
if session and session.state:
unit_str = session.state.get("temperature_unit")
if unit_str:
try:
return TemperatureUnit(unit_str.lower())
except ValueError:
pass
# 默认使用摄氏度
return TemperatureUnit.CELSIUS
def set_temperature_unit(self, app_name: str, user_id: str,
session_id: str, unit: TemperatureUnit) -> bool:
"""设置用户的温度单位偏好"""
session = self.session_service.get_session(app_name, user_id, session_id)
if session and session.state is not None:
session.state["temperature_unit"] = unit.value
session.state["preference_updated_at"] = datetime.now().isoformat()
session.state["preference_source"] = "explicit_setting"
return True
return False
def detect_and_update_preference(self, app_name: str, user_id: str,
session_id: str, user_input: str) -> Optional[TemperatureUnit]:
"""检测并更新用户偏好"""
detected_unit = PreferenceDetector.detect_temperature_unit(user_input)
if detected_unit:
success = self.set_temperature_unit(
app_name, user_id, session_id, detected_unit
)
if success:
return detected_unit
return None
class TemperatureConverter:
"""温度转换器"""
@staticmethod
def convert_all_units(temp_celsius: float) -> Dict[str, Temperature]:
"""将摄氏度转换为所有单位"""
return {
"celsius": Temperature(temp_celsius, TemperatureUnit.CELSIUS),
"fahrenheit": Temperature(temp_celsius * 9/5 + 32, TemperatureUnit.FAHRENHEIT),
"kelvin": Temperature(temp_celsius + 273.15, TemperatureUnit.KELVIN)
}
@staticmethod
def format_for_display(temp_celsius: float, preferred_unit: TemperatureUnit) -> str:
"""以用户偏好的单位格式化温度显示"""
all_units = TemperatureConverter.convert_all_units(temp_celsius)
preferred_temp = all_units[preferred_unit.value]
# 构建包含所有单位的显示字符串
display_parts = [str(preferred_temp)]
# 添加其他单位作为参考
for unit_name, temp in all_units.items():
if temp.unit != preferred_unit:
display_parts.append(f"({str(temp)})")
return " ".join(display_parts)
@staticmethod
def get_temperature_description(temp_celsius: float) -> str:
"""根据温度值获取描述性文字"""
if temp_celsius < 0:
return "极寒"
elif temp_celsius < 10:
return "寒冷"
elif temp_celsius < 20:
return "凉爽"
elif temp_celsius < 30:
return "温暖"
else:
return "炎热"
class SmartWeatherAgent:
"""智能天气代理(整合所有功能)"""
def __init__(self, session_service):
self.session_service = session_service
self.preference_manager = UserPreferenceManager(session_service)
self.converter = TemperatureConverter()
# 模拟天气数据库
self.weather_db = {
"New York": {"temp_c": 22, "condition": "晴朗", "humidity": 50},
"London": {"temp_c": 15, "condition": "雨天", "humidity": 80},
"Tokyo": {"temp_c": 26, "condition": "多云", "humidity": 65},
"Sydney": {"temp_c": 28, "condition": "晴朗", "humidity": 45},
"Paris": {"temp_c": 18, "condition": "阴天", "humidity": 70}
}
async def process_query(self, app_name: str, user_id: str,
session_id: str, user_input: str) -> Dict[str, Any]:
"""处理用户查询"""
# 1. 检测并更新偏好
detected_unit = self.preference_manager.detect_and_update_preference(
app_name, user_id, session_id, user_input
)
# 2. 获取当前偏好
current_unit = self.preference_manager.get_temperature_unit(
app_name, user_id, session_id
)
# 3. 提取城市名称(简化实现)
city = self._extract_city(user_input)
# 4. 查询天气
weather_info = await self._get_weather(city)
# 5. 格式化响应
response = self._format_response(
city, weather_info, current_unit, detected_unit
)
# 6. 更新会话历史
self._update_session_history(
app_name, user_id, session_id, city, weather_info
)
return response
def _extract_city(self, text: str) -> Optional[str]:
"""从文本中提取城市名称(简化版)"""
# 实际中应该使用更复杂的NLP或实体识别
city_keywords = ["纽约", "伦敦", "东京", "悉尼", "巴黎"]
for city in city_keywords:
if city in text:
return city
return None
async def _get_weather(self, city: str) -> Optional[Dict[str, Any]]:
"""获取天气信息"""
# 城市名称映射
city_map = {
"纽约": "New York",
"伦敦": "London",
"东京": "Tokyo",
"悉尼": "Sydney",
"巴黎": "Paris"
}
english_city = city_map.get(city)
if not english_city or english_city not in self.weather_db:
return None
return self.weather_db[english_city]
def _format_response(self, city: str, weather_info: Dict[str, Any],
current_unit: TemperatureUnit,
detected_unit: Optional[TemperatureUnit]) -> Dict[str, Any]:
"""格式化响应"""
if not weather_info:
return {
"success": False,
"message": f"抱歉,找不到{city}的天气信息。"
}
# 温度转换和格式化
temp_c = weather_info["temp_c"]
formatted_temp = self.converter.format_for_display(temp_c, current_unit)
temp_description = self.converter.get_temperature_description(temp_c)
# 构建响应
response = {
"success": True,
"city": city,
"temperature": {
"value": temp_c,
"display": formatted_temp,
"description": temp_description,
"unit": current_unit.value
},
"condition": weather_info["condition"],
"humidity": weather_info["humidity"],
"timestamp": datetime.now().isoformat()
}
# 如果检测到新的偏好,添加确认信息
if detected_unit:
response["preference_updated"] = True
response["new_unit"] = detected_unit.value
response["confirmation"] = f"已为您设置为使用{detected_unit.value}显示温度"
else:
response["preference_updated"] = False
return response
def _update_session_history(self, app_name: str, user_id: str,
session_id: str, city: str, weather_info: Dict):
"""更新会话历史"""
session = self.session_service.get_session(app_name, user_id, session_id)
if session and session.state is not None:
# 初始化历史记录(如果不存在)
if "query_history" not in session.state:
session.state["query_history"] = []
# 添加新记录
history_entry = {
"city": city,
"temperature_c": weather_info["temp_c"],
"condition": weather_info["condition"],
"timestamp": datetime.now().isoformat()
}
session.state["query_history"].append(history_entry)
# 限制历史记录长度
if len(session.state["query_history"]) > 50:
session.state["query_history"] = session.state["query_history"][-50:]
# 更新最后查询城市
session.state["last_queried_city"] = city
3.2 集成到ADK智能体
现在我们将这个系统集成到ADK智能体中:
from google.adk.agents import Agent
from google.adk.tools import BaseTool, ToolContext
from google.adk.runners import Runner
class TemperatureAwareWeatherTool(BaseTool):
"""温度感知的天气工具"""
def __init__(self, weather_agent: SmartWeatherAgent):
super().__init__(
name="get_weather_with_preference",
description="获取考虑用户偏好的天气信息"
)
self.weather_agent = weather_agent
async def run(self, city: str, tool_context: ToolContext = None) -> Dict[str, Any]:
"""工具执行方法"""
# 从上下文中提取会话信息
app_name = "weather_app"
user_id = "default_user" # 实际中应从上下文中获取
session_id = "default_session" # 实际中应从上下文中获取
# 构建模拟用户输入
user_input = f"{city}的天气"
# 处理查询
result = await self.weather_agent.process_query(
app_name, user_id, session_id, user_input
)
# 更新工具上下文中的状态
if tool_context and tool_context.state is not None:
if result.get("preference_updated"):
tool_context.state["last_preference_update"] = datetime.now().isoformat()
# 记录查询历史
if "query_history" not in tool_context.state:
tool_context.state["query_history"] = []
tool_context.state["query_history"].append({
"query": user_input,
"result": result,
"timestamp": datetime.now().isoformat()
})
return result
# 创建智能体
def create_temperature_aware_weather_agent(session_service) -> Agent:
"""创建温度感知的天气智能体"""
# 初始化天气代理
weather_agent = SmartWeatherAgent(session_service)
# 创建工具
weather_tool = TemperatureAwareWeatherTool(weather_agent)
# 创建智能体
agent = Agent(
name="TemperatureAwareWeatherAgent",
description="能够记住用户温度偏好的智能天气助手",
model="gemini-2.0-flash",
tools=[weather_tool],
instruction="""
你是一个智能天气助手,具有以下特殊能力:
1. **温度偏好记忆**:能够记住用户喜欢的温度单位(摄氏度或华氏度)
2. **自动检测**:能从用户对话中自动检测温度单位偏好
3. **智能转换**:自动在所有温度单位间转换
4. **个性化响应**:以用户偏好的格式显示天气信息
核心工作流程:
1. 当用户查询天气时,使用get_weather_with_preference工具
2. 工具会自动处理温度单位转换和偏好记忆
3. 如果检测到新的温度偏好,在响应中确认
4. 始终以用户偏好的单位显示温度
温度偏好检测规则:
- 当用户提到"摄氏度"、"摄氏"、"°C"、"C"时,设置为摄氏度
- 当用户提到"华氏度"、"华氏"、"°F"、"F"时,设置为华氏度
- 偏好设置是持久化的,会记住在整个会话中
响应格式示例:
用户:"我喜欢用华氏度"
你:"好的,已为您设置为使用华氏度显示温度。"
用户:"纽约天气怎么样?"
你:"纽约当前天气晴朗,温度72°F(22.2°C),湿度50%。感觉温暖舒适。"
用户:"切换到摄氏度显示"
你:"已切换到摄氏度显示。"
用户:"伦敦呢?"
你:"伦敦当前雨天,温度15.6°C(60°F),湿度80%。感觉凉爽。"
重要提示:
- 不要询问用户偏好,从对话中自动检测
- 每次响应都要显示用户偏好的温度单位
- 同时显示其他单位作为参考(括号中)
- 提供温度的感觉描述(极寒、寒冷、凉爽、温暖、炎热)
""",
output_key="weather_response"
)
return agent
3.3 测试与验证
让我们创建测试脚本来验证系统的功能:
import asyncio
from google.adk.sessions import InMemorySessionService
async def test_temperature_preference_system():
"""测试温度偏好系统"""
# 初始化服务
session_service = InMemorySessionService()
# 创建智能体
agent = create_temperature_aware_weather_agent(session_service)
# 创建Runner
runner = Runner(
agent=agent,
session_service=session_service,
app_name="weather_app"
)
# 测试用例
test_cases = [
{
"input": "纽约天气怎么样?",
"description": "初始查询,应使用默认摄氏度"
},
{
"input": "我喜欢用华氏度",
"description": "设置偏好为华氏度"
},
{
"input": "伦敦的天气呢?",
"description": "第二次查询,应使用华氏度"
},
{
"input": "切换到摄氏度",
"description": "更改偏好为摄氏度"
},
{
"input": "东京天气",
"description": "第三次查询,应使用摄氏度"
},
{
"input": "用华氏度显示悉尼天气",
"description": "查询时指定单位"
}
]
print("=== 温度偏好系统测试 ===")
print()
# 执行测试
for i, test_case in enumerate(test_cases, 1):
print(f"测试 {i}: {test_case['description']}")
print(f"输入: {test_case['input']}")
# 运行智能体
response = await runner.run_async(
user_id="test_user",
session_id="test_session",
new_message=test_case["input"]
)
# 提取响应
if hasattr(response, 'weather_response'):
print(f"响应: {response.weather_response}")
elif hasattr(response, 'content'):
print(f"响应: {response.content}")
# 显示当前会话状态
session = session_service.get_session(
"weather_app", "test_user", "test_session"
)
if session and session.state:
current_unit = session.state.get("temperature_unit", "celsius")
print(f"当前温度单位偏好: {current_unit}")
print("-" * 50)
print()
async def test_edge_cases():
"""测试边界情况"""
session_service = InMemorySessionService()
agent = create_temperature_aware_weather_agent(session_service)
runner = Runner(
agent=agent,
session_service=session_service,
app_name="weather_app"
)
print("=== 边界情况测试 ===")
print()
# 测试1:混合单位查询
print("测试1: 混合单位查询")
response = await runner.run_async(
user_id="edge_user",
session_id="edge_session",
new_message="今天25C,明天77F,哪个更热?"
)
print(f"响应: {getattr(response, 'weather_response', getattr(response, 'content', '无响应'))}")
print()
# 测试2:无效城市
print("测试2: 无效城市查询")
response = await runner.run_async(
user_id="edge_user",
session_id="edge_session",
new_message="火星天气怎么样?"
)
print(f"响应: {getattr(response, 'weather_response', getattr(response, 'content', '无响应'))}")
print()
# 测试3:复杂偏好表达
print("测试3: 复杂偏好表达")
test_inputs = [
"请用华氏温度单位",
"我要看摄氏度的",
"显示为°F",
"用C表示温度"
]
for input_text in test_inputs:
response = await runner.run_async(
user_id="edge_user",
session_id="edge_session",
new_message=input_text
)
print(f"输入: {input_text}")
print(f"响应: {getattr(response, 'weather_response', getattr(response, 'content', '无响应'))}")
print()
async def main():
"""主测试函数"""
print("开始温度偏好智能体系统测试")
print("=" * 60)
await test_temperature_preference_system()
await test_edge_cases()
print("测试完成!")
if __name__ == "__main__":
asyncio.run(main())
3.4 性能监控与优化
对于生产系统,我们需要监控智能体的性能:
import time
from dataclasses import dataclass, field
from typing import Dict, List, Optional
from datetime import datetime
import statistics
@dataclass
class PerformanceMetrics:
"""性能指标收集"""
total_queries: int = 0
successful_queries: int = 0
failed_queries: int = 0
total_response_time: float = 0.0
response_times: List[float] = field(default_factory=list)
preference_changes: int = 0
cache_hits: int = 0
cache_misses: int = 0
def record_query(self, success: bool, response_time: float):
"""记录查询指标"""
self.total_queries += 1
if success:
self.successful_queries += 1
else:
self.failed_queries += 1
self.total_response_time += response_time
self.response_times.append(response_time)
def record_preference_change(self):
"""记录偏好变更"""
self.preference_changes += 1
def record_cache_hit(self):
"""记录缓存命中"""
self.cache_hits += 1
def record_cache_miss(self):
"""记录缓存未命中"""
self.cache_misses += 1
@property
def success_rate(self) -> float:
"""计算成功率"""
if self.total_queries == 0:
return 0.0
return self.successful_queries / self.total_queries
@property
def average_response_time(self) -> float:
"""计算平均响应时间"""
if self.total_queries == 0:
return 0.0
return self.total_response_time / self.total_queries
@property
def p95_response_time(self) -> float:
"""计算95百分位响应时间"""
if not self.response_times:
return 0.0
return statistics.quantiles(self.response_times, n=20)[18] # 95th percentile
def get_summary(self) -> Dict[str, any]:
"""获取性能摘要"""
return {
"total_queries": self.total_queries,
"success_rate": f"{self.success_rate:.1%}",
"average_response_time_ms": f"{self.average_response_time * 1000:.1f}",
"p95_response_time_ms": f"{self.p95_response_time * 1000:.1f}",
"preference_changes": self.preference_changes,
"cache_hit_rate": f"{self.cache_hits / max(self.cache_hits + self.cache_misses, 1):.1%}",
"timestamp": datetime.now().isoformat()
}
class MonitoredSmartWeatherAgent(SmartWeatherAgent):
"""带监控的智能天气代理"""
def __init__(self, session_service):
super().__init__(session_service)
self.metrics = PerformanceMetrics()
self.query_cache = {}
async def process_query(self, app_name: str, user_id: str,
session_id: str, user_input: str) -> Dict[str, Any]:
"""带监控的查询处理"""
start_time = time.time()
# 检查缓存
cache_key = f"{user_id}:{session_id}:{user_input}"
if cache_key in self.query_cache:
self.metrics.record_cache_hit()
result = self.query_cache[cache_key]
response_time = time.time() - start_time
self.metrics.record_query(True, response_time)
return result
self.metrics.record_cache_miss()
try:
# 调用父类方法
result = await super().process_query(
app_name, user_id, session_id, user_input
)
# 记录偏好变更
if result.get("preference_updated"):
self.metrics.record_preference_change()
# 缓存结果(仅缓存成功查询)
if result.get("success"):
self.query_cache[cache_key] = result
# 简单缓存清理
if len(self.query_cache) > 100:
# 移除最旧的条目
oldest_key = next(iter(self.query_cache))
del self.query_cache[oldest_key]
response_time = time.time() - start_time
self.metrics.record_query(True, response_time)
return result
except Exception as e:
response_time = time.time() - start_time
self.metrics.record_query(False, response_time)
return {
"success": False,
"error": str(e),
"timestamp": datetime.now().isoformat()
}
def get_performance_report(self) -> Dict[str, any]:
"""获取性能报告"""
return self.metrics.get_summary()
4. 生产环境部署与最佳实践
构建带记忆的智能体团队不仅仅是技术实现,还需要考虑生产环境的实际需求。以下是部署和运维的关键考虑因素。
4.1 部署架构
对于生产环境,建议采用以下架构:
┌─────────────────┐ ┌─────────────────┐ ┌─────────────────┐
│ 客户端请求 │───▶│ API网关/负载 │───▶│ 智能体路由层 │
│ (Web/移动端) │ │ 均衡器 │ │ │
└─────────────────┘ └─────────────────┘ └────────┬────────┘
│
▼
┌─────────────────┐ ┌─────────────────┐ ┌─────────────────┐
│ 会话存储 │◀───│ 智能体执行器 │───▶│ 大模型API │
│ (Redis/DB) │ │ │ │ (Gemini/等) │
└─────────────────┘ └────────┬────────┘ └─────────────────┘
│
▼
┌─────────────────┐ ┌─────────────────┐
│ 监控与日志 │◀───│ 结果处理与 │
│ (Prometheus/ │ │ 格式化 │
│ ELK) │ └─────────────────┘
└─────────────────┘
4.2 配置管理
使用配置文件管理不同环境的设置:
# config/production.yaml
agent_system:
# 模型配置
models:
default: "gemini-2.0-flash"
fallback: "gpt-4o"
backup: "claude-3-sonnet"
# 会话配置
session:
storage: "redis" # 或 "database", "memory"
ttl: 3600 # 会话过期时间(秒)
cleanup_interval: 300 # 清理间隔(秒)
# 缓存配置
cache:
enabled: true
type: "redis" # 或 "memory", "memcached"
ttl: 300 # 缓存过期时间(秒)
max_size: 10000 # 最大缓存条目数
# 性能配置
performance:
max_concurrent_agents: 10
request_timeout: 30 # 秒
retry_attempts: 3
circuit_breaker:
failure_threshold: 5
reset_timeout: 60
# 监控配置
monitoring:
enabled: true
metrics_port: 9090
log_level: "INFO"
trace_sampling_rate: 0.1
# 安全配置
security:
input_validation: true
output_sanitization: true
rate_limiting:
enabled: true
requests_per_minute: 100
content_filter:
enabled: true
level: "moderate"
# 智能体团队配置
agents:
weather_agent:
enabled: true
tools: ["weather_lookup", "unit_conversion"]
model: "gemini-2.0-flash"
temperature: 0.7
max_tokens: 1000
greeting_agent:
enabled: true
tools: ["greeting"]
model: "gemini-2.0-flash"
temperature: 0.9
max_tokens: 500
routing_agent:
enabled: true
tools: []
model: "gemini-2.0-flash"
temperature: 0.3
max_tokens: 500
4.3 错误处理与弹性
健壮的生产系统需要完善的错误处理机制:
from typing import Optional, Dict, Any
import asyncio
from functools import wraps
import logging
logger = logging.getLogger(__name__)
class CircuitBreaker:
"""断路器模式实现"""
def __init__(self, failure_threshold: int = 5, reset_timeout: int = 60):
self.failure_threshold = failure_threshold
self.reset_timeout = reset_timeout
self.failure_count = 0
self.last_failure_time = None
self.state = "CLOSED" # CLOSED, OPEN, HALF_OPEN
def can_execute(self) -> bool:
"""检查是否允许执行"""
if self.state == "OPEN":
# 检查是否应该尝试恢复
if self.last_failure_time:
time_since_failure = time.time() - self.last_failure_time
if time_since_failure > self.reset_timeout:
self.state = "HALF_OPEN"
return True
return False
return True
def record_success(self):
"""记录成功"""
if self.state == "HALF_OPEN":
self.state = "CLOSED"
self.failure_count = 0
self.last_failure_time = None
def record_failure(self):
"""记录失败"""
self.failure_count += 1
self.last_failure_time = time.time()
if self.failure_count >= self.failure_threshold:
self.state = "OPEN"
logger.warning(f"断路器打开,失败次数: {self.failure_count}")
def retry_with_backoff(max_retries: int = 3, base_delay: float = 1.0):
"""带指数退避的重试装饰器"""
def decorator(func):
@wraps(func)
async def wrapper(*args, **kwargs):
last_exception = None
for attempt in range(max_retries + 1):
try:
return await func(*args, **kwargs)
except Exception as e:
last_exception = e
if attempt == max_retries:
break
# 计算退避时间
delay = base_delay * (2 ** attempt)
jitter = delay * 0.1 # 添加10%的抖动
actual_delay = delay + random.uniform(-jitter, jitter)
logger.warning(
f"尝试 {func.__name__} 失败 (尝试 {attempt + 1}/{max_retries + 1}): "
f"{str(e)}. {actual_delay:.2f}秒后重试..."
)
await asyncio.sleep(actual_delay)
# 所有重试都失败
logger.error(
f"{func.__name__} 在{max_retries + 1}次尝试后失败: {str(last_exception)}"
)
raise last_exception
return wrapper
return decorator
class ResilientAgentSystem:
"""具有弹性的智能体系统"""
def __init__(self, config: Dict[str, Any]):
self.config = config
self.circuit_breakers = {}
self.fallback_agents = {}
# 初始化断路器
for agent_name in config.get("agents", {}):
self.circuit_breakers[agent_name] = CircuitBreaker(
failure_threshold=config.get("circuit_breaker", {}).get("failure_threshold", 5),
reset_timeout=config.get("circuit_breaker", {}).get("reset_timeout", 60)
)
@retry_with_backoff(max_retries=3, base_delay=1.0)
async def execute_agent_with_resilience(
self,
agent_name: str,
input_text: str,
context: Dict[str, Any]
) -> Dict[str, Any]:
"""带弹性的智能体执行"""
# 检查断路器
circuit_breaker = self.circuit_breakers.get(agent_name)
if circuit_breaker and not circuit_breaker.can_execute():
logger.warning(f"断路器阻止执行 {agent_name},使用降级策略")
return await self._execute_fallback(agent_name, input_text, context)
try:
# 执行主智能体
result = await self._execute_primary_agent(agent_name, input_text, context)
# 记录成功
if circuit_breaker:
circuit_breaker.record_success()
return result
except Exception as e:
# 记录失败
if circuit_breaker:
circuit_breaker.record_failure()
logger.error(f"智能体 {agent_name} 执行失败: {str(e)}")
# 尝试降级
try:
return await self._execute_fallback(agent_name, input_text, context)
except Exception as fallback_error:
logger.error(f"降级策略也失败: {str(fallback_error)}")
raise
async def _execute_primary_agent(
self,
agent_name: str,
input_text: str,
context: Dict[str, Any]
) -> Dict[str, Any]:
"""执行主智能体"""
# 这里调用实际的智能体执行逻辑
# 简化实现
await asyncio.sleep(0.1) # 模拟处理时间
# 模拟随机失败(测试用)
if random.random() < 0.1: # 10%失败率
raise Exception("模拟智能体执行失败")
return {
"success": True,
"agent": agent_name,
"response": f"主智能体处理: {input_text}",
"timestamp": datetime.now().isoformat()
}
async def _execute_fallback(
self,
agent_name: str,
input_text: str,
context: Dict[str, Any]
) -> Dict[str, Any]:
"""执行降级智能体"""
# 这里调用降级智能体
# 简化实现
await asyncio.sleep(0.05) # 降级服务通常更快
return {
"success": True,
"agent": f"{agent_name}_fallback",
"response": f"降级智能体处理: {input_text}",
"timestamp": datetime.now().isoformat(),
"degraded": True
}
4.4 监控与可观测性
生产系统需要全面的监控:
import prometheus_client
from prometheus_client import Counter, Histogram, Gauge
import time
from typing import Dict, Any
class AgentMetrics:
"""智能体指标收集"""
def __init__(self):
# 请求计数器
self.requests_total = Counter(
'agent_requests_total',
'Total number of agent requests',
['agent_name', 'status']
)
# 响应时间直方图
self.request_duration = Histogram(
'agent_request_duration_seconds',
'Agent request duration in seconds',
['agent_name'],
buckets=[0.1, 0.5, 1.0, 2.0, 5.0, 10.0]
)
# 活跃请求数
self.requests_in_progress = Gauge(
'agent_requests_in_progress',
'Number of agent requests in progress',
['agent_name']
)
# 缓存指标
self.cache_hits = Counter(
'agent_cache_hits_total',
'Total number of cache hits',
['agent_name']
)
self.cache_misses = Counter(
'agent_cache_misses_total',
'Total number of cache misses',
['agent_name']
)
# 错误计数器
self.errors_total = Counter(
'agent_errors_total',
'Total number of agent errors',
['agent_name', 'error_type']
)
def record_request_start(self, agent_name: str):
"""记录请求开始"""
self.requests_in_progress.labels(agent_name=agent_name).inc()
def record_request_end(
self,
agent_name: str,
success: bool,
duration: float,
used_cache: bool = False
):
"""记录请求结束"""
self.requests_in_progress.labels(agent_name=agent_name).dec()
status = "success" if success else "failure"
self.requests_total.labels(
agent_name=agent_name,
status=status
).inc()
self.request_duration.labels(agent_name=agent_name).observe(duration)
if used_cache:
self.cache_hits.labels(agent_name=agent_name).inc()
else:
self.cache_misses.labels(agent_name=agent_name).inc()
def record_error(self, agent_name: str, error_type: str):
"""记录错误"""
self.errors_total.labels(
agent_name=agent_name,
error_type=error_type
).inc()
class MonitoredAgentExecutor:
"""带监控的智能体执行器"""
def __init__(self, metrics: AgentMetrics):
self.metrics = metrics
self.agent_instances = {}
async def execute_agent(
self,
agent_name: str,
input_text: str,
context: Dict[str, Any]
) -> Dict[str, Any]:
"""执行智能体并记录指标"""
start_time = time.time()
self.metrics.record_request_start(agent_name)
try:
# 检查缓存
cache_key = self._generate_cache_key(agent_name, input_text, context)
cached_result = self._get_from_cache(cache
更多推荐
所有评论(0)