多 Agent 协作框架:任务分解、通信协议与冲突解决的架构设计
多 Agent 协作框架:任务分解、通信协议与冲突解决的架构设计

一、单 Agent 的能力天花板:复杂任务的分解困境
单个 AI Agent 在处理简单、明确的任务时表现良好,但面对多步骤、多领域、需要并行处理的复杂任务时,会暴露三个结构性缺陷:上下文窗口有限导致长链推理断裂,单线程执行无法利用并行性,单一角色难以覆盖跨领域知识。
多 Agent 协作的核心思路是将复杂任务分解为子任务,分配给具有不同专长的 Agent 并行执行,再通过协调机制汇总结果。但这引入了新的工程难题:任务如何分解与分配?Agent 之间如何通信?并行执行的结果如何合并?当多个 Agent 产生冲突输出时如何裁决?
二、多 Agent 协作的核心架构:编排器-工作器模式
多 Agent 系统的架构模式可归纳为三种:顺序链式(Pipeline)、层级式(Hierarchical)与对等式(Peer-to-Peer)。本文聚焦于工程实践中最常用的层级式架构——编排器-工作器(Orchestrator-Worker)模式。
graph TB
subgraph 编排器-工作器架构
O[Orchestrator<br/>任务分解与结果汇总] --> W1[Worker A<br/>代码生成]
O --> W2[Worker B<br/>测试验证]
O --> W3[Worker C<br/>文档编写]
W1 -->|代码片段| O
W2 -->|测试报告| O
W3 -->|文档草稿| O
end
subgraph 通信层
MB[消息总线<br/>基于 Channel 的异步通信]
end
O <-.->|发布/订阅| MB
W1 <-.->|发布/订阅| MB
W2 <-.->|发布/订阅| MB
W3 <-.->|发布/订阅| MB
subgraph 冲突裁决
CR[冲突检测器] -->|策略选择| CR1[优先级裁决]
CR -->|策略选择| CR2[投票裁决]
CR -->|策略选择| CR3[人工介入]
end
O -->|冲突上报| CR
style O fill:#e1f5fe
style MB fill:#fff3e0
style CR fill:#ffebee
编排器(Orchestrator) 负责三个核心职责:将用户请求分解为子任务图(DAG),将子任务分配给合适的工作器,汇总工作器的输出并解决冲突。
工作器(Worker) 是具有特定专长的 Agent 实例,每个工作器拥有独立的系统提示、工具集与上下文。工作器之间不直接通信,所有交互通过编排器中转。
消息总线 基于 Rust 的 tokio::sync::mpsc 实现,提供异步、类型安全的进程内通信。
三、Rust 实现多 Agent 协作框架
3.1 任务图与依赖管理
use std::collections::{HashMap, HashSet};
use serde::{Deserialize, Serialize};
/// 子任务定义
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Task {
pub id: String,
pub description: String,
pub required_role: AgentRole,
pub dependencies: HashSet<String>, // 依赖的前置任务 ID
pub priority: u32,
pub status: TaskStatus,
pub result: Option<TaskResult>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Hash)]
pub enum AgentRole {
CodeGenerator,
TestRunner,
DocWriter,
Reviewer,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub enum TaskStatus {
Pending,
Running,
Completed,
Failed(String),
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TaskResult {
pub task_id: String,
pub output: String,
pub confidence: f64, // 置信度 [0.0, 1.0]
pub artifacts: Vec<String>, // 产出物标识
}
/// 任务图:DAG 结构,管理任务间的依赖关系
pub struct TaskGraph {
tasks: HashMap<String, Task>,
}
impl TaskGraph {
pub fn new() -> Self {
Self { tasks: HashMap::new() }
}
pub fn add_task(&mut self, task: Task) {
self.tasks.insert(task.id.clone(), task);
}
/// 获取当前可执行的任务(所有依赖已完成)
pub fn ready_tasks(&self) -> Vec<&Task> {
self.tasks.values()
.filter(|t| t.status == TaskStatus::Pending)
.filter(|t| {
t.dependencies.iter().all(|dep_id| {
self.tasks.get(dep_id)
.map(|dep| dep.status == TaskStatus::Completed)
.unwrap_or(false)
})
})
.collect()
}
/// 检测循环依赖
pub fn detect_cycles(&self) -> bool {
let mut visited = HashSet::new();
let mut rec_stack = HashSet::new();
for task_id in self.tasks.keys() {
if self.dfs_cycle(task_id, &mut visited, &mut rec_stack) {
return true;
}
}
false
}
fn dfs_cycle(
&self,
node: &str,
visited: &mut HashSet<String>,
rec_stack: &mut HashSet<String>,
) -> bool {
if rec_stack.contains(node) {
return true;
}
if visited.contains(node) {
return false;
}
visited.insert(node.to_string());
rec_stack.insert(node.to_string());
if let Some(task) = self.tasks.get(node) {
for dep in &task.dependencies {
if self.dfs_cycle(dep, visited, rec_stack) {
return true;
}
}
}
rec_stack.remove(node);
false
}
}
3.2 异步消息通信
use tokio::sync::{mpsc, oneshot};
/// Agent 间通信消息
#[derive(Debug)]
pub enum AgentMessage {
/// 编排器 → 工作器:分配任务
AssignTask {
task: Task,
context: String, // 前置任务的输出摘要
reply_to: oneshot::Sender<TaskResult>,
},
/// 工作器 → 编排器:请求额外信息
RequestContext {
task_id: String,
query: String,
reply_to: oneshot::Sender<String>,
},
/// 工作器 → 编排器:报告进度
Progress {
task_id: String,
percentage: u8,
note: String,
},
}
/// 工作器实例
pub struct Worker {
role: AgentRole,
receiver: mpsc::Receiver<AgentMessage>,
}
impl Worker {
pub fn new(role: AgentRole, receiver: mpsc::Receiver<AgentMessage>) -> Self {
Self { role, receiver }
}
pub async fn run(&mut self) {
while let Some(msg) = self.receiver.recv().await {
match msg {
AgentMessage::AssignTask { task, context, reply_to } => {
let result = self.execute_task(&task, &context).await;
// 回复结果给编排器,忽略发送失败(编排器可能已超时)
let _ = reply_to.send(result);
}
AgentMessage::RequestContext { query, reply_to, .. } => {
// 工作器主动请求额外上下文
let response = format!("上下文响应:{}", query);
let _ = reply_to.send(response);
}
AgentMessage::Progress { task_id, percentage, note } => {
log::info!(
"[{}] 任务 {} 进度: {}% - {}",
self.role.variant_name(),
task_id, percentage, note
);
}
}
}
}
async fn execute_task(&self, task: &Task, context: &str) -> TaskResult {
// 实际生产中此处调用 LLM API
// 此处用模拟逻辑示意
TaskResult {
task_id: task.id.clone(),
output: format!(
"[{}] 完成任务:{},上下文:{}",
self.role.variant_name(),
task.description,
&context[..context.len().min(100)]
),
confidence: 0.85,
artifacts: vec![format!("artifact_{}", task.id)],
}
}
}
3.3 编排器与冲突裁决
/// 编排器:任务分解、分配与冲突裁决
pub struct Orchestrator {
task_graph: TaskGraph,
workers: HashMap<AgentRole, mpsc::Sender<AgentMessage>>,
conflict_strategy: ConflictStrategy,
}
#[derive(Debug, Clone)]
pub enum ConflictStrategy {
/// 高置信度优先
HighestConfidence,
/// 指定角色优先
RolePriority(Vec<AgentRole>),
/// 人工介入
HumanReview,
}
impl Orchestrator {
pub fn new(conflict_strategy: ConflictStrategy) -> Self {
Self {
task_graph: TaskGraph::new(),
workers: HashMap::new(),
conflict_strategy,
}
}
/// 注册工作器
pub fn register_worker(&mut self, role: AgentRole, sender: mpsc::Sender<AgentMessage>) {
self.workers.insert(role, sender);
}
/// 执行完整任务流程
pub async fn execute(&mut self, user_request: &str) -> Result<String, String> {
// Step 1: 分解任务(实际生产中调用 LLM)
self.decompose_task(user_request)?;
if self.task_graph.detect_cycles() {
return Err("任务图存在循环依赖".to_string());
}
// Step 2: 按依赖顺序执行
let mut completed_results: HashMap<String, TaskResult> = HashMap::new();
loop {
let ready = self.task_graph.ready_tasks();
if ready.is_empty() {
break;
}
// 并行分配就绪任务
let mut handles = Vec::new();
for task in ready {
let task_id = task.id.clone();
let role = task.required_role.clone();
let sender = self.workers.get(&role)
.ok_or_else(|| format!("无可用工作器:{:?}", role))?;
// 构造上下文:汇总前置任务的输出
let context = self.build_context(&task.dependencies, &completed_results);
let (reply_tx, reply_rx) = oneshot::channel();
sender.send(AgentMessage::AssignTask {
task: task.clone(),
context,
reply_to: reply_tx,
}).await.map_err(|e| format!("发送任务失败:{}", e))?;
handles.push(tokio::spawn(async move {
let task_id = task_id;
match reply_rx.await {
Ok(result) => (task_id, Some(result)),
Err(_) => (task_id, None),
}
}));
}
// 收集结果
for handle in handles {
let (task_id, result) = handle.await.map_err(|e| format!("任务执行异常:{}", e))?;
match result {
Some(r) => {
completed_results.insert(task_id.clone(), r);
if let Some(task) = self.task_graph.tasks.get_mut(&task_id) {
task.status = TaskStatus::Completed;
task.result = Some(completed_results[&task_id].clone());
}
}
None => {
if let Some(task) = self.task_graph.tasks.get_mut(&task_id) {
task.status = TaskStatus::Failed("工作器无响应".to_string());
}
}
}
}
}
// Step 3: 汇总结果
Ok(self.merge_results(&completed_results))
}
fn build_context(
&self,
deps: &HashSet<String>,
results: &HashMap<String, TaskResult>,
) -> String {
deps.iter()
.filter_map(|id| results.get(id))
.map(|r| r.output.clone())
.collect::<Vec<_>>()
.join("\n")
}
fn merge_results(&self, results: &HashMap<String, TaskResult>) -> String {
results.values()
.map(|r| r.output.clone())
.collect::<Vec<_>>()
.join("\n\n---\n\n")
}
fn decompose_task(&mut self, request: &str) -> Result<(), String> {
// 简化示意:将请求分解为代码生成 + 测试 + 文档三个子任务
let code_task = Task {
id: "code".to_string(),
description: format!("生成代码:{}", request),
required_role: AgentRole::CodeGenerator,
dependencies: HashSet::new(),
priority: 1,
status: TaskStatus::Pending,
result: None,
};
let test_task = Task {
id: "test".to_string(),
description: format!("编写测试:{}", request),
required_role: AgentRole::TestRunner,
dependencies: HashSet::from(["code".to_string()]),
priority: 2,
status: TaskStatus::Pending,
result: None,
};
let doc_task = Task {
id: "doc".to_string(),
description: format!("编写文档:{}", request),
required_role: AgentRole::DocWriter,
dependencies: HashSet::from(["code".to_string()]),
priority: 2,
status: TaskStatus::Pending,
result: None,
};
self.task_graph.add_task(code_task);
self.task_graph.add_task(test_task);
self.task_graph.add_task(doc_task);
Ok(())
}
}
四、多 Agent 协作的架构权衡
4.1 通信开销与一致性延迟
编排器-工作器模式下,所有通信必须经过编排器中转,这引入了不可避免的延迟。对于需要频繁交换中间结果的任务(如代码生成与审查的迭代循环),通信开销可能占总执行时间的 30% 以上。对等式通信(P2P)可以减少中转延迟,但增加了协调复杂度与一致性风险。
4.2 任务分解的粒度权衡
任务粒度过细会导致编排开销膨胀(每个子任务的分配、通信、汇总都有固定成本),粒度过粗则无法充分利用并行性。实测发现,子任务数量在 3-7 个时,总执行时间最优。超过 10 个子任务后,编排开销开始显著影响性能。
4.3 冲突裁决的可靠性
当多个工作器对同一问题给出不同答案时,冲突裁决策略的可靠性直接决定系统输出质量。高置信度优先策略假设 LLM 的置信度校准是准确的,但研究表明 LLM 的置信度校准存在系统性偏差。角色优先策略假设特定角色更可靠,但忽略了具体问题的差异。人工介入策略最可靠,但破坏了自动化流程。
4.4 错误传播与容错
DAG 中的某个节点失败后,其所有下游任务都无法执行。简单的"失败即中止"策略过于保守,而"跳过失败任务继续执行"可能导致下游任务基于不完整的输入产生错误输出。生产环境需要为每个任务定义降级策略(如使用缓存结果、使用简化版本)。
五、总结
多 Agent 协作框架通过编排器-工作器模式,将复杂任务分解为子任务 DAG 并行执行,解决了单 Agent 的上下文窗口限制与单线程执行瓶颈。核心设计包括:基于 DAG 的任务图管理依赖关系,基于 Channel 的异步消息通信,以及基于策略的冲突裁决机制。
落地路线建议:第一,从顺序链式协作开始(代码生成 → 审查),验证通信协议的可靠性;第二,引入 DAG 任务图支持并行执行,重点优化编排开销;第三,实现冲突裁决机制,从高置信度优先逐步演进到混合策略;第四,为关键任务建立降级与容错机制,提升系统的鲁棒性。
更多推荐


所有评论(0)