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

cover

一、单 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 任务图支持并行执行,重点优化编排开销;第三,实现冲突裁决机制,从高置信度优先逐步演进到混合策略;第四,为关键任务建立降级与容错机制,提升系统的鲁棒性。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐