
1. 项目概述从单体Agent到复杂编排的必然演进如果你最近在折腾大模型应用开发尤其是基于LangChain、LlamaIndex或者自研框架构建多Agent系统那么“调度”和“并发”这两个词一定让你又爱又恨。爱的是它们代表了智能体Agent系统从简单的单轮对话玩具迈向能够处理复杂、长流程、多任务协同的“准生产级”应用的关键一步恨的是随之而来的状态管理、资源竞争、响应延迟等问题足以让任何一个开发者掉不少头发。我最初构建的Agent系统就是一个典型的“一问一答”式单体。用户提问Agent思考调用工具返回结果。一切都很美好直到需要处理这样的场景“帮我分析这份财报PDF提取关键财务指标然后与过去三年的数据进行趋势对比最后生成一份摘要报告。” 这个任务里PDF解析、数据提取、趋势计算、报告生成每一步都可能耗时数秒甚至数十秒让用户在前端傻等一个“思考中”的转圈圈体验极差。更不用说当多个用户同时发起类似请求时系统资源被单个长任务独占并发能力几乎为零。这就是“Agent 调度与并发”要解决的核心问题。它不是一个炫技的功能而是高可用、高性能Agent系统的基石。本次分享我将围绕一个实战中提炼出的架构模式展开读写分离、SSE早返回与异步SideCar。这个组合拳能有效解决Agent系统在长任务处理、实时反馈和高并发下的核心痛点。简单来说读写分离让计算和状态管理解耦SSEServer-Sent Events早返回让用户即时感知进度而异步SideCar则负责以非阻塞方式执行那些耗时或并行的子任务。下面我们就来拆解这个架构的每一块拼图。2. 核心架构设计读写分离、SSE与SideCar三位一体2.1 读写分离状态与计算的解耦之道在传统的同步Agent调用中HTTP请求线程会阻塞等待整个Agent工作流执行完毕。这个过程包含了LLM的多次往返思考、执行、工具调用可能是慢速的API或数据库查询、以及内部状态更新。所有这一切都绑定在同一个请求上下文里这是导致响应慢、资源利用率低的根本原因。读写分离的思路借鉴了数据库设计的经典模式但在Agent系统中我们分离的是“状态写入”和“任务执行”。写操作状态管理由一个轻量级的、快速的协调器Coordinator负责。它的核心职责是接收请求快速响应用户的初始请求。创建任务生成一个唯一的任务IDTask ID并将任务元信息用户输入、目标、初始参数持久化到任务状态存储如Redis、PostgreSQL中。这个操作必须非常快。触发执行向一个任务队列如RabbitMQ、Redis Streams、Kafka投递一个消息消息体包含Task ID和必要参数然后立即返回这个Task ID给客户端。状态更新监听任务执行过程中的状态事件由执行器发出并更新任务状态存储中的对应记录如“进行中正在解析PDF”、“完成数据提取”、“失败网络超时”。读操作任务执行由一个或多个执行器Executor负责。它们从任务队列中消费消息获取任务从队列中拉取任务消息。执行工作流根据Task ID从状态存储中加载完整上下文执行真正的Agent逻辑——调用LLM、运行工具函数、进行链式或图式推理。发布状态在执行的关键节点开始、子步骤完成、成功、失败将状态事件发布到事件流如Redis Pub/Sub、Server-Sent Events通道或直接回调协调器。注意这里的“读写”主要针对任务状态这个核心资源。协调器负责快速写入初始状态和更新状态执行器负责读取状态以执行业务逻辑并将结果状态写回。客户端则通过Task ID频繁“读”取最新状态。这样做的好处是立竿见影的用户请求在毫秒级内得到响应拿到了Task ID后端耗时的计算被异步化系统吞吐量不再受限于单个任务的执行时间。同时状态集中存储为实现任务监控、重试、回调提供了统一的基础。2.2 SSE早返回打造流式交互体验仅仅异步化还不够。用户拿到一个Task ID后如果只能通过不断轮询Polling来获取进度不仅增加服务器压力体验也不够“丝滑”。SSEServer-Sent Events技术就是为了解决这个问题而生。它允许服务器主动向客户端推送数据是实现“进度条”和“实时日志”的理想选择。在我们的架构中当协调器返回Task ID的同时会建立一个SSE连接通道。这个通道与Task ID绑定。随后执行器在运行过程中的每一个状态更新都会通过协调器或直接经由事件系统推送到对应的SSE通道。客户端监听这个通道就能实时看到“任务已接收正在排队…”“开始执行步骤1/4解析文档中…”“步骤1完成开始步骤2提取财务数据…”“步骤2完成发现关键指标营收同比增长15%…”“任务执行完毕最终报告已生成。”“早返回”的精髓就在这里我们不需要等待所有步骤完成而是在第一个耗时操作开始前就把Task ID和SSE连接建立好并返回给用户。后续的每一个进展都像直播一样推送给前端。这极大地提升了用户体验让用户感知到系统在“努力工作”而非“卡死”。实现上需要注意SSE连接的管理、超时重连、以及在海量并发下的连接数问题。通常我们会为每个活跃任务维持一个SSE连接并在任务完成后延迟关闭一段时间以允许客户端接收最终状态。2.3 异步SideCar并行化与资源隔离的利器“SideCar”模式通常指一个与主应用协同工作的辅助进程。在Agent系统中我将那些可并行、耗时长、或需要特殊环境的子任务抽象为独立的“SideCar任务执行器”。为什么需要它想象一下Agent工作流中的一个步骤“调用搜索引擎API获取最近三天的行业新闻”。这个操作可能因网络波动而耗时且它与Agent的核心推理逻辑是独立的。如果让主执行器同步等待这个调用会阻塞整个工作流。异步SideCar的工作流程如下主执行器运行到需要调用SideCar的节点时例如一个特定的工具调用。主执行器不直接执行该工具而是向一个SideCar任务队列发布一个子任务消息包含子任务类型和参数然后立即继续执行工作流的其他分支如果有或进入等待。一个专有的SideCar执行器集群可能部署在更适合处理I/O密集型或具有特定依赖的环境里消费该队列执行具体的耗时操作如网络请求、文件处理、模型推理。SideCar执行器完成后将结果发布到结果回调队列或通过事件系统通知主执行器。主执行器在合适的同步点或通过异步事件驱动获取子任务结果并将其整合到主工作流的上下文中。这种模式带来了多重优势并行化多个SideCar任务可以同时执行加速整体流程。非阻塞主执行器不被慢速I/O阻塞能更高效地调度LLM计算这通常是更宝贵的资源。资源隔离SideCar可以独立扩缩容。例如专门处理图像识别的SideCar可以使用GPU实例而主执行器使用CPU实例资源利用更合理。容错性单个SideCar任务失败可以通过重试队列处理不影响主工作流的状态管理主工作流会收到失败事件并决定如何处置。将读写分离、SSE和异步SideCar组合起来就形成了一个健壮的Agent异步调度系统。协调器作为总控快速响应并管理状态主执行器负责核心推理流程编排SideCar处理脏活累活SSE则将这一切的进展实时地呈现给用户。3. 核心组件实现与关键技术选型3.1 任务状态存储与协调器实现协调器的核心是快和稳。我推荐使用Redis作为首选的任务状态存储。原因如下它内存级的读写速度能满足高并发状态更新的需求它丰富的数据结构Hash, Sorted Set非常适合存储任务元信息、进度和中间结果它的过期TTL特性可以自动清理已完成的任务数据此外Redis的Pub/Sub或Streams可以直接用于事件驱动。一个典型的任务状态Hash结构可能如下{ “task_id”: “550e8400-e29b-41d4-a716-446655440000” “status”: “running” // pending, running, success, failed, cancelled “progress”: 40, // 百分比进度 “current_step”: “正在分析数据趋势” “created_at”: 1691234567 “updated_at”: 1691234570 “user_input”: “分析财报并对比趋势” “result_url”: “” // 最终结果存储地址 “error_msg”: “” }协调器可以用任何轻量级框架实现如Flask/FastAPI (Python)、Express (Node.js)、Spring Boot (Java)。它的主要端点有两个POST /api/task接收请求生成Task ID写入Redis向主任务队列投递消息建立SSE通道关联并立即返回{“task_id”: “...” “sse_url”: “/api/events?task_id...”}。GET /api/eventsSSE端点根据task_id将客户端连接到对应的事件流。关键实现细节任务ID生成使用UUID或雪花算法Snowflake确保全局唯一。队列选择Redis Streams或List是简单轻量的选择如果需要更强大的特性如优先级、死信队列RabbitMQ或Apache Kafka更合适。状态更新原子性使用Redis的WATCH/MULTI/EXEC命令或Lua脚本来确保并发下的状态更新安全。3.2 基于事件驱动的执行器设计执行器是业务逻辑的核心。它需要从队列中拉取任务加载上下文并驱动Agent工作流。这里的关键是将工作流的每一步都转化为事件。我倾向于使用有限状态机FSM或工作流引擎来管理Agent的执行步骤。每个步骤如parse_documentextract_datacall_llmcall_tool都是一个状态。执行器从一个状态转移到下一个状态每完成一个状态就发布一个状态事件。例如使用Python的transitions库或asyncio状态机可以清晰地定义状态流转class AgentTask: states [pending, parsing, extracting, reasoning, summarizing, success, failed] def __init__(self, task_id): self.task_id task_id self.machine Machine(modelself, statesAgentTask.states, initialpending) # 定义状态转移和回调函数 self.machine.add_transition(start_parsing, pending, parsing, afterdo_parse) self.machine.add_transition(on_parse_success, parsing, extracting, afterdo_extract) # ... 其他转移do_parsedo_extract这些方法会执行实际逻辑并在完成后触发下一个状态转移同时通过Redis Pub/Sub发布事件如{“task_id”: “...” “event”: “step_complete” “step”: “parsing” “data”: {...}}。执行器与SideCar的交互当do_extract方法需要调用一个慢速工具时它不会同步调用而是发布一个SideCar任务到特定队列如sidecar_tasks:web_search。将自身状态置为waiting_for_sidecar这是一个内部状态不一定暴露给用户。然后就可以去处理队列中的其他任务了如果是多线程/协程模型。当SideCar结果返回时触发一个on_sidecar_result事件使工作流从waiting_for_sidecar状态转移到下一个状态如reasoning。3.3 SSE服务端与客户端的实现要点服务端实现SSE核心是保持一个长连接并按照特定格式发送数据。以FastAPI为例from sse_starlette.sse import EventSourceResponse import asyncio app.get(“/api/events”) async def event_stream(task_id: str): async def event_generator(): # 1. 订阅与该task_id相关的Redis频道 pubsub redis_client.pubsub() await pubsub.subscribe(f“task_events:{task_id}”) # 2. 发送初始连接确认事件 yield {“event”: “connected” “data”: {“task_id”: task_id}} # 3. 循环监听Redis消息并转换为SSE格式 while True: message await pubsub.get_message(ignore_subscribe_messagesTrue, timeout10.0) if message and message[‘type’] ‘message’: event_data json.loads(message[‘data’]) # SSE格式 “event: event_type\ndata: json_data\n\n” yield { “event”: event_data.get(“event” “message”) “data”: json.dumps(event_data.get(“data”)) } # 检查任务是否已完成若完成发送结束事件并退出循环 task_status await get_task_status(task_id) if task_status in [‘success’ ‘failed’ ‘cancelled’]: yield {“event”: “end” “data”: json.dumps({“status”: task_status})} break await asyncio.sleep(0.1) return EventSourceResponse(event_generator())客户端实现则很简单使用EventSourceAPI即可const eventSource new EventSource(/api/events?task_id${taskId}); eventSource.onmessage (event) { const data JSON.parse(event.data); updateUI(data); // 更新进度条、日志显示等 }; eventSource.onerror (err) { console.error(“SSE连接错误” err); eventSource.close(); // 可以考虑回退到轮询 };避坑指南连接限制浏览器对同一域名下的SSE连接数有限制通常6个。对于需要同时监控多个任务的复杂前端需要考虑连接池管理或使用WebSocket更复杂但功能更强。代理与超时确保你的反向代理如Nginx配置了合适的超时时间以支持长连接。心跳机制在长时间没有数据推送时定期发送注释行以:开头作为心跳防止连接因超时被关闭。4. 异步SideCar的通信与协同模式SideCar模式的核心是松耦合的进程间通信。我实践下来有两种主流的通信方式各有优劣。4.1 基于消息队列的请求-响应模式这是最直观的方式如上文所述主执行器向任务队列发消息SideCar消费并执行然后将结果发回结果队列。优点解耦彻底SideCar可以独立部署、伸缩、升级。队列自带缓冲、重试、负载均衡能力。缺点增加了系统复杂性需要管理多个队列响应路径变长延迟稍高需要处理结果与主任务的关联通过correlation_id或task_id。技术选型轻量级/快速上手Redis Streams或Redis的RPUSH/BRPOP模式。Redis性能好但消息持久化、复杂路由能力较弱。生产级/高可靠RabbitMQ。它提供了强大的消息确认、持久化、路由键、死信队列等特性非常适合对可靠性要求高的场景。高吞吐/流处理Apache Kafka。如果SideCar任务量极大或者结果数据流需要被多个消费者处理如同时更新数据库和发送通知Kafka是更好的选择。4.2 基于RPC/HTTP的异步调用模式在这种模式下主执行器通过异步HTTP客户端如aiohttphttpx调用一个SideCar服务的API端点但不等待其响应或者只等待一个“已接收”的确认。SideCar服务在后台处理任务处理完成后通过一个预设的回调URLCallback URL主动通知主执行器。优点架构更简单直观符合常见的微服务调用模式调试方便。缺点主执行器需要维护回调端点SideCar服务需要有网络权限来回调如果SideCar服务挂掉需要额外的机制保证任务不丢失。实现示例 主执行器async def call_sidecar_tool(tool_name, params, callback_url): sidecar_url f“http://sidecar-service/{tool_name}” payload { “params”: params “callback”: callback_url # 例如 “http://executor-host/callback/task_id” “task_id”: current_task_id } # 使用异步客户端设置较短超时只关心请求是否成功发送 async with httpx.AsyncClient(timeout5.0) as client: try: resp await client.post(sidecar_url, jsonpayload) resp.raise_for_status() # 请求已成功发送至SideCar主流程可以继续 except Exception as e: # 记录失败可能触发重试或任务失败状态 logger.error(f“调用SideCar {tool_name} 失败: {e}”)SideCar服务在完成工作后向callback_url发送POST请求告知结果。选择建议对于内部系统、任务关键型的操作我推荐使用消息队列模式可靠性更高。对于与外部系统集成、或SideCar本身是第三方服务的情况回调模式可能更标准。5. 实战中的问题排查与性能优化5.1 常见问题与诊断流程即使架构设计得再完美在实际运行中也会遇到各种问题。下面是一个快速诊断清单问题现象可能原因排查步骤用户收到Task ID后SSE连接无任何进度更新。1. 协调器未正确发布初始事件。2. 执行器未消费任务队列。3. 执行器运行中但未发布状态事件。4. Redis Pub/Sub订阅失败或网络问题。1. 检查Redis中对应task_id的状态是否已从pending变为running。2. 查看任务队列的待处理消息数确认消费者是否在线。3. 查看执行器日志确认工作流是否启动有无报错。4. 在Redis命令行用SUBSCRIBE命令手动订阅任务事件频道看是否有消息。任务状态长时间卡在某个步骤。1. SideCar任务执行超时或失败。2. 主执行器等待SideCar结果的逻辑有bug。3. 工作流状态机陷入死循环。1. 检查SideCar执行器的日志和队列或回调记录。2. 检查主执行器中等待SideCar结果的状态处理代码确认超时机制是否生效。3. 在任务状态中增加更细粒度的“子步骤”日志定位卡住的具体位置。高并发下新任务响应变慢甚至超时。1. 协调器成为瓶颈数据库/Redis连接池耗尽。2. 任务队列堆积执行器消费不过来。3. SSE连接数过多耗尽服务器资源。1. 监控协调器的CPU、内存和数据库连接数。考虑水平扩展协调器无状态易于扩展。2. 监控队列长度动态扩缩容执行器实例。3. 考虑对SSE连接使用更高效的服务器如专用于推送的节点或评估改用WebSocket。SideCar任务执行成功但主任务未收到结果。1. 结果消息丢失队列未持久化。2. 回调URL调用失败网络、认证、服务不可用。3. 主执行器处理结果的代码有bug。1. 队列模式检查消息确认机制确保消息被“ack”。启用死信队列查看是否有失败消息。2. 回调模式检查SideCar服务的出站日志查看回调请求的状态码。检查主执行器回调端点的可访问性和日志。3. 在主执行器结果处理逻辑中添加详尽的日志和异常捕获。5.2 性能调优与伸缩策略执行器水平扩展执行器应该是无状态的状态在Redis中因此可以轻松地启动多个实例通过消息队列的竞争消费模式来分摊负载。使用Kubernetes HPA或基于队列长度的自定义伸缩策略。SideCar专业化与分组根据工具类型创建不同的SideCar队列和执行器集群。例如web_search队列由一组擅长HTTP请求的实例消费pdf_parse队列由一组安装了特定OCR库的实例消费。这样可以优化资源利用和隔离故障。连接池与资源复用无论是数据库连接、Redis连接还是HTTP客户端在执行器和SideCar中都必须使用连接池避免频繁创建销毁连接的开销。对于LLM API调用考虑使用支持批处理和流式响应的客户端。状态存储优化对于大型中间结果如解析后的全文文本、生成的图片不要直接存在Redis中应使用对象存储如S3、MinIO或文件系统在Redis中只存储引用地址。对任务状态Hash中的字段进行精简只保留必要信息。使用Redis的EXPIRE命令为已完成的任务状态设置合理的过期时间如24小时自动清理。SSE连接优化对于超大规模并发可以考虑使用专门的“推送网关”服务来管理SSE连接该网关与业务逻辑分离并通过内部消息系统如Redis Pub/Sub接收事件进行推送。这避免了业务服务器被大量长连接拖垮。5.3 监控与可观测性建设一个健壮的系统离不开监控。你需要关注以下指标业务指标任务创建速率、平均处理时间、成功率、各步骤耗时分布。系统指标协调器/执行器/SideCar的CPU、内存使用率Redis内存使用率及命中率消息队列的长度和消费延迟。链路追踪为每个task_id在整个系统中的流转协调器 - 队列 - 执行器 - SideCar - 回调注入追踪标识如OpenTelemetry Trace ID便于在出现问题时快速定位瓶颈和故障点。实现上可以将任务状态变更、步骤开始/结束等关键事件作为结构化日志输出并接入ELK或Loki等日志系统。同时将上述指标暴露给Prometheus并在Grafana中绘制仪表盘。6. 架构演进与高级模式探讨当基本的三件套读写分离、SSE、SideCar运行稳定后可以考虑向更高级的模式演进。动态工作流编排目前的执行器可能内嵌了固定的工作流逻辑。可以将其抽象出来使用像Apache Airflow、Prefect或Temporal这样的工作流编排引擎来定义和管理Agent的工作流。这样你可以可视化地编排复杂的、有条件分支的、可重试的任务流并将执行器变为通用的“Worker”只负责执行原子操作调用LLM、调用工具。这大大提升了复杂业务流程的维护性和可观测性。基于事件的协同Agent将每个Agent或工具都视为一个独立的事件发布/订阅者。它们通过一个全局的事件总线如Redis Pub/Sub或NATS进行通信。一个Agent完成任务后发布一个事件监听该事件的另一个Agent被触发执行。这种模式可以实现极其灵活和松耦合的多Agent协同但事件流的管理和调试复杂度也会显著增加。资源感知调度目前的调度可能只是简单的FIFO先进先出。可以引入更智能的调度器考虑任务的优先级、所需资源是否需要GPU、需要多少内存、以及当前系统的负载情况将任务分派到最合适的执行器或SideCar集群上。这需要更复杂的调度算法和资源元数据管理。回归到我们最初的起点实现Agent的调度与并发本质上是将“同步阻塞的智能思考”转变为“异步流式的服务协作”。读写分离奠定了异步化的基础SSE早返回重塑了用户体验而异步SideCar则释放了并行处理的潜力。这套模式并非银弹它会引入分布式系统固有的复杂性但在应对长耗时、多步骤、高并发的真实AI应用场景时它提供的解耦、可伸缩和实时反馈能力是构建可靠、可用、好用的智能体系统的必经之路。从我踩过的坑来看前期在状态管理、事件定义和错误处理上多花些时间设计远比后期在混乱的代码中挣扎要划算得多。