ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

深入理解 AI Agent · 多 Agent 编排 #02:Agent 之间怎么“说话“——通信、状态与冲突解决

深入理解 AI Agent · 多 Agent 编排 #02:Agent 之间怎么“说话“——通信、状态与冲突解决 多 Agent 编排解决了谁先谁后、谁干什么的问题但真正决定系统能否在生产环境跑稳的是三个更底层的机制Agent 之间怎么传消息、怎么共享状态、意见不一致时谁说了算。上一篇聊了多 Agent 编排的四种模式——串行、并行、条件分支、竞争。模式定义的是拓扑但拓扑之上还需要三样东西让协作真正运转传话Agent A 怎么把消息发给 Agent B共享记忆多个 Agent 怎么读到同一份上下文达成共识A 说 X、B 说 Y最终输出听谁的这三件事从简单到复杂构成了 Agent 协作的三个层次。本文从 Dream-SaaS 的真实代码出发逐层拆解。一、通信模式对比从成本视角看四种选择Agent 之间的通信本质上是一个工程权衡问题。没有最好的模式只有最合适的选择。1.1 同步调用最简单但级联失败是定时炸弹同步调用是最直觉的方式——A 调 B阻塞等结果。HTTP RPC、方法直接调用都属于此类。它的优点是简单调用链清晰错误处理直接调试方便。但隐患也很明显级联失败B 依赖 CC 依赖 DD 超时 10 秒整条链路全部阻塞。在微服务领域这叫雪崩在 Agent 编排领域同样存在。线程耗尽Agent 调用通常涉及 LLM 请求单次耗时 2-10 秒。如果用同步阻塞Tomcat 工作线程很快被占满。无法并行串行模式还好并行模式下同步等待意味着线程数 并行度资源开销大。同步调用适合同 JVM 内、延迟可控的场景。一旦跨进程或跨服务就需要考虑异步。1.2 异步消息解耦的代价是复杂度异步消息——A 发消息后不阻塞B 处理完后通过回调或轮询通知 A。解耦了发送方和接收方但引入了新的复杂度消息可靠性消息发出去了B 收到没有处理成功没有需要 ACK 和重试机制。顺序性先发的消息一定先到吗在网络分区场景下不一定。调试困难请求链路被打散日志难以串联需要 traceId 贯穿。异步消息适合跨服务、高吞吐、可容忍延迟的场景。1.3 事件总线 vs 消息队列别混淆概念这两个词经常被混用但侧重点不同维度事件总线Event Bus消息队列Message Queue模型发布-订阅一对多点对点一对一或消费组耦合发布方不关心谁订阅发送方指定目标或 Topic典型实现Guava EventBus、Spring ApplicationEventRabbitMQ、Kafka、RocketMQ适用场景同进程内领域事件通知跨服务可靠消息投递持久化通常不持久化持久化、支持回溯在 Agent 场景中事件总线适合同 JVM 内的 Agent 间通知如任务完成事件消息队列适合跨服务的可靠任务分发。用错了会出问题用事件总线做跨服务通信进程一重启消息全丢用消息队列做同进程通知杀鸡用牛刀且引入运维负担。1.4 共享状态高并发下的双刃剑共享状态严格来说不是通信但它是 Agent 间交换信息的一种方式——多个 Agent 读写同一个状态存储内存 Map、Redis、数据库通过你写我读传递信息。优点速度快内存读写微秒级、实现简单、天然支持多消费者。致命缺点竞态条件。两个并行 Agent 同时写同一个 key后写的覆盖先写的结果不可预测。下文状态共享一章会展开作用域隔离的应对方案。1.5 成本-复杂度对比通信方式实现成本运维复杂度可靠性延迟适用场景同步调用低低中级联失败风险低同 JVM、短链路异步消息中中高需 ACK/重试中跨服务、可容忍延迟事件总线低低低不持久化极低同进程事件通知消息队列高高高中跨服务可靠投递共享状态低中中竞态风险极低高频读、低频写1.6 Dream-SaaS 实战HybridAgentChannel 双通道设计Dream-SaaS 同时存在同 JVM 内的 Agent 协作和跨服务的 Agent 调用。如果全部走远程 HTTP同 JVM 内的调用白白增加网络开销如果全部走内存跨服务的 Agent 又无法通信。解决方案是HybridAgentChannel——同 JVM 走 InMemory跨服务走 A2A HTTP运行时根据消息标记自动路由// HybridAgentChannel.java — 同JVM InMemory跨服务 A2A HTTP Slf4j public class HybridAgentChannel implements AgentChannel { private final InMemoryAgentChannel inMemory; private final A2AClient a2aClient; private final AgentCollaborationProperties properties; Override public void send(AgentMessage message) { if (message null) return; if (shouldUseA2a(message)) { delegateRemote(message); // 跨服务走 A2A return; } inMemory.send(message); // 同 JVM 走内存 } private void delegateRemote(AgentMessage message) { try { String base properties.getA2a().getRemoteBaseUrl(); String taskId a2aClient.sendTask( base, message.toAgent(), message.messageType(), message.payload() ! null ? message.payload() : Map.of()); log.info(a2a delegated taskId{} toAgent{}, taskId, message.toAgent()); } catch (Exception e) { // 远程调用失败自动降级到本地内存通道 log.warn(a2a delegate failed, fallback in-memory: {}, e.toString()); inMemory.send(message); } } }设计要点有三路由判定shouldUseA2a()检查配置开关、远程地址和消息 payload 中的transporta2a标记。只有显式标记为远程的消息才走 A2A默认走内存避免不必要的网络开销。自动降级远程调用抛异常时catch 住并 fallback 到inMemory.send()。这意味着网络抖动不会导致消息丢失而是退化为本地执行。降级有代价fallback 到本地意味着远程 Agent 的专业能力不可用会由本地的通用 Agent 兜底。这个代价需要业务方知晓不能静默降级。二、JSON-RPC 2.0 与 A2A 协议跨服务 Agent 通信需要一个标准协议。Dream-SaaS 选择了 JSON-RPC 2.0 作为 A2AAgent-to-Agent协议的底座。2.1 为什么选 JSON-RPC 而非 REST/gRPC维度RESTgRPCJSON-RPC 2.0协议复杂度中资源建模高需 proto 定义低method params浏览器/调试友好好差二进制好纯 JSON类型安全弱强弱流式支持需 SSE/WebSocket原生支持不支持需扩展Agent 场景适配需为每个操作建模 REST 资源需维护 proto 文件method 即动作天然适合 RPC 语义Agent 间的交互本质是发指令、拿结果是 RPC 语义而非资源 CRUD。JSON-RPC 的method params id模型足够轻量且与 A2A 规范的消息结构天然对齐。gRPC 性能更好但在 Agent 场景中LLM 推理耗时通常在秒级HTTP/JSON 的序列化开销毫秒级完全可以忽略。2.2 A2AProtocol 消息模型A2A 协议的核心消息结构非常精简// A2AProtocol.java — A2A JSON-RPC 2.0 消息模型 public final class A2AProtocol { public static final String JSONRPC 2.0; // 核心消息request / result / error 三态合一 public record A2AMessage( String jsonrpc, String method, Object params, String id, Object result, A2AError error) { public static A2AMessage request(String method, Object params, String id) { return new A2AMessage(JSONRPC, method, params, id, null, null); } public static A2AMessage result(String id, Object result) { return new A2AMessage(JSONRPC, null, null, id, result, null); } public static A2AMessage error(String id, int code, String message) { return new A2AMessage(JSONRPC, null, null, id, null, new A2AError(code, message)); } } // 任务状态PENDING → RUNNING → COMPLETED / FAILED public record TaskState(String taskId, String status, Object output, String error) {} // 三个核心方法 public static final String METHOD_TASKS_SEND tasks/send; public static final String METHOD_TASKS_GET tasks/get; public static final String METHOD_AGENTS_CARD agents/card; public static final String STATUS_PENDING PENDING; public static final String STATUS_RUNNING RUNNING; public static final String STATUS_COMPLETED COMPLETED; public static final String STATUS_FAILED FAILED; }三个方法的职责tasks/send提交任务立即返回taskId异步不阻塞等待执行结果。tasks/get根据taskId查询任务状态和输出客户端通过轮询获取结果。agents/card获取 Agent 的能力描述AgentCard用于服务发现和能力协商。任务状态机是线性的四态模型PENDING → RUNNING → COMPLETED/FAILED。没有 PAUSED、CANCELLED 等中间态——刻意保持简单因为 Agent 任务通常不支持暂停和恢复加了反而增加实现复杂度。2.3 A2AClient异步提交 轮询模式A2A 采用提交-轮询模式而非提交-等待模式因为 LLM 任务耗时长长时间持有 HTTP 连接不现实。客户端的pollTaskAsync方法用ScheduledExecutorService实现定时轮询用CompletableFuture做异步编排// A2AClient.java — 异步轮询获取任务结果 public CompletableFutureA2AProtocol.TaskState pollTaskAsync( String baseUrl, String taskId, int maxAttempts, long sleepMs) { int attempts Math.max(1, maxAttempts); long interval Math.max(50L, sleepMs); CompletableFutureA2AProtocol.TaskState cf new CompletableFuture(); AtomicInteger tryCount new AtomicInteger(0); AtomicReferenceScheduledFuture? handle new AtomicReference(); Runnable tick () - { if (cf.isDone()) return; int n tryCount.incrementAndGet(); try { A2AMessage req A2AMessage.request( METHOD_TASKS_GET, Map.of(taskId, taskId), UUID.randomUUID().toString()); A2AMessage res post(baseUrl, req); if (res.error() ! null) { cf.completeExceptionally( new IllegalStateException(res.error().message())); cancelPoll(handle); return; } TaskState state objectMapper.convertValue( res.result(), TaskState.class); if (state ! null (STATUS_COMPLETED.equals(state.status()) || STATUS_FAILED.equals(state.status()))) { cf.complete(state); // 终态完成 Future cancelPoll(handle); return; } if (n attempts) { cf.completeExceptionally( new IllegalStateException(A2A task poll timeout: taskId)); cancelPoll(handle); } } catch (Throwable t) { cf.completeExceptionally(t); cancelPoll(handle); } }; // 立即执行第一次之后按固定间隔轮询 handle.set(POLL_SCHED.scheduleAtFixedRate(tick, 0L, interval, MILLISECONDS)); cf.whenComplete((r, e) - cancelPoll(handle)); return cf; }几个实现细节值得注意轮询间隔下界 50msMath.max(50L, sleepMs)防止配置过小导致轮询风暴。首次立即执行scheduleAtFixedRate(tick, 0L, ...)的初始延迟为 0第一次查询不等待减少短任务的不必要延迟。自动取消cf.whenComplete注册回调无论成功、失败还是超时都会取消定时任务避免线程泄漏。轮询线程池隔离POLL_SCHED是独立的 2 线程调度池不与业务线程争抢资源。2.4 A2AServer异步任务处理与状态流转服务端收到tasks/send后不阻塞等待 Agent 执行完成而是将任务提交到 worker 线程池异步执行立即返回taskId// A2AServer.java — 收到任务后异步执行状态流转 PENDING→RUNNING→COMPLETED/FAILED private A2AMessage send(A2AMessage request) { MapString, Object params (MapString, Object) request.params(); String taskId UUID.randomUUID().toString(); String toAgent String.valueOf(params.getOrDefault(agent, )); MapString, Object payload (MapString, Object) params.get(payload); A2ALocalExecutor executor resolveExecutor(toAgent, payload); if (executor null) { taskStore.put(new TaskState(taskId, STATUS_FAILED, null, 无匹配的 executor)); return A2AMessage.error(request.id(), -32000, 无匹配的本地 A2A executor); } // 状态流转PENDING → RUNNING taskStore.put(new TaskState(taskId, STATUS_PENDING, null, null)); taskStore.put(new TaskState(taskId, STATUS_RUNNING, null, null)); linkRunFromPayload(taskId, payload); // 建立 runId/orchTaskId 索引 // 提交到 worker 池异步执行 String input extractInput(payload, params); String agentName executor.agentName(); workerPool.submit(() - runTask(taskId, agentName, executor, input, payload)); // 立即返回 taskId不等待执行结果 return A2AMessage.result(request.id(), Map.of(taskId, taskId, status, STATUS_RUNNING, agent, agentName)); } private void runTask(String taskId, String agentName, A2ALocalExecutor executor, String input, MapString,Object payload) { try { String output executor.execute(input, payload); taskStore.put(new TaskState(taskId, STATUS_COMPLETED, Map.of(output, output ! null ? output : , agent, agentName), null)); } catch (Exception e) { taskStore.put(new TaskState(taskId, STATUS_FAILED, null, e.getMessage())); } }worker 线程池使用BoundedDaemonExecutors.create(a2a-task-, 2, 8, 32)创建——核心 2 线程、最大 8 线程、队列 32有界设计防止任务积压导致 OOM。2.5 混合通道降级远程失败自动 fallback回到 HybridAgentChannel 的降级逻辑当 A2A 远程调用抛异常时catch 住后走inMemory.send()。这意味着即使远程服务不可用消息也不会丢而是由本地 Agent 处理。但这里有一个容易被忽略的问题降级可能导致重复执行。如果远程服务实际上收到了任务并在执行只是响应超时了客户端 fallback 到本地再执行一次同一个任务就被执行了两遍。对于幂等操作如查询这不是问题但对于有副作用的操作如发消息、写数据重复执行可能造成业务异常。Dream-SaaS 的做法是在 payload 中携带runId和orchestratorTaskId服务端通过这些标识做去重检查——如果同一个 taskId 已经在处理中直接返回已有状态不重复执行。三、状态共享的三种模式与陷阱通信解决了消息怎么传的问题状态共享解决的是数据怎么放的问题。多个 Agent 协作时中间结果、上下文、执行进度需要被各方访问。3.1 全局状态桶快但有毒最简单的方案一个全局MapString, Object所有 Agent 都往里写、往里读。优点实现零成本读写速度极快。致命问题竞态条件。两个并行 Agent 同时写同一个 key后写覆盖先写。更隐蔽的是读-改-写非原子性——Agent A 读到count1加 1 后写回 2Agent B 也读到 1加 1 写回 2。最终结果是 2 而不是 3。在 Agent 场景中这种 bug 尤其难排查因为 LLM 的输出本身就有不确定性结果不对很难区分是模型问题还是状态污染。3.2 作用域隔离安全但复杂将状态按作用域隔离——每个 Agent 只能写自己的命名空间跨 Agent 的数据通过显式消息传递。例如state.agentA.result ... state.agentB.result ... state.merged.output ... // 只有汇总节点能写这从根本上消除了并行写冲突但代价是状态结构需要提前规划灵活性降低。Agent 之间共享数据需要经过编排层中转增加了一层间接性。3.3 事件溯源 快照可追溯但开销大不直接保存当前状态而是保存所有状态变更事件Event Sourcing。当前状态通过回放事件派生定期做快照Snapshot加速回放。优点完整审计轨迹可以回溯任意时间点的状态对调试 Agent 行为极有价值。缺点实现复杂度高存储开销大事件 schema 变更需要版本管理。3.4 核心观点状态即契约状态即契约——Agent 之间共享的每一个状态字段都是一个隐含的接口约定。字段名、类型、语义、生命周期任何一方擅自变更都会导致协作失败。这和微服务中的数据库共享是反模式同理。当两个 Agent 直接读写同一个状态键时它们之间的耦合比任何 API 调用都紧——API 有版本号、有 schema 校验而共享状态只有一个约定俗成的 key 名。实践中应该最小化共享状态能通过消息传递的就不要共享存储。共享状态必须有 schema哪怕是一个简单的 Java record 或 JSON Schema。写操作收口只允许特定角色如汇总节点写共享状态其他角色只读。状态有 TTL过期自动清理避免内存泄漏和脏数据干扰。3.5 Dream-SaaS 实战A2ATaskStore 双索引设计A2ATaskStore 负责任务状态的存储和查询采用内存热缓存 可选 Redis 持久化的双层架构// A2ATaskStore.java — 内存热缓存 可选Redis持久化 public class A2ATaskStore { // 主存储taskId → TaskState private final MapString, TaskState memory new ConcurrentHashMap(); // 双索引用于编排排障从业务标识反查 taskId private final MapString, String runIndexMemory new ConcurrentHashMap(); private final MapString, String orchTaskIndexMemory new ConcurrentHashMap(); private final ObjectProviderStringRedisTemplate redisProvider; public void put(TaskState state) { memory.put(state.taskId(), state); if (!redisActive()) return; try { String json objectMapper.writeValueAsString(state); redis().opsForValue().set(key(state.taskId()), json, ttl()); } catch (Exception e) { log.warn([A2ATaskStore] redis put failed taskId{}: {}, state.taskId(), e.getMessage()); } } /** runId → taskId 索引编排排障时从运行ID反查任务 */ public void linkRun(String runId, String taskId) { runIndexMemory.put(runId.trim(), taskId.trim()); if (!redisActive()) return; try { redis().opsForValue().set(runKey(runId.trim()), taskId, ttl()); } catch (Exception e) { log.warn(linkRun failed: {}, e.getMessage()); } } }设计要点内存优先get()先查内存 ConcurrentHashMap命中直接返回微秒级。未命中才查 Redis查到后回填内存。这保证了同实例内的查询性能。Redis 可选通过ObjectProviderStringRedisTemplate注入Redis 不可用时自动降级为纯内存模式不影响核心功能。redisActive()还检查task-store配置项支持memory/redis/auto三种模式。双索引设计除了主存储taskId → TaskState还维护了runId → taskId和orchTaskId → taskId两个索引。runId是整个编排运行的唯一标识orchTaskId是编排层子任务 ID。当线上出问题需要排障时可以从任意一个业务标识反查到 A2A 任务查看执行状态和输出。这在微服务链路追踪中相当于多入口关联 ID。TTL 过期Redis key 设置了 TTL默认 24 小时可配置避免任务状态无限堆积。内存缓存没有主动过期但因为 JVM 重启即清空且 Redis 是持久层的 source of truth不会造成长期问题。Redis 失败不影响主流程所有 Redis 操作都被 try-catch 包裹失败只打 warn 日志不抛异常。这是缓存的正确姿态——缓存挂了不能拖垮主业务。四、当 Agent 意见不合冲突解决的五层机制并行执行多个 Agent 时最棘手的问题不是怎么让它们跑起来而是结果不一样怎么办。Dream-SaaS 设计了五层冲突解决机制从轻到重依次为投票、优先级、仲裁 Agent、回退策略、最终一致性。4.1 第一层投票VOTE多个 Agent 输出相同语义的结果时取多数票。但直接字符串比较会遇到一个问题LLM 的输出有随机性同一个答案可能因为多了一个空格、换了一行、大小写不同就被判定为不同导致永远无法达成多数共识。ResultMerger 的mergeVote方法通过normalizeVoteKey规范化后再计票// ResultMerger.java — 规范化空白后多数票减少LLM微调导致的永不共识 private String mergeVote(ListTaskResult results) { MapString, ListTaskResult buckets new LinkedHashMap(); for (TaskResult r : results) { String key normalizeVoteKey(r.output()); // 规范化后再分桶 buckets.computeIfAbsent(key, k - new ArrayList()).add(r); } Map.EntryString, ListTaskResult winner buckets.entrySet().stream() .max(Comparator.comparingInt(e - e.getValue().size())) .orElse(null); if (winner null || winner.getValue().isEmpty()) return 无共识; // 返回票数最高的原文不是规范化后的文本 return winner.getValue().getFirst().output(); } // trim 合并连续空白 转小写 static String normalizeVoteKey(String output) { if (output null) return ; return output.trim().replaceAll(\\s, ).toLowerCase(); }规范化做了三件事trim()去除首尾空白、replaceAll(\\s, )将连续空白包括换行、制表符合并为单个空格、toLowerCase()统一小写。这样 Yes 和 yes 和 yes\n 会被归为同一桶。但返回的是原文而非规范化文本——规范化只用于计票输出保持原样避免丢失格式信息。4.2 第二层优先级PRIORITY投票假设所有 Agent 地位平等但现实中不同 Agent 的权威性不同。例如代码审查场景中Senior Reviewer 的意见权重应该高于 Junior Reviewer。PRIORITY 策略按agentNames白名单顺序取权威结果// ResultMerger.java — 按 agentNames 白名单顺序取权威结果 private String mergePriority(ListTaskResult results, ListString priorityAgentNames) { if (priorityAgentNames ! null !priorityAgentNames.isEmpty()) { for (String name : priorityAgentNames) { if (name null || name.isBlank()) continue; for (TaskResult r : results) { if (name.equals(r.agentName())) { log.info([ResultMerger] priority selected: {}, name); return r.output(); // 按白名单顺序命中即返回 } } } } // 白名单中没有任何成功结果取首个成功 TaskResult top results.getFirst(); return top.output(); }这里的优先级不是分数加权而是权威顺序——白名单中排第一的 Agent 如果成功了直接取它的结果不看其他 Agent。只有白名单中的 Agent 全部失败时才回退到取首个成功结果。这种设计适合主-备模式主 Agent 正常时用主的主挂了才用备的。4.3 第三层仲裁 Agent投票和优先级都是机械合并不理解语义。当结果冲突复杂到无法用规则解决时引入一个专门的仲裁 Agent由它阅读所有分歧方的输出给出最终裁决。仲裁 Agent 的优势是灵活——它能理解语义、权衡利弊、给出理由。但它有两个硬伤延迟翻倍原本 N 个 Agent 并行执行现在还要等仲裁 Agent 读完 N 份结果再生成裁决总延迟 max(子 Agent) 仲裁 Agent。裁决质量不独立于模型能力仲裁 Agent 自身同样依赖 LLM 推理存在产生幻觉或判断偏差的可能。一旦裁决者的输出质量不稳定正确候选反而可能被推翻成错误结论。Dream-SaaS 用RoleAgentOutputContract对角色输出做硬约束防止执行角色越权做任务分解// RoleAgentOutputContract.java — 防止执行角色越权做任务分解 final class RoleAgentOutputContract { static final String RULES 【输出硬约束】 - 你是执行角色不是任务规划器。 - 直接给出对本角色职责的结论/话术/清单使用中文。 - 禁止输出任务拆解 JSON禁止 id/capabilities/dependencies 数组。 - 禁止空数组 []信息不足时说明缺什么并给出可执行的假设结论。 - 可引用前置任务结果但不要复述或续写规划稿。 ; static String withRules(String systemPrompt) { return (systemPrompt null ? : systemPrompt) RULES; } }这段 Prompt 约束会拼接到每个执行角色 Agent 的 system prompt 末尾从输入端防止越权。但 Prompt 约束不是 100% 可靠的——LLM 偶尔会不听话。因此还需要输出端的检测器做双重防护。反直觉观点仲裁 Agent 不是万能解药。很多团队一遇到结果冲突就想加个仲裁 Agent但仲裁本身引入了新的故障点和延迟。更稳健的做法是先用规则投票/优先级解决 80% 的冲突剩下 20% 真正需要语义理解的才交给仲裁 Agent同时对仲裁结果也要做校验——如果仲裁结果与所有候选结果差异过大应该标记为低置信度而非直接采信。4.4 第四层回退策略四档降级当所有 Agent 都失败了或者合并结果不可信时需要有兜底方案。Dream-SaaS 的ChatFallbackAdapter实现了SingleAgentFallback接口提供四档降级A2A 远程调用正常路径调用远程专业 Agent。InMemory 本地通道远程不可用时fallback 到本地同 JVM Agent。单 Agent 兜底编排整体失败时ChatFallbackAdapter作为通用对话 Agent 直接回答用户问题。它带有对话窗口记忆和 RAG 检索能力保证即使编排全挂了用户也能得到一个回答。错误提示连兜底 Agent 也失败时返回友好的错误信息而非 500 白屏。// ChatFallbackAdapter.java — 通用对话兜底实现 SingleAgentFallback Component public class ChatFallbackAdapter implements LocalAgentAdapter, SingleAgentFallback { private static final String SYSTEM_BASE 你是业务助手。可参考近期对话与长期记忆规则/事实作答用简洁中文。; Override public String answer(String userRequest) { return answer(userRequest, orch-default); } Override public String answer(String userRequest, String conversationId) { String sessionId (conversationId null || conversationId.isBlank()) ? orch-default : conversationId.trim(); String system buildSystem(structuredBean, sessionId); String answer chatClientService.sendMessageWithConversation( system, userRequest null ? : userRequest, sessionId); return answer null ? : answer; } }降级的核心原则是优雅降级Graceful Degradation每一档降级都应该让用户感知到功能少了但还能用而不是系统崩了。兜底 Agent 的回答质量可能不如专业 Agent但至少不会让用户面对一个错误页面。4.5 第五层最终一致性有些场景不需要实时一致。例如生成一份多维度分析报告三个 Agent 分别从不同角度分析合并时如果某一方结果还没回来可以先返回部分结果缺失部分标注生成中后续通过 WebSocket 或轮询补全。最终一致性不是不处理冲突而是按业务容忍度分级强一致如支付结果必须等所有 Agent 成功否则整体失败。弱一致如内容推荐允许部分失败返回最佳可用结果。最终一致如报告生成先返回部分结果异步补全。4.6 ResultMerger 四种合并策略总览策略适用场景延迟开销语义理解失败处理CONCAT结果互补需要全量展示无无失败结果单独列出VOTE同类任务结果应趋同无无取多数票忽略少数PRIORITY主备模式有权威 Agent无无白名单全失败取首个SYNTHESIZE需要综合多角色结论低确定性合成失败结果单独列出SYNTHESIZE 策略是确定性合成——不调用 LLM而是按固定模板生成综合结论 分角色详情。综合结论取每个角色输出的前 160 字作为摘要分角色详情保留完整输出。这比让 LLM 做合成更可控不会引入幻觉但灵活性有限。4.7 AgentOutputFailureDetector识别假成功这是一个容易被忽略但极其重要的组件。Agent 返回了一个字符串HTTP 状态码 200没有抛异常——但输出内容其实是错误的。Dream-SaaS 称之为假成功分两类传输/模型错误输出本身是异常堆栈或 HTTP 错误信息但被 Agent 包装成了正常返回。角色输出污染执行 Agent 没有给出本职结论而是把任务分解器的 JSON 规划原样吐回来了。// AgentOutputFailureDetector.java — 识别假成功 final class AgentOutputFailureDetector { private static final Pattern TASK_PLAN_ARRAY Pattern.compile( (?s)^\\s*\\[\\s*\\{.*\id\\\s*:.*\capabilities\\\s*:.*\dependencies\\\s*:.*}\\s*]\\s*$); static OptionalString detectFailure(String output) { if (looksLikeTransportOrModelFailure(output)) { return Optional.of(failureMessage(output)); } if (looksLikeTaskPlanLeak(output)) { return Optional.of(角色输出疑似任务规划JSON非本职结论已标失败); } return Optional.empty(); } // 检测传输错误Exception/Error 开头、HTTP 4xx/5xx、Spring AI 重试异常 static boolean looksLikeTransportOrModelFailure(String output) { if (output null) return false; String t output.trim(); if (t.isEmpty()) return false; if (t.startsWith(Exception:) || t.startsWith(Error:)) return true; String lower t.toLowerCase(); if (lower.contains(nontransientaiexception) || lower.contains(no response body available) || lower.contains(org.springframework.ai.retry)) return true; return t.length() 500 lower.matches((?s).*\\bhttp\\s[45]\\d\\d\\b.*); } // 检测任务规划泄漏数组包含 id/capabilities/dependencies static boolean looksLikeTaskPlanLeak(String output) { if (output null) return false; String t output.trim(); if (t.equals([])) return true; // 去掉 markdown 围栏 json ... if (t.startsWith()) { int nl t.indexOf(\n); int end t.lastIndexOf(); if (nl 0 end nl) t t.substring(nl 1, end).trim(); } if (!t.startsWith([) || !t.contains(\capabilities\) || !t.contains(\dependencies\)) return false; if (t.length() 8000) return false; return TASK_PLAN_ARRAY.matcher(t) || (t.contains(\id\) countTaskLikeObjects(t) 2); } }检测器在 ResultMerger 合并前被调用——如果检测到假成功该结果会被标记为失败不参与合并而是进入失败列表。这和RoleAgentOutputContract的 Prompt 约束形成双重防护Prompt 约束在输入端预防检测器在输出端兜底。为什么需要正则匹配而不是让 LLM 判断因为检测器运行在合并阶段如果再调一次 LLM 来判断这个输出是不是错误延迟和成本都不可接受。正则和关键词匹配虽然不够智能但速度快微秒级、确定性强覆盖了高频场景。五、结语Agent 间的通信、状态共享和冲突解决是多 Agent 编排从能跑到跑稳的关键跨越。几个关键结论通信选型没有银弹——同 JVM 走内存、跨服务走 A2A、事件通知走总线、可靠投递走队列。HybridAgentChannel 的双通道设计是务实选择但降级时要防重复执行。状态即契约——共享状态的耦合度高于 API 调用必须最小化、schema 化、写操作收口、设置 TTL。双索引设计让排障从大海捞针变成精准定位。冲突解决要分层——投票解决 80% 的机械冲突优先级处理主备场景仲裁 Agent 留给真正需要语义理解的少数情况。仲裁 Agent 本身也是故障点不是万能解药。假成功比明确失败更危险——HTTP 200 不代表结果正确必须在输出端检测异常内容。Prompt 约束 输出检测的双重防护比任何单一手段都可靠。下一篇将进入 Dream-SaaS 的真实编排架构——Supervisor Orchestrator 双层设计看编排层如何调度 Agent、管理 DAG、处理容错以及 Supervisor 三层路由策略的具体实现。下篇预告《深入理解 AI Agent七Supervisor 与 Orchestrator——Dream-SaaS 双层编排架构拆解》深入理解 AI Agent · 编排篇中
返回列表