程序员实战:多智能体 React 模式层级指挥优化与代码实现
在多智能体系统开发中,传统 “单智能体独立决策” 模式常面临任务协同混乱、资源浪费、响应延迟等问题。而基于 React 模式(思考 - 行动 - 反馈)的层级指挥架构,通过 “指挥者 - 执行者” 分层设计,能实现任务高效分发与协同。本文以程序员视角,结合 Python 实战代码,详解多智能体 React 模式的核心逻辑、层级指挥实现及性能优化技巧,帮你构建稳定高效的多智能体系统。
React 模式核心:多智能体的思考 - 行动闭环
多智能体 React 模式的本质是为每个智能体赋予 “思考(Think)- 行动(Act)- 反馈(Feedback)” 的闭环能力 —— 通过思考阶段分析任务目标、生成行动计划,行动阶段调用工具或执行任务,反馈阶段根据结果调整策略。这一闭环是层级指挥的基础,确保智能体不盲目执行,而是动态适配任务需求。
基础 React Agent 实现
以 “数据处理智能体” 为例,我们通过自定义类实现 React 闭环,核心包含think(生成行动步骤)、act(执行工具调用)、feedback(结果分析与策略调整)三个方法:
from typing import List, Dict, Optional
import time
from enum import Enum
# 定义行动类型枚举(工具/任务类型)
class ActionType(Enum):
DATA_LOAD = "data_load" # 数据加载
DATA_CLEAN = "data_clean" # 数据清洗
DATA_ANALYZE = "data_analyze" # 数据分析
RESULT_SAVE = "result_save" # 结果保存
# 定义行动结果类
class ActionResult:
def __init__(self, success: bool, data: Optional[Dict] = None, error: Optional[str] = None):
self.success = success # 行动是否成功
self.data = data # 行动返回数据
self.error = error # 错误信息(失败时)
self.timestamp = time.time() # 执行时间戳
# 基础React智能体类
class ReactAgent:
def __init__(self, agent_id: str, name: str):
self.agent_id = agent_id # 智能体唯一ID
self.name = name # 智能体名称
self.task_history: List[Dict] = [] # 任务执行历史
self.current_task: Optional[Dict] = None # 当前任务
self.available_tools = { # 智能体可用工具
ActionType.DATA_LOAD: self._load_data,
ActionType.DATA_CLEAN: self._clean_data,
ActionType.DATA_ANALYZE: self._analyze_data,
ActionType.RESULT_SAVE: self._save_result
}
def think(self, task: Dict) -> List[ActionType]:
"""
思考阶段:分析任务目标,生成行动步骤
:param task: 任务信息(含目标、参数)
:return: 行动类型列表(执行顺序)
"""
self.current_task = task
task_goal = task.get("goal", "")
print(f"[{self.name}] 思考任务:{task_goal}")
# 根据任务目标生成行动步骤(示例:数据处理任务流程)
if "分析" in task_goal and "数据" in task_goal:
# 数据处理流程:加载→清洗→分析→保存
action_steps = [
ActionType.DATA_LOAD,
ActionType.DATA_CLEAN,
ActionType.DATA_ANALYZE,
ActionType.RESULT_SAVE
]
print(f"[{self.name}] 生成行动步骤:{[step.value for step in action_steps]}")
return action_steps
elif "清洗" in task_goal:
# 仅清洗任务:加载→清洗→保存
return [ActionType.DATA_LOAD, ActionType.DATA_CLEAN, ActionType.RESULT_SAVE]
else:
raise ValueError(f"[{self.name}] 无法识别任务目标:{task_goal}")
def act(self, action_steps: List[ActionType]) -> ActionResult:
"""
行动阶段:按步骤执行工具调用
:param action_steps: 行动步骤列表
:return: 最终行动结果
"""
if not self.current_task:
return ActionResult(success=False, error="无当前任务,需先调用think()")
task_params = self.current_task.get("params", {})
intermediate_results = {} # 中间结果存储
try:
# 按顺序执行每个行动步骤
for step in action_steps:
print(f"\n[{self.name}] 执行行动:{step.value}")
# 获取工具函数并执行
tool_func = self.available_tools.get(step)
if not tool_func:
raise Exception(f"无可用工具:{step.value}")
# 执行工具,传入参数与中间结果
step_result = tool_func(task_params, intermediate_results)
if not step_result.success:
raise Exception(f"行动失败:{step.value},原因:{step_result.error}")
# 保存中间结果(供后续步骤使用)
intermediate_results[step.value] = step_result.data
print(f"[{self.name}] 行动成功:{step.value},中间结果:{list(intermediate_results[step.value].keys())[:2]}...")
# 所有步骤执行完成,返回最终结果
final_result = {
"task_id": self.current_task.get("task_id"),
"intermediate": intermediate_results,
"status": "completed"
}
return ActionResult(success=True, data=final_result)
except Exception as e:
return ActionResult(success=False, error=str(e))
def feedback(self, action_result: ActionResult) -> Dict:
"""
反馈阶段:分析行动结果,调整策略(如重试、优化步骤)
:param action_result: 行动结果
:return: 反馈后的策略调整建议
"""
print(f"\n[{self.name}] 反馈阶段:{'行动成功' if action_result.success else '行动失败'}")
feedback_info = {
"agent_id": self.agent_id,
"task_id": self.current_task.get("task_id"),
"success": action_result.success,
"adjustment": None# 策略调整建议
}
if action_result.success:
# 成功:记录历史,无调整
self.task_history.append({
"task": self.current_task,
"result": action_result.data,
"time": time.strftime("%Y-%m-%d %H:%M:%S")
})
feedback_info["adjustment"] = "任务完成,无需调整"
else:
# 失败:分析原因,生成调整建议
error_msg = action_result.error
if "数据格式错误" in error_msg:
feedback_info["adjustment"] = "建议检查数据格式,重新执行【数据加载→清洗】步骤"
elif "工具调用失败" in error_msg:
feedback_info["adjustment"] = "建议重试工具调用,或切换备用工具"
else:
feedback_info["adjustment"] = f"需人工介入处理:{error_msg}"
print(f"[{self.name}] 反馈建议:{feedback_info['adjustment']}")
return feedback_info
# ------------------------------
# 以下为工具函数(模拟实际数据处理逻辑)
# ------------------------------
def _load_data(self, task_params: Dict, intermediate: Dict) -> ActionResult:
"""工具:加载数据"""
data_path = task_params.get("data_path", "data.csv")
try:
# 模拟数据加载(实际为读取文件/数据库)
time.sleep(0.5)
data = {
"id": [1, 2, 3, 4],
"value": [10, None, 30, 40], # 含缺失值(用于后续清洗)
"category": ["A", "B", "A", "C"]
}
return ActionResult(success=True, data={"raw_data": data, "path": data_path})
except Exception as e:
return ActionResult(success=False, error=f"数据加载失败:{str(e)}")
def _clean_data(self, task_params: Dict, intermediate: Dict) -> ActionResult:
"""工具:清洗数据(处理缺失值、异常值)"""
raw_data = intermediate.get(ActionType.DATA_LOAD.value, {}).get("raw_data")
if not raw_data:
return ActionResult(success=False, error="无原始数据,需先执行数据加载")
try:
# 模拟数据清洗:填充缺失值、去重
time.sleep(0.5)
cleaned_data = raw_data.copy()
# 填充缺失值(用均值)
value_mean = sum(v for v in cleaned_data["value"] if v is not None) / 3
cleaned_data["value"] = [v if v is not None else value_mean for v in cleaned_data["value"]]
return ActionResult(success=True, data={"cleaned_data": cleaned_data})
except Exception as e:
return ActionResult(success=False, error=f"数据清洗失败:{str(e)}")
def _analyze_data(self, task_params: Dict, intermediate: Dict) -> ActionResult:
"""工具:数据分析(统计、聚合)"""
cleaned_data = intermediate.get(ActionType.DATA_CLEAN.value, {}).get("cleaned_data")
if not cleaned_data:
return ActionResult(success=False, error="无清洗后数据,需先执行数据清洗")
try:
# 模拟数据分析:计算均值、分类统计
time.sleep(0.5)
value_mean = sum(cleaned_data["value"]) / len(cleaned_data["value"])
category_count = {}
for cat in cleaned_data["category"]:
category_count[cat] = category_count.get(cat, 0) + 1
analysis_result = {
"value_mean": round(value_mean, 2),
"category_count": category_count,
"sample_count": len(cleaned_data["id"])
}
return ActionResult(success=True, data={"analysis_result": analysis_result})
except Exception as e:
return ActionResult(success=False, error=f"数据分析失败:{str(e)}")
def _save_result(self, task_params: Dict, intermediate: Dict) -> ActionResult:
"""工具:保存结果"""
analysis_result = intermediate.get(ActionType.DATA_ANALYZE.value, {}).get("analysis_result")
if not analysis_result:
# 若无分析结果(仅清洗任务),取清洗后数据
analysis_result = intermediate.get(ActionType.DATA_CLEAN.value, {}).get("cleaned_data")
try:
# 模拟结果保存(实际为写入文件/数据库)
time.sleep(0.5)
save_path = task_params.get("save_path", "result.json")
return ActionResult(success=True, data={"save_path": save_path, "result": analysis_result})
except Exception as e:
return ActionResult(success=False, error=f"结果保存失败:{str(e)}")
# 测试基础React Agent
if __name__ == "__main__":
# 初始化数据处理智能体
data_agent = ReactAgent(agent_id="agent_001", name="数据处理智能体")
# 定义任务:分析销售数据(含目标、参数)
task = {
"task_id": "task_001",
"goal": "分析2025年1月销售数据,计算均值与分类统计",
"params": {
"data_path": "sales_202501.csv",
"save_path": "sales_analysis_202501.json"
}
}
# 执行React闭环
try:
# 1. 思考:生成行动步骤
action_steps = data_agent.think(task)
# 2. 行动:执行步骤
action_result = data_agent.act(action_steps)
# 3. 反馈:分析结果
feedback = data_agent.feedback(action_result)
print(f"\n最终反馈:{feedback}")
except Exception as e:
print(f"执行失败:{str(e)}")
上述代码实现了 React 模式的核心闭环:think方法根据任务目标生成结构化行动步骤,act方法按顺序调用工具执行,feedback方法根据结果生成调整建议。这种设计的优势在于:智能体不依赖固定脚本,而是通过 “思考” 动态适配任务变化,为层级指挥提供了 “可拆解、可协同” 的基础。
层级指挥设计:任务分发与协同控制
当多智能体协同处理复杂任务(如 “市场分析” 需数据处理、报告生成、可视化三个子任务)时,仅靠单个 React Agent 无法高效完成 —— 此时需要 “层级指挥” 架构:由 “指挥者 Agent” 负责任务拆解、智能体分配、进度监控,“执行者 Agent” 负责执行具体子任务并反馈结果。
层级指挥架构实现
我们设计两类智能体:CommanderAgent(指挥者)和ExecutorAgent(执行者,继承自 ReactAgent)。指挥者的核心能力是 “任务拆解→智能体匹配→结果聚合”,执行者专注于子任务的 React 闭环执行。
from typing import Dict, List, Optional
import time
from collections import defaultdict
# 执行者智能体(继承ReactAgent,专注子任务执行)
class ExecutorAgent(ReactAgent):
def __init__(self, agent_id: str, name: str, specialty: str):
super().__init__(agent_id, name)
self.specialty = specialty # 执行者专长(如“数据处理”“报告生成”“可视化”)
# 按专长过滤可用工具(示例:可视化执行者仅保留相关工具)
if self.specialty == "可视化":
self.available_tools = {
ActionType.DATA_LOAD: self._load_data,
ActionType.DATA_ANALYZE: self._analyze_data, # 仅基础分析
ActionType.RESULT_SAVE: self._save_visual_result # 新增可视化工具
}
# 可视化执行者专属工具:生成图表
def _save_visual_result(self, task_params: Dict, intermediate: Dict) -> ActionResult:
"""工具:生成可视化图表并保存"""
analysis_result = intermediate.get(ActionType.DATA_ANALYZE.value, {}).get("analysis_result")
if not analysis_result:
return ActionResult(success=False, error="无分析数据,无法生成可视化")
try:
time.sleep(0.8)
# 模拟生成图表(如柱状图、折线图)
visual_data = {
"chart_type": "bar",
"x_axis": list(analysis_result["category_count"].keys()),
"y_axis": list(analysis_result["category_count"].values()),
"title": task_params.get("chart_title", "数据可视化图表")
}
save_path = task_params.get("visual_path", "visual_chart.png")
return ActionResult(success=True, data={"visual_path": save_path, "chart_data": visual_data})
except Exception as e:
return ActionResult(success=False, error=f"可视化生成失败:{str(e)}")
# 指挥者智能体(负责任务拆解、智能体分配、结果聚合)
class CommanderAgent:
def __init__(self, commander_id: str, name: str):
self.commander_id = commander_id
self.name = name
self.executors: Dict[str, ExecutorAgent] = {} # 管理的执行者列表(key:agent_id)
self.task_queue: List[Dict] = [] # 待分配任务队列
self.task_status: Dict[str, Dict] = defaultdict(dict) # 任务状态(key:task_id)
def register_executor(self, executor: ExecutorAgent):
"""注册执行者智能体"""
self.executors[executor.agent_id] = executor
print(f"[{self.name}] 注册执行者:{executor.name}(专长:{executor.specialty})")
def decompose_task(self, main_task: Dict) -> List[Dict]:
"""
任务拆解:将主任务拆分为子任务
:param main_task: 主任务(如“市场分析”)
:return: 子任务列表
"""
main_goal = main_task.get("goal", "")
print(f"\n[{self.name}] 拆解主任务:{main_goal}")
# 示例:“市场分析”拆分为3个子任务
if "市场分析" in main_goal:
task_params = main_task.get("params", {})
sub_tasks = [
# 子任务1:数据处理(分配给数据处理执行者)
{
"task_id": f"{main_task['task_id']}_sub1",
"goal": "加载并分析市场数据,计算关键指标",
"type": "data_processing",
"params": {
"data_path": task_params.get("data_path", "market_data.csv"),
"save_path": "market_analysis_data.json"
},
"main_task_id": main_task["task_id"]
},
# 子任务2:报告生成(分配给报告执行者)
{
"task_id": f"{main_task['task_id']}_sub2",
"goal": "基于分析数据生成市场分析报告(Markdown格式)",
"type": "report_generation",
"params": {
"data_path": "market_analysis_data.json",
"report_path": "market_analysis_report.md"
},
"main_task_id": main_task["task_id"]
},
# 子任务3:可视化(分配给可视化执行者)
{
"task_id": f"{main_task['task_id']}_sub3",
"goal": "基于分析数据生成市场趋势图表",
"type": "visualization",
"params": {
"data_path": "market_analysis_data.json",
"visual_path": "market_trend_chart.png",
"chart_title": "2025年Q1市场趋势"
},
"main_task_id": main_task["task_id"]
}
]
print(f"[{self.name}] 拆解完成,子任务数:{len(sub_tasks)}")
return sub_tasks
else:
raise ValueError(f"[{self.name}] 不支持的主任务类型:{main_goal}")
def assign_task(self, sub_tasks: List[Dict]) -> Dict:
"""
任务分配:根据子任务类型匹配执行者
:param sub_tasks: 子任务列表
:return: 分配结果(子任务ID→执行者ID)
"""
assignment = {}
# 建立“子任务类型→执行者列表”映射
specialty_map = defaultdict(list)
for executor_id, executor in self.executors.items():
specialty_map[executor.specialty.lower()].append(executor_id)
# 为每个子任务分配执行者
for sub_task in sub_tasks:
task_type = sub_task["type"]
# 匹配专长(如“data_processing”匹配“数据处理”)
matched_executors = []
if task_type == "data_processing":
matched_executors = specialty_map.get("数据处理", [])
elif task_type == "report_generation":
matched_executors = specialty_map.get("报告生成", [])
elif task_type == "visualization":
matched_executors = specialty_map.get("可视化", [])
if not matched_executors:
raise Exception(f"无匹配执行者的子任务:{sub_task['goal']}(类型:{task_type})")
# 简单负载均衡:选择第一个可用执行者(实际可扩展为轮询/权重)
executor_id = matched_executors[0]
assignment[sub_task["task_id"]] = executor_id
# 记录任务状态
self.task_status[sub_task["task_id"]] = {
"status": "assigned",
"executor_id": executor_id,
"sub_task": sub_task
}
print(f"[{self.name}] 分配子任务:{sub_task['task_id']} → 执行者:{self.executors[executor_id].name}")
return assignment
def execute_hierarchical(self, main_task: Dict) -> Dict:
"""
层级执行:拆解→分配→执行→聚合
:param main_task: 主任务
:return: 主任务最终结果
"""
start_time = time.time()
self.task_status[main_task["task_id"]] = {"status": "started"}
try:
# 1. 任务拆解
sub_tasks = self.decompose_task(main_task)
self.task_queue.extend(sub_tasks)
# 2. 任务分配
assignment = self.assign_task(sub_tasks)
# 3. 执行子任务(并行执行,此处简化为串行,实际可加线程池)
sub_task_results = {}
for sub_task_id, executor_id in assignment.items():
executor = self.executors[executor_id]
sub_task = self.task_status[sub_task_id]["sub_task"]
print(f"\n=== 执行子任务:{sub_task_id}(执行者:{executor.name})===")
# 执行者执行React闭环
action_steps = executor.think(sub_task)
action_result = executor.act(action_steps)
feedback = executor.feedback(action_result)
# 更新任务状态与结果
self.task_status[sub_task_id]["status"] = "completed" if action_result.success else "failed"
self.task_status[sub_task_id]["result"] = action_result.data
self.task_status[sub_task_id]["feedback"] = feedback
sub_task_results[sub_task_id] = action_result.data
# 4. 结果聚合:整合所有子任务结果
aggregated_result = {
"main_task_id": main_task["task_id"],
"status": "completed",
"sub_task_results": sub_task_results,
"total_time": round(time.time() - start_time, 2),
"executor_count": len(self.executors)
}
self.task_status[main_task["task_id"]] = {"status": "completed", "result": aggregated_result}
print(f"\n[{self.name}] 主任务执行完成!总耗时:{aggregated_result['total_time']}秒")
print(f"聚合结果:子任务数={len(sub_task_results)},状态=completed")
return aggregated_result
except Exception as e:
error_msg = str(e)
self.task_status[main_task["task_id"]] = {"status": "failed", "error": error_msg}
print(f"[{self.name}] 主任务执行失败:{error_msg}")
return {"status": "failed", "error": error_msg}
# 测试层级指挥架构
if __name__ == "__main__":
# 1. 初始化指挥者
commander = CommanderAgent(commander_id="commander_001", name="市场分析指挥者")
# 2. 初始化执行者(3个不同专长)
data_executor = ExecutorAgent(
agent_id="exec_001",
name="数据处理执行者",
specialty="数据处理"
)
report_executor = ExecutorAgent(
agent_id="exec_002",
name="报告生成执行者",
specialty="报告生成"
)
visual_executor = ExecutorAgent(
agent_id="exec_003",
name="可视化执行者",
specialty="可视化"
)
# 3. 注册执行者到指挥者
commander.register_executor(data_executor)
commander.register_executor(report_executor)
commander.register_executor(visual_executor)
# 4. 定义主任务:2025年Q1市场分析
main_task = {
"task_id": "main_task_001",
"goal": "完成2025年Q1市场分析,含数据计算、报告生成、趋势可视化",
"params": {
"data_path": "market_q1_2025.csv",
"report_title": "2025年Q1市场分析报告"
}
}
# 5. 执行层级指挥
commander.execute_hierarchical(main_task)
层级指挥架构的核心优势在于:
- 任务解耦:指挥者专注 “做什么、谁来做”,执行者专注 “怎么做”,避免单智能体职责过载;
- 弹性扩展:新增子任务类型时,仅需注册对应专长的执行者,无需修改指挥者核心逻辑;
- 可控性强:指挥者可实时监控子任务状态,出现失败时能快速定位到具体执行者与步骤。
实战优化策略:性能与可靠性提升
在实际部署中,多智能体 React 模式可能面临 “执行延迟”“任务失败重试”“资源竞争” 等问题。针对这些痛点,我们从 “并行执行”“错误重试”“缓存优化” 三个维度进行实战优化,确保系统高效稳定运行。
1. 并行执行:提升多子任务效率
传统串行执行子任务效率低(如 3 个子任务各耗时 2 秒,串行需 6 秒,并行仅需 2 秒)。通过 Pythonconcurrent.futures.ThreadPoolExecutor实现执行者并行执行,指挥者统一监控进度。
# 指挥者并行执行优化(修改execute_hierarchical方法)
from concurrent.futures import ThreadPoolExecutor, as_completed
def execute_hierarchical_parallel(self, main_task: Dict, max_workers: int = 3) -> Dict:
"""并行执行子任务,提升效率"""
start_time = time.time()
self.task_status[main_task["task_id"]] = {"status": "started"}
try:
# 1. 任务拆解与分配(同前)
sub_tasks = self.decompose_task(main_task)
assignment = self.assign_task(sub_tasks)
if not assignment:
raise Exception("无分配的子任务")
# 2. 并行执行子任务
sub_task_results = {}
# 创建线程池(max_workers:最大并行数)
with ThreadPoolExecutor(max_workers=max_workers) as executor:
# 提交所有子任务到线程池
future_tasks = {}
for sub_task_id, executor_id in assignment.items():
sub_task = self.task_status[sub_task_id]["sub_task"]
# 提交任务:每个子任务对应一个线程
future = executor.submit(
self._run_sub_task, # 子任务执行函数
sub_task_id=sub_task_id,
executor_id=executor_id,
sub_task=sub_task
)
future_tasks[future] = sub_task_id
# 监控所有子任务完成情况
for future in as_completed(future_tasks):
sub_task_id = future_tasks[future]
try:
# 获取子任务结果
result, feedback = future.result()
sub_task_results[sub_task_id] = result
self.task_status[sub_task_id].update({
"status": "completed",
"result": result,
"feedback": feedback
})
print(f"[{self.name}] 子任务并行完成:{sub_task_id}")
except Exception as e:
error_msg = str(e)
self.task_status[sub_task_id].update({
"status": "failed",
"error": error_msg
})
print(f"[{self.name}] 子任务并行失败:{sub_task_id},原因:{error_msg}")
# 3. 结果聚合(同前)
aggregated_result = {
"main_task_id": main_task["task_id"],
"status": "completed" if all(s["status"] == "completed" for s in self.task_status.values()) else "partially_completed",
"sub_task_results": sub_task_results,
"total_time": round(time.time() - start_time, 2),
"parallel_workers": max_workers
}
print(f"\n[{self.name}] 并行执行完成!总耗时:{aggregated_result['total_time']}秒(串行预计:{len(sub_tasks)*2}秒)")
return aggregated_result
except Exception as e:
error_msg = str(e)
self.task_status[main_task["task_id"]] = {"status": "failed", "error": error_msg}
return {"status": "failed", "error": error_msg}
def _run_sub_task(self, sub_task_id: str, executor_id: str, sub_task: Dict) -> tuple:
"""子任务执行函数(供线程池调用)"""
executor = self.executors[executor_id]
try:
action_steps = executor.think(sub_task)
action_result = executor.act(action_steps)
feedback = executor.feedback(action_result)
if not action_result.success:
raise Exception(f"子任务执行失败:{feedback['adjustment']}")
return action_result.data, feedback
except Exception as e:
raise Exception(f"执行者{executor.name}执行失败:{str(e)}")
2. 错误重试:提升任务可靠性
执行者执行子任务时可能因 “工具调用超时”“数据格式错误” 等临时问题失败,通过 “重试机制” 自动重新执行,避免人工介入。
# 执行者错误重试优化(修改act方法)
def act_with_retry(self, action_steps: List[ActionType], max_retries: int = 2) -> ActionResult:
"""带重试机制的行动方法"""
retry_count = 0
while retry_count < max_retries:
try:
# 调用原act方法执行
return super().act(action_steps)
except Exception as e:
retry_count += 1
if retry_count >= max_retries:
# 达到最大重试次数,返回失败
error_msg = f"执行失败(已重试{max_retries}次):{str(e)}"
return ActionResult(success=False, error=error_msg)
# 重试延迟(指数退避:1s→2s→4s...)
delay = 2 ** (retry_count - 1)
print(f"[{self.name}] 执行失败,{delay}秒后重试(第{retry_count}次):{str(e)}")
time.sleep(delay)
return ActionResult(success=False, error="未执行任何重试")
# 指挥者中调用重试方法(修改_run_sub_task)
def _run_sub_task_with_retry(self, sub_task_id: str, executor_id: str, sub_task: Dict) -> tuple:
executor = self.executors[executor_id]
try:
action_steps = executor.think(sub_task)
# 调用带重试的act方法
action_result = executor.act_with_retry(action_steps, max_retries=2)
feedback = executor.feedback(action_result)
if not action_result.success:
raise Exception(f"子任务执行失败:{feedback['adjustment']}")
return action_result.data, feedback
except Exception as e:
raise Exception(f"执行者{executor.name}执行失败:{str(e)}")
3. 缓存优化:避免重复计算
多智能体执行相似任务时(如多个子任务加载同一数据文件),会重复调用相同工具,浪费资源。通过 “结果缓存” 存储工具执行结果,后续调用直接复用。
# 工具结果缓存优化(基于functools.lru_cache)
from functools import lru_cache
import json
# 为执行者工具添加缓存装饰器(修改ExecutorAgent)
class ExecutorAgentWithCache(ExecutorAgent):
def __init__(self, agent_id: str, name: str, specialty: str):
super().__init__(agent_id, name, specialty)
# 为工具函数添加缓存(键:参数的JSON字符串,避免不可哈希类型)
self._load_data = self._cache_tool(self._load_data)
def _cache_tool(self, tool_func):
"""工具函数缓存装饰器"""
@lru_cache(maxsize=128) # 最大缓存128个结果
def cached_tool(task_params_json: str, intermediate_json: str) -> ActionResult:
# 将JSON字符串转为字典
task_params = json.loads(task_params_json)
intermediate = json.loads(intermediate_json) if intermediate_json else {}
# 调用原工具函数
result = tool_func(task_params, intermediate)
# 缓存结果(需确保ActionResult可哈希,此处简化为返回数据)
return result
# 包装函数:处理参数序列化
def wrapper(task_params: Dict, intermediate: Dict) -> ActionResult:
task_params_json = json.dumps(task_params, sort_keys=True)
intermediate_json = json.dumps(intermediate, sort_keys=True)
return cached_tool(task_params_json, intermediate_json)
return wrapper
缓存优化的关键在于:选择 “输入稳定、计算耗时” 的工具(如数据加载、复杂分析)进行缓存,避免缓存频繁变化的结果(如实时数据查询)。
总结:多智能体 React 模式的落地要点
基于 React 模式的多智能体层级指挥架构,核心是 “闭环思考 + 分层协同 + 实战优化”。对程序员而言,落地时需关注三个要点:
- 职责边界清晰:指挥者专注任务拆解与分配,执行者专注子任务 React 闭环,避免 “一锅粥” 式设计;
- 优化贴合场景:并行执行适合 CPU 密集型子任务,重试机制适合网络 / IO 不稳定场景,缓存适合重复计算场景;
- 可观测性强:通过日志、任务状态监控,确保问题可定位(如哪个执行者、哪个步骤失败)。
这种架构不仅适用于市场分析、数据处理等场景,还可扩展到智能客服(多轮对话 + 工具调用)、自动驾驶(多传感器数据协同)等领域。随着多智能体技术的发展,层级指挥 + React 模式将成为复杂任务处理的核心范式 —— 掌握其实现与优化技巧,能让你在智能系统开发中更具竞争力。
更多推荐


所有评论(0)