尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Spring AI Graph并行执行与HITL实战:构建高可靠智能体工作流

Spring AI Graph并行执行与HITL实战:构建高可靠智能体工作流 1. 项目概述当Spring AI Graph遇上并行与人工干预如果你正在用Spring AI构建一个稍微复杂点的智能体Agent比如一个能自动分析日志、执行系统命令并生成报告的自动化工具你很可能已经遇到了两个绕不开的“坎”效率和可靠性。单线程的、线性的执行流程在面对多个独立任务时慢得像蜗牛而一个完全自主的AI在面对边界模糊或高风险操作时又让人不敢完全放手。这正是“Spring AI Graph从0到Supervisor二并行执行HITL实战”这个标题背后要解决的核心痛点。简单来说这个项目探讨的是如何利用Spring AI Graph框架构建一个更强大、更实用的AI智能体系统。它不再是一个简单的问答机器人而是一个具备工作流编排能力的“数字员工”。其中“并行执行”是为了榨干硬件性能让多个任务同时跑起来极大提升处理吞吐量“HITL”Human-In-The-Loop人在回路则是引入关键的人工审核或决策环节在自动化流程中设置安全阀确保AI的行为可控、结果可靠。想象一下你有一个智能体需要同时清理十个服务器的缓存、检查服务状态并拉取最新日志这些任务互不依赖并行执行能节省大量时间。而在执行“重启服务”或“删除文件”这类高风险命令前自动暂停并弹窗让你确认这就是HITL的价值。本文将从一个一线开发者的视角手把手带你深入Spring AI Graph的核心拆解如何从零开始设计并实现支持并行执行和HITL机制的Supervisor监督者智能体。我们会抛开那些空洞的概念直接进入代码和设计思路分享我在实际项目中趟过的坑、总结的最佳实践以及如何平衡自动化与安全性。无论你是刚开始接触Spring AI还是已经在用Graph构建复杂工作流这里都有你能直接“抄作业”的干货。2. 核心架构与设计思路拆解在动手写代码之前我们必须把架构想清楚。Spring AI Graph的核心思想是将AI智能体的行为建模为一个有向图图中的节点代表一个处理单元如调用大模型、执行工具函数、条件判断边代表执行流的方向。默认情况下图的执行是顺序的沿着边一条路走到黑。要实现并行和HITL我们需要在这个模型上动手术。2.1 理解Spring AI Graph的执行模型与瓶颈Spring AI Graph的执行引擎本质上是单线程遍历图。它根据每个节点的执行结果和边的条件决定下一个激活的节点。这种模型对于描述清晰的决策流程比如“先分析问题再选择工具最后执行”非常优雅但遇到以下场景就力不从心了批量独立任务需要处理100个文档每个文档的摘要生成任务彼此毫无关联。多路信息获取为了回答一个问题需要同时查询数据库、调用外部API和搜索向量知识库。实时监控与响应需要同时监听多个消息队列或数据流。如果强行用顺序执行耗时将是所有任务耗时的总和这完全无法接受。因此并行化不是可选项而是生产环境必须项。2.2 并行执行的设计策略分支、聚合与资源隔离实现并行并不是让Graph引擎本身变成多线程的。更可行的策略是将需要并行处理的任务封装为一个特殊的“并行节点”。这个节点的职责是任务分派接收一批输入数据创建多个独立的子任务。并行执行利用Java的并发工具如CompletableFuture、线程池同时执行这些子任务。结果聚合等待所有子任务完成将结果收集、合并传递给图的下一个节点。这里的关键设计点在于资源隔离。每个子任务可能都会调用AI模型消耗Token或外部服务。我们必须确保线程安全使用的工具类或客户端必须是线程安全的或者为每个任务创建独立实例。流量控制避免对下游服务如OpenAI API造成突发流量冲击需要通过信号量或限流器控制并发数。错误处理一个子任务失败不应该导致整个并行节点崩溃。我们需要定义策略是快速失败、忽略错误还是重试在我的实践中我通常会定义一个ParallelNode它继承自AbstractNode在其execute方法中实现上述逻辑。输入可以是一个列表输出也是一个列表或聚合后的对象。2.3 HITL人在回路的集成模式中断、通知与恢复HITL机制的核心是在自动化流程中插入一个等待人工干预的状态。在Graph中这可以建模为一个特殊的“人工审核节点”。该节点的行为模式是中断执行当执行流到达此节点时图的整体状态被持久化保存到数据库或Redis然后当前执行线程挂起或返回一个“等待中”的状态。通知用户通过预设的渠道如WebSocket、邮件、企业内部IM机器人向相关人员发送审核请求附上当前上下文、AI建议的操作和需要人工决策的点。接收决策提供一个接口如REST API供人工审核后提交决策结果批准、拒绝、修改参数。恢复执行根据人工决策从持久化的状态中恢复Graph的执行并注入决策结果使流程继续向下进行。这里最棘手的是状态管理。Spring AI Graph提供了State的概念但默认可能是在内存中。为了实现HITL我们必须配置一个持久化的StateStore例如使用Redis或数据库。这样当流程暂停时整个对话历史和上下文都能被完整保存并在恢复时无损还原。2.4 Supervisor的角色编排者与决策者“Supervisor”在这个语境下不再是简单的单个节点而是一个更高层次的编排框架。它负责动态图构建根据初始任务决定是否需要并行分支在哪里插入HITL节点。异常监控与处理监控并行任务的执行情况处理子任务失败决定是否触发降级流程或人工介入。生命周期管理管理HITL流程的发起、超时和决策结果应用。你可以把Supervisor理解为这个智能体系统的“大脑皮层”它不处理具体的脏活累活那些由工具节点做但它规划工作流、评估风险、并在关键时刻呼叫人类“前额叶”来做决断。实现上Supervisor本身也可以是一个Graph或者是一个管理多个Graph实例的服务。3. 关键组件与核心代码实现理论说得再多不如一行代码。接下来我们深入到实现层面看看关键组件如何落地。我会以构建一个“服务器运维诊断智能体”为例它需要并行检查多个服务器指标并在执行重启操作前请求人工确认。3.1 构建可并行执行的Graph节点首先我们定义一个ParallelProcessingNode。假设我们需要并行执行ping命令检查一组服务器的连通性。Component public class ParallelProcessingNode extends AbstractNode { private final TaskExecutor taskExecutor; // Spring管理的线程池 private final ServerDiagnosisTool diagnosisTool; // 假设的服务器诊断工具 public ParallelProcessingNode(TaskExecutor taskExecutor, ServerDiagnosisTool diagnosisTool) { this.taskExecutor taskExecutor; this.diagnosisTool diagnosisTool; } Override public MonoMessage execute(State state) { // 1. 从状态中获取需要检查的服务器列表 ListString serverIps (ListString) state.get(servers); // 2. 创建并行任务列表 ListMonoDiagnosisResult tasks serverIps.stream() .map(ip - Mono.fromCallable(() - diagnosisTool.pingServer(ip)) .subscribeOn(Schedulers.fromExecutor(taskExecutor)) // 指定在线程池执行 .onErrorReturn(new DiagnosisResult(ip, Error, Ping failed)) // 错误处理 ) .collect(Collectors.toList()); // 3. 等待所有任务完成 return Mono.zip(tasks, resultsArray - { ListDiagnosisResult results Arrays.stream(resultsArray) .map(obj - (DiagnosisResult) obj) .collect(Collectors.toList()); // 4. 将聚合结果放入新的消息中 return new Message(results, Map.of(contentType, diagnosis_results)); }); } }关键点解析使用了Spring的TaskExecutor来注入线程池避免手动管理线程。采用Project Reactor的Mono和Schedulers进行响应式编排这与Spring AI Graph的响应式接口MonoMessage天然契合。onErrorReturn确保了单个服务器检查失败不会导致整个并行节点失败而是返回一个标识错误的结果。最终使用Mono.zip等待所有并行任务完成并进行结果聚合。3.2 实现HITL审核节点与状态持久化接下来实现一个HumanApprovalNode。这需要与持久化状态存储配合。首先配置一个Redis状态存储以Spring Boot 3为例spring: ai: graph: state-store: type: redis redis: prefix: ai_graph_state:然后定义HITL节点Component public class HumanApprovalNode extends AbstractNode { private final StateStore stateStore; private final ApprovalNotificationService notificationService; public HumanApprovalNode(StateStore stateStore, ApprovalNotificationService notificationService) { this.stateStore stateStore; this.notificationService notificationService; } Override public MonoMessage execute(State state) { // 1. 生成一个唯一的审核ID String approvalId approval_ System.currentTimeMillis() _ ThreadLocalRandom.current().nextInt(); // 2. 将当前Graph状态持久化关联审核ID state.put(approvalId, approvalId); stateStore.put(approvalId, state); // 假设stateStore支持直接存储State对象 // 3. 构建审核请求内容从状态中提取AI建议的操作 String proposedAction (String) state.get(proposedAction); String requestContext String.format(AI建议执行操作%s。上下文%s, proposedAction, state.get(context)); // 4. 发送通知给人工审核者 notificationService.sendApprovalRequest(approvalId, requestContext); // 5. 返回一个特殊消息指示流程在此暂停等待外部事件驱动恢复 // 这里可以返回一个包含审核ID和等待状态的消息由外部的Supervisor或控制器监听后续结果 MapString, Object metadata Map.of( type, HUMAN_APPROVAL_PENDING, approvalId, approvalId, message, Waiting for human approval. ); return Mono.just(new Message(, metadata)); } }关键点解析StateStore是核心它负责将运行中的Graph状态包括所有变量、消息历史保存起来。这里假设其接口是put(key, state)。ApprovalNotificationService是一个自定义服务可以通过企业微信、钉钉、邮件等方式发送消息。该节点执行后Graph引擎会认为这个节点完成了但实际业务逻辑在此中断。我们需要一个外部机制如一个独立的REST控制器来监听人工决策并根据approvalId恢复对应的Graph执行。3.3 定义Supervisor的决策与恢复逻辑Supervisor需要监听人工审核的结果。我们创建一个REST端点来处理RestController RequestMapping(/api/approval) public class ApprovalController { private final Graph graph; // 你的Spring AI Graph实例 private final StateStore stateStore; private final TaskExecutor taskExecutor; PostMapping(/decide) public ResponseEntityString handleDecision(RequestBody ApprovalDecision decision) { // decision包含 approvalId, decision(APPROVE/REJECT), comment等字段 // 1. 从状态存储中恢复Graph状态 State savedState stateStore.get(decision.getApprovalId()); if (savedState null) { return ResponseEntity.status(404).body(Approval request not found or expired.); } // 2. 将人工决策注入到恢复的状态中 savedState.put(humanDecision, decision.getDecision()); savedState.put(humanComment, decision.getComment()); // 3. 异步恢复Graph的执行 Mono.fromRunnable(() - { try { // 这里需要根据决策决定恢复后跳转到哪个节点。 // 一种常见做法是在原Graph中设计分支根据humanDecision变量决定路径。 graph.execute(savedState).block(); // 恢复执行 } catch (Exception e) { // 记录恢复执行失败的日志 log.error(Failed to resume graph after approval, e); } }) .subscribeOn(Schedulers.fromExecutor(taskExecutor)) .subscribe(); // 4. 清理已使用的状态可选 stateStore.remove(decision.getApprovalId()); return ResponseEntity.ok(Decision received, resuming workflow.); } }关键点解析恢复执行的关键是获取到之前保存的State对象并用它重新调用graph.execute()。恢复执行是异步的避免阻塞HTTP请求线程。需要在原始的Graph定义中在HITL节点之后设计条件分支ConditionalNode或RouterNode来检查humanDecision变量从而走向“执行操作”或“取消操作”的不同分支。3.4 工具节点的线程安全与幂等性设计在并行环境下工具节点的设计必须格外小心。以上面的ServerDiagnosisTool为例Component public class ServerDiagnosisTool { // 使用HttpClient它本身是线程安全的 private final HttpClient httpClient HttpClient.newHttpClient(); public DiagnosisResult pingServer(String ip) { // 实现ping逻辑可能是HTTP请求或ICMP // 关键这个方法应该是无状态的或者使用的资源是线程安全的。 // 避免使用非线程安全的类如SimpleDateFormat。 try { HttpRequest request HttpRequest.newBuilder() .uri(URI.create(http:// ip :8080/health)) .timeout(Duration.ofSeconds(5)) .build(); HttpResponseString response httpClient.send(request, HttpResponse.BodyHandlers.ofString()); return new DiagnosisResult(ip, OK, response.body()); } catch (Exception e) { return new DiagnosisResult(ip, ERROR, e.getMessage()); } } }设计原则无状态工具类最好设计为无状态的所有数据通过参数传入。使用线程安全客户端如HttpClient、RedisTemplate配置正确的话、JdbcTemplate。避免共享可变数据绝对不要在工具类中使用静态的、可变的成员变量。幂等性工具方法执行多次应产生相同的结果这对于失败重试和并行执行的一致性至关重要。4. 完整工作流编排与配置实战现在我们把所有零件组装起来定义一个完整的Graph。我们将使用Spring AI的DSL领域特定语言或Bean方式来定义图。这里以编程式风格为例Configuration public class ServerOpsGraphConfig { Bean public Graph serverOpsGraph(ParallelProcessingNode parallelNode, HumanApprovalNode approvalNode, ExecuteActionNode actionNode, SendReportNode reportNode) { return new GraphBuilder() .start(start) .on(收到运维请求).to(analyzeRequest) .node(analyzeRequest) .action(state - { // 解析请求生成服务器列表和潜在操作 state.put(servers, List.of(192.168.1.101, 192.168.1.102)); state.put(potentialAction, restart_nginx); return Mono.just(new Message(分析完成)); }) .to(parallelCheck) .node(parallelCheck) .action(parallelNode) // 并行检查所有服务器 .to(evaluateResults) .node(evaluateResults) .action(state - { // 评估并行检查的结果 ListDiagnosisResult results (ListDiagnosisResult) state.get(diagnosis_results); boolean allHealthy results.stream().allMatch(r - OK.equals(r.status())); state.put(allHealthy, allHealthy); if (allHealthy restart_nginx.equals(state.get(potentialAction))) { state.put(proposedAction, 在健康服务器上重启Nginx); return Mono.just(new Message(建议执行重启, Map.of(next, requireApproval))); } else { return Mono.just(new Message(无需执行操作, Map.of(next, sendReport))); } }) .to(decisionRouter) .node(decisionRouter) .router(state - { String nextStep (String) state.getLastMessage().getMetadata().get(next); return nextStep; }) .route(requireApproval).to(humanApproval) .route(sendReport).to(generateReport) .node(humanApproval) .action(approvalNode) // 执行HITL流程在此暂停 .to(postApprovalRouter) // 审核完成后会恢复执行并跳转到这里 .node(postApprovalRouter) .router(state - { String decision (String) state.get(humanDecision); return APPROVE.equals(decision) ? executeAction : sendReport; }) .route(executeAction).to(performAction) .route(sendReport).to(generateReport) .node(performAction) .action(actionNode) // 执行实际的重启命令 .to(generateReport) .node(generateReport) .action(reportNode) // 生成并发送报告 .end() .build(); } }流程解读开始-分析请求确定目标服务器和可能操作。并行检查同时ping所有服务器收集状态。评估结果如果所有服务器健康且需要重启则提议操作并路由到requireApproval。人工审核HumanApprovalNode将状态持久化并发送通知流程暂停。外部决策人工通过/api/approval/decide接口提交决定。恢复与路由控制器恢复Graph执行根据humanDecision路由到执行动作或直接生成报告。执行动作/生成报告完成后续步骤。这个Graph清晰地展示了并行执行和HITL如何嵌入到一个线性的工作流中并通过路由节点实现了条件分支。5. 性能优化、错误处理与实战心得将并行和HITL引入生产环境会带来一系列新的挑战。下面分享一些踩坑后总结的经验。5.1 并行执行的性能调优与资源控制盲目并行会导致资源耗尽。以下配置和策略至关重要配置专用线程池不要在Graph中共享默认的TaskScheduler。Configuration public class ThreadPoolConfig { Bean(graphParallelExecutor) public TaskExecutor graphParallelExecutor() { ThreadPoolTaskExecutor executor new ThreadPoolTaskExecutor(); executor.setCorePoolSize(5); // 根据实际负载调整 executor.setMaxPoolSize(20); executor.setQueueCapacity(100); executor.setThreadNamePrefix(graph-parallel-); executor.initialize(); return executor; } }在ParallelProcessingNode中注入这个BeanQualifier(graphParallelExecutor) TaskExecutor taskExecutor。限制并发度对于调用外部API如OpenAI的子任务必须做限流。可以使用Resilience4j的Bulkhead或简单的Semaphore。private final Semaphore openAiApiSemaphore new Semaphore(5); // 最大同时5个请求 public MonoString callOpenAi(String prompt) { return Mono.fromCallable(() - { openAiApiSemaphore.acquire(); try { // 调用OpenAI API return openAiClient.chat(prompt); } finally { openAiApiSemaphore.release(); } }).subscribeOn(Schedulers.boundedElastic()); }设置超时每个并行任务必须有超时控制防止个别慢任务拖死整个并行节点。MonoDiagnosisResult task Mono.fromCallable(() - diagnosisTool.pingServer(ip)) .timeout(Duration.ofSeconds(10)) // 设置10秒超时 .onErrorResume(e - Mono.just(new DiagnosisResult(ip, TIMEOUT, e.getMessage()))) .subscribeOn(Schedulers.fromExecutor(taskExecutor));5.2 HITL流程的可靠性保障设计HITL流程中断了自动化的连续性其可靠性直接关系到用户体验和业务连续性。状态存储的选型与序列化选型Redis性能好适合临时状态数据库更持久适合长周期审核。根据审核超时时间选择。序列化确保你的State对象中的所有数据包括自定义对象都能被正确序列化和反序列化。使用JSON序列化时注意复杂对象的类型擦除问题。我推荐使用Jackson并注册所有可能的类型。审核请求的超时与清理不可能无限期等待人工审核。需要有一个后台任务定期扫描“等待中”的超时状态并执行默认操作如拒绝或通知升级。Scheduled(fixedDelay 300000) // 每5分钟运行一次 public void cleanupStaleApprovals() { // 从状态存储中查找超过30分钟的pending状态 // 执行默认决策并恢复对应的Graph走向失败或取消分支 // 清理状态 }幂等的恢复操作/api/approval/decide接口可能被重复调用如网络重试。需要确保基于同一个approvalId的恢复操作只执行一次。可以在状态中增加一个resumed标志位或者在恢复前使用分布式锁。5.3 调试与监控复杂Graph的实践技巧当Graph变得复杂尤其是加入了并行和中断后调试会变得困难。结构化日志与TraceId为每个Graph执行实例生成一个唯一的traceId并在所有相关日志包括并行子任务、HITL通知中打印它。这能让你在日志海洋中轻松串联起一次完整执行的所有步骤。在关键节点如并行节点开始/结束、HITL节点暂停记录结构化的JSON日志便于被ELK等系统收集和分析。可视化Graph执行状态虽然Spring AI Graph没有官方的UI但你可以通过暴露端点来查询StateStore中的状态。一个简单的做法是将Graph的结构节点和边以及当前执行实例的状态当前节点、变量通过API返回前端用D3.js或类似库画出一个动态的执行流程图。单元测试策略测试单个节点对ParallelProcessingNode、HumanApprovalNode等核心节点编写单元测试模拟输入状态验证输出和行为。集成测试片段使用GraphTestUtils如果Spring AI提供或自己模拟执行引擎测试包含几个节点的子图例如测试从“评估结果”到“人工审核”的路由逻辑。Mock外部依赖在测试中务必Mock掉AI模型调用、外部HTTP请求、邮件发送等让测试快速且稳定。5.4 从“能用”到“好用”的经验之谈HITL的交互设计比技术实现更重要发给人的审核请求信息必须清晰、 actionable。不要只扔一句“AI建议重启请审核”。而应该提供AI建议的操作、做出此建议的依据例如“因为检测到内存泄漏”、预计影响“重启将导致服务中断约10秒”、可供选择的选项“批准”、“拒绝”、“延迟到凌晨执行”。好的交互能极大提升人工处理效率和体验。并行任务的粒度要适中不要为每个微小操作都创建一个并行任务。线程创建、调度、结果聚合都有开销。通常将一批同质化的、耗时的、独立的数据处理单元进行并行化收益最大。例如处理100个文件每个文件处理需要1秒并行化后效果显著处理10个文件每个只需10毫秒并行化可能反而更慢。设计降级和熔断策略当并行调用的外部服务不稳定时要有后备方案。比如并行查询三个知识库如果一个超时是直接失败还是使用另外两个的结果继续在Graph中可以通过在并行节点后增加一个“结果融合与降级”节点来实现复杂的错误处理逻辑。状态变量的命名规范随着Graph变复杂状态中会有很多变量。建立命名规范如input_、temp_、output_前缀或者按模块划分能极大提升代码可读性和维护性。避免使用过于通用的键名如data、list。
返回列表