在多智能体系统开发中,传统 “单智能体独立决策” 模式常面临任务协同混乱、资源浪费、响应延迟等问题。而基于 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)

层级指挥架构的核心优势在于:

  1. 任务解耦:指挥者专注 “做什么、谁来做”,执行者专注 “怎么做”,避免单智能体职责过载;
  1. 弹性扩展:新增子任务类型时,仅需注册对应专长的执行者,无需修改指挥者核心逻辑;
  1. 可控性强:指挥者可实时监控子任务状态,出现失败时能快速定位到具体执行者与步骤。

实战优化策略:性能与可靠性提升

在实际部署中,多智能体 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 模式的多智能体层级指挥架构,核心是 “闭环思考 + 分层协同 + 实战优化”。对程序员而言,落地时需关注三个要点:

  1. 职责边界清晰:指挥者专注任务拆解与分配,执行者专注子任务 React 闭环,避免 “一锅粥” 式设计;
  1. 优化贴合场景:并行执行适合 CPU 密集型子任务,重试机制适合网络 / IO 不稳定场景,缓存适合重复计算场景;
  1. 可观测性强:通过日志、任务状态监控,确保问题可定位(如哪个执行者、哪个步骤失败)。

这种架构不仅适用于市场分析、数据处理等场景,还可扩展到智能客服(多轮对话 + 工具调用)、自动驾驶(多传感器数据协同)等领域。随着多智能体技术的发展,层级指挥 + React 模式将成为复杂任务处理的核心范式 —— 掌握其实现与优化技巧,能让你在智能系统开发中更具竞争力。

更多推荐