MOSAIC框架:如何通过自适应聚合与推理并发优化多智能体协作效率
发布时间:2026/8/20 8:40:41 作者:尧图编辑部 阅读量:1,286

1. 项目概述从“智能体排队”到“交响乐团指挥”最近在折腾大语言模型应用落地的朋友估计都绕不开一个头疼的问题怎么让多个智能体Agent高效地协同工作我们常常会设计一个“主控”Agent让它去调度几个具备不同专长的“子”Agent比如一个负责查资料一个负责写代码一个负责做总结。想法很美好但一跑起来就发现这流程慢得让人心焦。Agent A 干完活把结果传给主控主控再思考一下分派给 Agent BB 干完再传回来……整个链条是串行的大量的时间花在了等待网络I/O、模型推理和上下文传递上。这感觉就像你去银行办业务明明开了五个窗口但所有人还是得在同一个窗口前排队效率低下。“MOSAIC: Efficient Mixture-of-Agent Scheduling via Adaptive Aggregation and Inference Concurrency”这个项目瞄准的就是这个痛点。它不是一个新模型架构而是一套智能体工作流的调度与执行优化框架。你可以把它想象成一个交响乐团的指挥但这位指挥的厉害之处在于他不仅能看懂总谱任务还能实时预判每位乐手智能体的状态和乐句难度动态调整演奏顺序甚至让不同的声部提前准备、同时演奏最终让整首曲子复杂任务以最短的时间、最流畅的方式完成。它的核心价值在于在不改变单个智能体能力的前提下通过系统级的调度优化显著提升多智能体系统的整体吞吐量和响应速度。这对于构建复杂的AI应用如自动化报告生成、多步骤问题求解、交互式创作平台等意味着更低的延迟和更高的成本效益。无论是研究者在实验多智能体协作的新范式还是工程师在部署对实时性有要求的AI服务MOSAIC提供的思路和工具都极具参考价值。2. 核心设计思路打破串行枷锁的两把钥匙传统的多智能体流水线是“请求-响应”式的串行思维MOSAIC的设计哲学则是“并发”与“预测”。它主要从两个维度进行优化这也是其名称中“Adaptive Aggregation”自适应聚合和“Inference Concurrency”推理并发的由来。2.1 自适应聚合让数据流动更聪明串行链路最大的浪费在于“空等”。Agent B 必须等到 Agent A 的完整输出后才能开始自己的工作。但很多时候Agent A 的输出是流式的token by token或者其输出的前半部分已经包含了足够 Agent B 启动工作的信息。自适应聚合的核心思想是不必等待上游智能体生成全部结果一旦其输出中包含了满足下游智能体启动条件的信息就立即触发下游智能体的工作。这需要对智能体间的数据依赖关系有精细的建模。举个例子我们有一个任务“分析某公司财报总结其风险并给出投资建议”。我们设计了三个智能体信息提取Agent从财报PDF中提取关键财务数据。风险分析Agent基于财务数据计算比率、识别异常。报告生成Agent综合前两步结果撰写分析报告。在传统模式下报告生成Agent必须等风险分析Agent完全结束后才能开始。但在MOSAIC框架下我们可以这样设计信息提取Agent一旦提取出“营业收入”和“净利润”这两个数据就可以立即被发送给风险分析Agent开始计算“净利率”。与此同时信息提取Agent继续提取“资产负债”数据。风险分析Agent在计算出“净利率”后如果发现异常如大幅下滑这个“初步风险信号”可以立即触发报告生成Agent开始撰写报告的风险提示部分。而报告生成Agent的“公司概述”部分可能只需要公司名称和年份这些信息在信息提取Agent工作的最初期就已经可获得。这样三个Agent的工作在时间线上出现了大量的重叠而不是一个接一个的接力赛。注意实现自适应聚合的关键是定义清晰的“数据就绪”事件和“触发条件”。这需要开发者对任务流程和每个Agent的输入需求有深刻理解。过度激进的聚合可能导致下游Agent基于不完整或噪声信息做出错误判断需要在“速度”和“准确性”之间做权衡。2.2 推理并发让计算资源忙起来即使通过自适应聚合让数据流提前了如果计算资源GPU/CPU在同一时间只能服务一个智能体的推理请求那么硬件利用率依然低下速度瓶颈只是从“等待数据”转移到了“等待算力”。推理并发的目标是让多个智能体的模型推理过程在硬件层面尽可能并行执行。这听起来简单但在大语言模型场景下挑战不小因为单个模型的推理过程通常占用大量显存且对输入序列长度敏感。MOSAIC在这方面可能借鉴或实现了多种技术连续批处理这是推理服务器如vLLM, TGI的常见优化。当多个推理请求先后到达时服务器不是立即单独处理每一个而是稍作等待将多个请求的输入张量在批次维度上进行拼接一次性送入模型计算。这能极大提升GPU的计算单元利用率。MOSAIC的调度器需要智能地将不同智能体的推理请求组合成高效的批次。模型并行与流水线并行对于超大规模的模型单个GPU可能放不下。MOSAIC可能需要协调将单个智能体的模型拆分到多个设备上模型并行或者将不同智能体的模型部署到不同的设备上形成流水线。抢占式调度与资源预留对于优先级不同的智能体任务调度器需要能够动态分配计算资源。例如负责最终输出的“主控”Agent可能需要更高的优先级以确保响应延迟可控。将自适应聚合与推理并发结合才是MOSAIC威力最大的地方。调度器不仅要知道“哪个Agent的数据准备好了”还要知道“当前GPU上正在运行什么队列里还有什么下一个最应该调度哪个Agent进来以形成高效的连续批处理”。这需要一个全局的、感知系统负载的调度策略。3. 系统架构与关键组件拆解一个典型的MOSAIC风格系统可能包含以下几个核心组件我们可以将其类比为一个现代化的物流调度中心。3.1 任务解析与DAG构建器这是系统的“大脑”前额叶负责理解用户提交的复杂任务。输入是一个自然语言指令或结构化任务描述输出是一个有向无环图。节点代表一个原子性的智能体操作。每个节点包含智能体类型、所需输入模式、预期输出模式、执行参数如模型温度、最大生成长度。边代表数据依赖关系。从节点A指向节点B的边意味着B的某个输入依赖于A的某个输出。动态性这个DAG可能不是完全静态的。例如一个“决策Agent”的输出“需要深入分析XX部分”可能会动态创建新的节点并添加到图中。实操要点构建DAG时要尽量将任务拆解为功能单一、输入输出明确的原子节点。过于复杂的节点会成为并发瓶颈。同时要为每条依赖边标注详细的数据契约比如“需要‘净利润’数值字段”或“需要‘风险段落’文本片段”这是实现自适应聚合的基础。3.2 自适应聚合调度器这是系统的“中枢神经系统”是MOSAIC的核心创新所在。它持续监控DAG中所有节点的状态等待、就绪、运行、完成和数据的生产/消费情况。其工作流程如下事件监听监听各个智能体节点的输出流。每个节点不再是生成完整结果后一次性抛出而是需要支持“流式”或“分阶段”的事件发布。例如发布{“event”: “partial_output”, “data_type”: “financial_figure”, “name”: “revenue”, “value”: “1.2B”}。依赖关系检查调度器内部维护一个“数据依赖表”。当收到一个数据事件时它立刻查询有哪些下游节点正在等待这类数据这些节点的其他输入依赖是否已经满足触发决策如果一个下游节点的所有输入依赖因最新事件而得到满足则该节点状态立即从“等待”变为“就绪”并被加入就绪队列。如果只是部分满足则更新该节点的“已满足依赖列表”继续等待。优先级排序就绪队列中的节点并非平等。调度器会根据一些启发式规则排序例如关键路径节点优先在DAG中处于最长路径决定最终完成时间上的节点应优先执行。产出数据多的节点优先一个节点的输出是多个下游节点的输入提前执行它能解锁更多并行性。计算量小的节点优先快速执行完小任务可以更快地释放资源。3.3 并发推理执行引擎这是系统的“肌肉”负责以最高效的方式在硬件上执行“就绪”的智能体节点。它与调度器紧密耦合。请求池化执行引擎从调度器接收一批“就绪”的节点任务。它会将这些任务转化为对底层模型服务如 OpenAI API 本地部署的 vLLM 服务的推理请求。连续批处理执行引擎不会立即发送单个请求而是会实施一个短暂的等待窗口例如几毫秒到几十毫秒将窗口期内到达的、且模型相同/兼容的请求聚合在一起形成一个批次Batch再发送给推理后端。这能大幅提升GPU利用率。异构资源管理系统中可能存在不同规模的模型如 GPT-4-large 和 GPT-3.5-turbo部署在不同的硬件或服务上。执行引擎需要根据节点的模型要求将请求路由到合适的后端并管理不同后端的负载。结果回馈与流式传播当推理结果开始返回时同样可能是流式的执行引擎将结果拆分成对应节点的输出并立即以“数据事件”的形式反馈给调度器从而可能触发新一轮的节点调度形成正向循环。3.4 上下文管理与数据总线这是系统的“血液循环系统”。在多轮并发执行中数据中间结果的传递和管理至关重要。统一数据格式所有智能体的输入和输出都被强制或封装为一种统一的中间表示格式例如JSON并遵循预定义的模式Schema。这简化了数据路由和解析。上下文存储需要一个共享的、低延迟的存储如内存缓存Redis或高性能内存结构来存储所有中间结果。每个数据片段都有一个全局唯一的ID供下游节点引用。版本与一致性在动态DAG中一个节点的输出可能被多个下游节点消费。如果该节点因为输入变化而重新执行其输出版本会更新。数据总线需要管理这种版本依赖确保下游节点使用的是正确版本的数据避免脏读。4. 实战构建一个简易的MOSAIC风格任务流水线理论说了这么多我们动手设计一个简化版的系统来体会一下MOSAIC的思想。假设我们要实现一个“智能周报生成器”任务描述是“基于本周的JIRA tickets和Git commits生成一份技术团队周报需包含完成的工作、遇到的问题和下周计划。”4.1 第一步任务拆解与DAG设计我们设计四个智能体节点JIRA解析Agent (A1)输入JIRA API令牌/查询。输出本周已关闭的ticket列表每个包含标题、描述、负责人。Git解析Agent (A2)输入Git仓库地址、分支、时间范围。输出本周的commit列表每个包含哈希、作者、消息、更改文件。信息融合与摘要Agent (A3)输入A1的输出ticket列表、A2的输出commit列表。输出一份结构化的中间摘要按人/模块分类的工作项。报告生成Agent (A4)输入A3的输出结构化摘要。输出最终的人类可读的周报Markdown文本。依赖关系非常清晰A3 依赖 A1 和 A2 A4 依赖 A3。这是一个简单的“V”型DAG。4.2 第二步实现自适应聚合的契机在基础设计中A3必须等待A1和A2都完全执行完毕。但我们能否优化优化点1A1和A2之间没有依赖它们可以完全并行启动。这是最基础的并发。优化点2A1和A2的输出都是“列表”。假设A1处理速度较快先输出了前5个ticket。A3是否可以先基于这5个ticket和尚未到来的commit开始做一些预处理工作比如先提取这5个ticket的关键词或进行初步分类。这要求A3被设计成可以处理“部分输入”。优化点3A4是否必须等待A3生成完整的结构化摘要也许A3可以按“人员”维度流式输出。当A3输出“张三本周完成了【任务A 任务B】”时A4就可以开始撰写关于张三的周报部分。为了实现优化点2和3我们需要重新定义节点接口和数据事件。数据事件定义示例// A1 产生的事件 { “agent_id”: “A1”, “event_type”: “data_chunk”, “chunk_id”: “tickets_batch_1”, “data”: [ {“id”: “PROJ-101”, “title”: “修复登录BUG”, “assignee”: “张三”, “status”: “closed”}, {“id”: “PROJ-102”, “title”: “设计文档评审”, “assignee”: “李四”, “status”: “closed”} ], “is_last_chunk”: false }A3的依赖声明需要细化依赖A1的输出但消费模式可以是“增量消费”。依赖A2的输出消费模式也是“增量消费”。A3内部实现一个缓冲区分别缓存来自A1和A2的增量数据。一旦两个缓冲区都有数据就可以启动一轮“部分融合”处理产生一个“部分摘要”事件发给A4。4.3 第三步调度与执行的核心逻辑伪代码示意下面用伪代码展示调度器核心循环的概念class AdaptiveScheduler: def __init__(self, dag): self.dag dag self.node_status {node: ‘pending’ for node in dag.nodes} self.data_buffer {} # 存储各节点产生的数据块 self.ready_queue [] def on_data_event(self, event): # 1. 存储数据 self.data_buffer[event.agent_id] self.data_buffer.get(event.agent_id, []) [event.data] # 2. 检查依赖此数据源的下游节点 for downstream_node in self.dag.get_downstream_nodes(event.agent_id): if self.node_status[downstream_node] ! ‘pending’: continue # 3. 检查该下游节点的所有输入依赖是否都已部分满足 all_inputs_ready True for upstream_node in self.dag.get_upstream_nodes(downstream_node): if upstream_node not in self.data_buffer or len(self.data_buffer[upstream_node]) 0: all_inputs_ready False break # 4. 如果满足则激活该节点 if all_inputs_ready: self.node_status[downstream_node] ‘ready’ # 准备该节点的输入数据可能是多个上游节点的最新数据块 input_data self.prepare_input_for_node(downstream_node) # 加入就绪队列并附带输入数据和优先级分数 self.ready_queue.append((downstream_node, input_data, self.calculate_priority(downstream_node))) def scheduling_loop(self): while not self.is_dag_finished(): if self.ready_queue: # 按优先级排序例如关键路径优先 self.ready_queue.sort(keylambda x: x[2], reverseTrue) # 取出优先级最高的一个或多个任务 tasks_to_run self.select_tasks_for_batch(self.ready_queue) # 交给并发执行引擎 execution_engine.submit_batch(tasks_to_run) # 从就绪队列移除 for task in tasks_to_run: self.ready_queue.remove(task) self.node_status[task.node] ‘running’ else: # 等待新数据事件 time.sleep(0.001) def prepare_input_for_node(self, node): # 这里实现从data_buffer中提取该节点所需的数据 # 可能涉及数据的拼接、转换等 inputs {} for upstream in self.dag.get_upstream_nodes(node): # 取该上游节点最新的数据块或所有累积的数据块 upstream_data self.data_buffer[upstream][-1] # 示例取最新块 inputs[upstream] upstream_data return inputs4.4 第四步集成并发推理引擎执行引擎execution_engine接收到一批tasks_to_run后按模型分组将任务按它们所需的底层模型如 ‘gpt-3.5-turbo’ ‘claude-3-sonnet’分组。构建批处理请求对同一模型的每组任务将它们各自的输入prompt context组合成一个批处理列表。调用模型API/服务使用该模型对应的客户端以批处理模式发送请求。对于不支持原生批处理的API则需要在客户端模拟如使用异步请求。处理流式响应如果模型支持流式输出执行引擎需要将token流实时分拆并立即通过on_data_event回调函数发回给调度器从而实现更极致的“流水线”效果。如果不支持流式则等待完整响应后再回传。5. 性能优化与常见问题排查在实际部署中你会遇到各种预期之外的问题。以下是一些关键的性能调优点和排错经验。5.1 性能瓶颈分析与优化瓶颈位置表现症状优化策略调度器逻辑CPU占用高但GPU利用率低。就绪队列经常为空或很短。检查DAG复杂度调度算法如优先级计算是否过于复杂。考虑使用更高效的数据结构如优先堆管理就绪队列。数据序列化/反序列化网络I/O或进程间通信延迟高中间数据量大时尤为明显。使用高效的二进制序列化协议如Protocol Buffers, MessagePack替代JSON。考虑使用共享内存传递大块数据。推理后端GPU利用率饱和请求排队严重单个请求延迟飙升。1.连续批处理调整执行引擎的批处理等待窗口大小。窗口太小批次小利用率低窗口太大首个请求的延迟高。需要权衡。2.模型量化与优化对本地部署的模型进行量化INT8/FP4使用更快的推理引擎如vLLM, TensorRT-LLM。3.垂直/水平扩展升级GPU硬件或将负载分摊到多个推理后端实例。节点粒度某些节点执行时间极长阻塞整个流水线。考虑将大节点进一步拆分为更细粒度的子节点增加并发机会。但拆分过细会增加调度和数据传输开销需要平衡。网络延迟智能体部署在不同区域或云服务上跨网络调用延迟显著。将通信频繁的智能体部署在同一个可用区或内网。对于与外部API如OpenAI的通信考虑使用连接池、预建立长连接。5.2 常见问题与调试技巧数据依赖死锁现象系统挂起没有节点在执行就绪队列为空但DAG未完成。排查检查DAG是否存在循环依赖理论上DAG不应有环但动态生成的节点可能引入。检查自适应聚合逻辑是否某个下游节点在等待一个永远不会产生的数据事件可能是数据事件的类型或字段名定义不匹配。工具为调度器添加详细的日志记录每个数据事件的产生和消费情况。可视化DAG的实时状态高亮显示处于“等待”状态的节点及其缺失的依赖。部分结果不一致或错误现象最终输出结果混乱包含过时或矛盾的信息。排查这是自适应聚合和并发执行带来的典型挑战。下游节点可能基于上游节点的“早期”、“不完整”输出做出了决策而上游节点后续输出了修正信息。解决引入“数据版本”或“epoch”概念。同一轮完整计算周期内的数据属于一个版本。下游节点只消费已标记为“稳定”或“最终”版本的数据。或者设计节点具有“可重入”性当接收到新的、更完整的上游数据时可以重新计算并发布新的结果调度器需处理这种“修正”传播。资源耗尽OOM现象进程崩溃报内存或显存不足错误。排查并发执行大量模型实例尤其是大模型极易导致OOM。解决在调度器中实现资源感知调度。为每个节点预估其内存/显存消耗可通过历史运行数据拟合。执行引擎在运行任务前检查当前可用资源是否满足。不满足则将其保留在就绪队列优先调度资源需求小的任务。实施任务“排队”和“降级”机制如用较小模型替代。调试与追踪困难现象问题复现难错误来源不清。解决建立强大的可观测性体系。为每个任务用户请求生成唯一Trace ID并贯穿所有智能体调用、数据事件。集中式日志记录每个关键步骤节点开始/结束、数据发布/消费。使用分布式追踪系统如OpenTelemetry可视化整个请求的生命周期清晰看到时间花在了哪里数据流向了何方。实操心得在项目初期不要过度追求极致的并发和聚合。先实现一个正确的、串行的版本然后逐步引入并发优化。每增加一层优化如聚合、批处理都要进行充分的测试确保在提升性能的同时结果的正确性和稳定性没有受损。性能优化是一个持续测量、假设、验证、迭代的过程。6. 扩展思考MOSAIC理念的边界与未来MOSAIC所代表的“通过智能调度提升多组件系统效率”的思想其实超越了多智能体协作的范畴。它可以应用于任何存在复杂依赖和计算密集型组件的流水线系统。与工作流引擎的融合现有的工作流引擎如Airflow, Prefect擅长管理定时、依赖任务但在细粒度、动态、低延迟的智能计算任务调度上并不擅长。MOSAIC的调度策略可以与这些引擎结合由工作流引擎管理宏观任务流由MOSAIC调度器管理微观的智能体执行。成本与延迟的权衡更激进的并发和聚合意味着更多的资源同时被占用尤其是昂贵的GPU。调度器是否可以引入“成本”维度例如在非高峰时段采用激进策略以尽快完成任务在高峰时段采用保守策略以节省资源这需要调度策略具备多目标优化能力。面向失败的弹性设计在并发环境中一个节点的失败不应该导致整个流水线崩溃。调度器需要能检测故障进行重试可能使用备用模型或者根据已有部分结果执行“降级”流程生成一个虽不完美但可用的输出。从我个人的实践来看引入MOSAIC这类优化思路后一个中等复杂度的多智能体任务流程端到端延迟降低30%-50%是完全可以期待的尤其是在I/O等待和轻量级模型交互密集的场景下。它的价值不在于创造新的智能而在于让已有的智能更高效地协同这对于真正构建可用、好用的AI应用至关重要。开始设计你的下一个多智能体系统时不妨从画出一个清晰的DAG开始然后问问自己哪些环节可以提前哪些计算可以并行你的“指挥棒”准备好了吗