ARTICLE DETAIL

资讯详情

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

Spring AI Graph并行执行与HITL集成实战:从串行到并发工作流设计

Spring AI Graph并行执行与HITL集成实战:从串行到并发工作流设计 1. 从“单线程”到“多线程”为什么我们需要并行执行在上一篇文章里我们搭建了一个基础的Spring AI Graph并实现了一个简单的Supervisor监督者模式。那个模型就像一个项目经理把任务Task分配给不同的专家Agent然后等待他们一个个完成汇报。这很清晰也很有效但效率上有个明显的瓶颈所有任务都是串行执行的。想象一下你有一个项目需要市场分析、技术调研和竞品报告。如果让三位专家按顺序工作市场分析做完才启动技术调研那整个项目的周期就会被拉得很长。在真实的AI应用场景中这种串行带来的延迟是不可接受的。比如一个智能客服需要同时理解用户意图、查询知识库、生成回复并检查合规性这些子任务如果串行用户体验会大打折扣。这就是“并行执行”要解决的问题。它允许Graph中的多个节点Node同时被激活和执行只要它们之间的依赖关系允许。在Spring AI Graph的语境下这不仅仅是技术上的“多线程”更是一种工作流编排思维的升级。我们需要从设计上就思考哪些任务可以同时进行它们之间的数据流如何同步状态如何管理结合网络热词中频繁出现的snap graph builder和graph engineering这正是一个典型的“图工程”问题。我们不再仅仅是编写链式调用的代码而是在设计一个由异步任务节点构成的、有向无环的工作流图。Supervisor的角色也随之进化从一个简单的任务分发者变成一个并发的协调者与仲裁者它需要知道哪些分支可以并行并在所有分支完成后进行结果的汇聚与决策。而HITL的引入则让这个并发的系统变得更加复杂和强大。HITL即 Human-In-The-Loop人在回路意味着在AI自动执行的过程中在关键节点引入人工干预、审核或决策。在并行执行的场景下HITL点可能出现在多个并行的分支上如何优雅地暂停部分流程、等待人工输入、再继续执行同时不影响其他独立分支的运行这是对Graph状态管理和消息路由机制的严峻考验。所以本篇的目标很明确我们要将一个串行的Spring AI Graph改造为支持并行执行的系统并在此基础上无缝集成HITL机制构建一个既能高效自动化又能关键环节受控的智能体Agent系统。这不仅仅是代码的堆砌更是一次对spring ai 实现 自主agent架构的深度实践。2. 核心概念拆解状态、边与并发模型在动手之前我们必须把Spring AI Graph中几个支撑并行和HITL的核心概念吃透。很多人在使用spring ai graph时遇到的困惑都源于对这些底层模型的一知半解。2.1 状态State与状态存储State Store这是Graph的“记忆中枢”。一个Graph的执行过程本质上就是其状态对象在不同节点间流转和演化的过程。这个状态对象通常是一个Map或一个POJO包含了当前工作流的所有上下文信息。当你看到网络热词spring ai状态存储时它指的就是持久化这个状态对象的机制。为什么需要持久化设想一个长时间运行的、包含HITL的流程服务器可能会重启或者流程需要暂停几天等待人工审批。如果没有状态存储所有中间结果都会丢失。Spring AI提供了基于内存、Redis或数据库的StateStore实现确保工作流状态的可恢复性。在并行执行中状态管理变得更加微妙。多个节点可能同时读取和修改状态的不同部分。虽然Spring AI Graph的节点执行在理论上是线程安全的通常通过锁或乐观锁机制但在设计状态结构时我们应有意识地进行“领域划分”避免不同并行分支频繁读写同一数据域以减少冲突。例如将市场分析结果放在state.put(“marketData”, ...)技术调研结果放在state.put(“techData”, ...)。2.2 边Edge与路由逻辑边决定了工作流的走向。在串行模型中边很简单上一个节点执行完根据其输出结果通常是一个String类型的next动作决定去往哪个节点。在并行模型中边的类型变得丰富条件边Conditional Edge最常用根据某个条件决定下一个节点。固定边Fixed Edge无条件指向某个节点。多边Multi-Edge / 扇出Fan-out这是实现并行的关键。一个节点可以同时拥有多条符合条件的边指向多个不同的下游节点。当Graph执行到这个节点时它会同时激活所有符合条件的下游节点实现并行分支。汇聚边这不是一个明确的类型而是一种模式。通常通过一个“汇聚节点”来实现该节点等待所有并行分支的完成通过检查状态中各个分支的结果是否就绪然后进行结果合并。lang graph annotated[list[basemessage], add_messages]这个热词虽然指向LangGraph的特定注解但其思想是相通的它描述了节点如何消费和产出消息列表。在Spring AI Graph中节点通过操作State来传递信息理解“状态即消息”的概念至关重要。2.3 Spring AI Graph的并发模型Spring AI Graph默认构建在Project Reactor或CompletableFuture等异步编程模型之上。当你调用Graph.Builder构建图并执行时它内部已经是一个异步流。实现并行的核心在于Graph.Builder的addNode方法和边的配置。你无法直接命令两个节点“同时运行”而是通过图的拓扑结构来声明这种并行可能性。执行引擎如GraphExecutor会解析这个结构当遇到可以并行的路径时它会利用底层的异步任务执行器如TaskExecutor来并发执行。一个常见的误解是认为需要自己写多线程代码。实际上你只需要正确地定义节点和边并发由框架托管。你的主要工作从编写线程安全代码转变为设计线程安全的状态结构和无冲突的节点逻辑。3. 实战构建一个并行执行的智能调研助手让我们通过一个具体案例将理论付诸实践。我们要构建一个“智能调研助手”当用户提出一个产品创意时它能并行执行以下任务市场分析分析该产品的潜在市场规模和用户群体。技术可行性评估评估实现该产品所需的技术栈和难度。竞品分析寻找市场上已有的类似产品并分析其优劣。最后由一个“报告生成”节点汇总以上三个并行任务的结果形成一份完整的调研报告。3.1 定义状态与节点首先我们定义工作流的状态。这是一个承载所有数据的容器。import java.util.Map; import java.util.concurrent.ConcurrentHashMap; public class ResearchState { // 初始输入 private String productIdea; // 并行分支的输出 private MapString, String parallelResults new ConcurrentHashMap(); // 最终输出 private String finalReport; // 构造函数、getter、setter 省略... // 注意parallelResults 使用 ConcurrentHashMap 以支持并发写入 }接下来定义三个并行的工作节点。每个节点都是一个FunctionResearchState, ResearchState。import org.springframework.ai.chat.client.ChatClient; import org.springframework.stereotype.Component; import java.util.function.Function; Component public class MarketAnalysisNode implements FunctionResearchState, ResearchState { private final ChatClient chatClient; public MarketAnalysisNode(ChatClient chatClient) { this.chatClient chatClient; } Override public ResearchState apply(ResearchState state) { String idea state.getProductIdea(); // 调用大模型进行市场分析 String analysis chatClient.prompt() .user(u - u.text(请对以下产品创意进行市场分析包括潜在规模、目标用户 idea)) .call() .content(); // 将结果存入并行结果Map键为节点名称 state.getParallelResults().put(market_analysis, analysis); System.out.println(市场分析节点执行完成。); return state; } } // 技术评估节点 TechFeasibilityNode 和竞品分析节点 CompetitorAnalysisNode 结构类似。 // 它们分别调用AI并将结果以 tech_feasibility 和 competitor_analysis 为键存入 parallelResults。关键点这三个节点之间没有直接的依赖关系它们都只依赖于初始的productIdea并且各自向parallelResults这个共享状态的不同键写入数据。使用ConcurrentHashMap保证了并发写入的安全性。3.2 构建支持并行的Graph这是最核心的一步。我们需要使用Graph.Builder来声明这种并行关系。import org.springframework.ai.graph.Graph; import org.springframework.ai.graph.builder.GraphBuilder; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; Configuration public class ParallelResearchGraphConfig { Bean public GraphResearchState parallelResearchGraph( MarketAnalysisNode marketNode, TechFeasibilityNode techNode, CompetitorAnalysisNode competitorNode, ReportGenerationNode reportNode) { return GraphBuilder.ResearchStatebuilder() .addNode(start, state - state) // 起始节点什么都不做只是入口 .addNode(marketAnalysis, marketNode) .addNode(techFeasibility, techNode) .addNode(competitorAnalysis, competitorNode) .addNode(generateReport, reportNode) // 报告生成节点 .addEdge(start, marketAnalysis) // 从start同时指向三个分析节点 .addEdge(start, techFeasibility) .addEdge(start, competitorAnalysis) // 关键三个分析节点完成后都指向报告生成节点 .addEdge(marketAnalysis, generateReport) .addEdge(techFeasibility, generateReport) .addEdge(competitorAnalysis, generateReport) .build(); } }这段代码的精髓在于边的定义。我们创建了一个“扇出-扇入”结构扇出start节点同时连接到三个分析节点。当start节点执行完毕后Graph引擎会看到三条出边并且没有条件限制因此它会同时尝试激活并执行marketAnalysis,techFeasibility,competitorAnalysis三个节点。扇入三个分析节点各自完成后都有一条边指向generateReport节点。这里存在一个潜在问题generateReport节点会被触发三次吗不会。Spring AI Graph的引擎足够智能它通常会处理这种“多源汇聚”的情况确保下游节点在收到所有必要输入或满足特定条件后才执行一次。但更稳妥的做法是在generateReport节点内部做判断。3.3 实现汇聚节点报告生成报告生成节点需要等待所有并行任务完成然后汇总结果。Component public class ReportGenerationNode implements FunctionResearchState, ResearchState { private final ChatClient chatClient; public ReportGenerationNode(ChatClient chatClient) { this.chatClient chatClient; } Override public ResearchState apply(ResearchState state) { MapString, String results state.getParallelResults(); // 等待策略检查所有预期的结果是否都已就绪 if (!results.containsKey(market_analysis) || !results.containsKey(tech_feasibility) || !results.containsKey(competitor_analysis)) { // 如果还有结果未就绪可以选择返回当前状态等待再次被触发 // 更复杂的场景可以使用条件边或外部信号控制 System.out.println(报告生成节点尚有并行任务未完成等待中...); return state; // 返回未修改的状态引擎可能稍后重试取决于配置 } // 所有结果已就绪开始汇总 String prompt String.format( 请基于以下分析生成一份全面的产品调研报告\n 市场分析%s\n 技术评估%s\n 竞品分析%s\n 报告需要包含总结与建议。, results.get(market_analysis), results.get(tech_feasibility), results.get(competitor_analysis) ); String report chatClient.prompt().user(u - u.text(prompt)).call().content(); state.setFinalReport(report); System.out.println(最终报告生成完成); return state; } }注意上述“等待策略”是一种简单的实现。在生产环境中对于严格的“与”汇聚AND-join建议利用Graph的条件边特性或者使用CountDownLatch等同步机制在状态中标记任务完成让generateReport节点的入边上添加条件判断只有所有前置任务完成时才激活该节点。Spring AI Graph的未来版本可能会提供更原生的汇聚模式支持。3.4 执行与测试最后我们注入这个Graph并执行它。RestController public class ResearchController { private final GraphResearchState researchGraph; public ResearchController(GraphResearchState researchGraph) { this.researchGraph researchGraph; } GetMapping(/research) public String conductResearch(RequestParam String idea) { ResearchState initialState new ResearchState(); initialState.setProductIdea(idea); ResearchState finalState researchGraph.execute(initialState); return finalState.getFinalReport(); } }当你调用/research?idea一个AI编程助手时观察日志你应该能看到三个分析节点的日志输出顺序是随机的或几乎同时开始最后才是报告生成的日志。这证明了并行执行已经生效。4. 融入HITL在并行流程中插入“人工审批点”现在我们升级需求。假设“竞品分析”这个环节AI给出的结果需要经过人工审核确认后才能生效用于生成最终报告。其他两个分析任务可以继续并行。这就是一个典型的HITL场景。我们需要改造流程competitorAnalysis节点执行后不直接进入generateReport。而是进入一个humanApproval节点该节点将AI结果暂存并通知人工例如发送邮件、生成待办项。人工在外部系统如管理后台审核后触发一个回调接口批准或驳回结果。如果批准流程继续到generateReport如果驳回可能返回competitorAnalysis重试或者结束流程。4.1 扩展状态与新增节点首先扩展状态以容纳HITL所需的信息。public class ResearchStateWithHITL extends ResearchState { // HITL相关状态 private String hitlTaskId; // 人工任务ID用于关联 private String competitorAnalysisDraft; // AI生成的竞品分析草稿 private Boolean competitorAnalysisApproved; // 是否被批准 private String humanFeedback; // 人工反馈意见 // ... getters and setters }新增HumanApprovalNode。这个节点不直接调用AI而是“暂停”流程创建人工任务。Component public class HumanApprovalNode implements FunctionResearchStateWithHITL, ResearchStateWithHITL { private final TaskService taskService; // 假设有一个管理人工任务的服务 Override public ResearchStateWithHITL apply(ResearchStateWithHITL state) { // 1. 将AI的草稿保存到状态 state.setCompetitorAnalysisDraft(state.getParallelResults().get(competitor_analysis)); // 2. 创建一个待办任务关联到当前Graph执行实例可通过stateId标识 String taskId taskService.createTask( “审核竞品分析报告”, state.getCompetitorAnalysisDraft(), state.getGraphExecutionId() // 需要Graph执行上下文提供ID ); state.setHitlTaskId(taskId); // 3. 清除并行结果中的该条目因为尚未批准 state.getParallelResults().remove(competitor_analysis); System.out.println(已创建人工审核任务ID: taskId “流程暂停等待。”); return state; } }4.2 改造Graph引入条件路由现在需要重构Graph让competitorAnalysis之后的路由变得有条件。Bean public GraphResearchStateWithHITL researchGraphWithHITL( MarketAnalysisNode marketNode, TechFeasibilityNode techNode, CompetitorAnalysisNode competitorNode, HumanApprovalNode approvalNode, ReportGenerationNode reportNode) { return GraphBuilder.ResearchStateWithHITLbuilder() .addNode(start, state - state) .addNode(marketAnalysis, marketNode) .addNode(techFeasibility, techNode) .addNode(competitorAnalysis, competitorNode) .addNode(humanApproval, approvalNode) .addNode(generateReport, reportNode) // 并行扇出 .addEdge(start, marketAnalysis) .addEdge(start, techFeasibility) .addEdge(start, competitorAnalysis) // 市场和技术分析直接去报告节点 .addEdge(marketAnalysis, generateReport) .addEdge(techFeasibility, generateReport) // 竞品分析后先去人工审核 .addEdge(competitorAnalysis, humanApproval) // 关键从人工审核节点出来的条件边 .addEdge(humanApproval, generateReport, state - Boolean.TRUE.equals(state.getCompetitorAnalysisApproved())) .addEdge(humanApproval, competitorAnalysis, state - Boolean.FALSE.equals(state.getCompetitorAnalysisApproved())) .build(); }解读competitorAnalysis完成后强制进入humanApproval。humanApproval节点有两条出边都是条件边。第一条边指向generateReport条件是competitorAnalysisApproved true。第二条边指向competitorAnalysis重试条件是competitorAnalysisApproved false。4.3 实现外部回调接口人工在后台完成审核后需要调用一个API来更新状态并重新触发Graph。RestController public class HitlCallbackController { private final StateStoreResearchStateWithHITL stateStore; // 状态存储 private final GraphResearchStateWithHITL graph; private final TaskService taskService; PostMapping(/api/hitl/callback) public ResponseEntityString handleApproval(RequestBody ApprovalRequest request) { // 1. 根据任务ID找到对应的状态 String taskId request.getTaskId(); String stateId taskService.getStateIdByTaskId(taskId); // 需要建立任务与状态ID的映射 ResearchStateWithHITL state stateStore.get(stateId); if (state null) { return ResponseEntity.notFound().build(); } // 2. 更新状态中的人工审批结果 state.setCompetitorAnalysisApproved(request.isApproved()); state.setHumanFeedback(request.getFeedback()); if (request.isApproved()) { // 如果批准将之前保存的草稿正式放入并行结果集 state.getParallelResults().put(competitor_analysis, state.getCompetitorAnalysisDraft()); } // 3. 保存更新后的状态 stateStore.put(stateId, state); // 4. 重新执行Graph从当前状态继续。 // 注意这里需要获取到当初执行Graph的同一个实例或能根据stateId恢复执行上下文的方法。 // Spring AI Graph 可能需要通过特定的 Executor 来恢复执行。 graph.resume(stateId); // 假设有这样一个resume方法 return ResponseEntity.ok(处理成功流程已继续。); } }核心难点上述代码中的graph.resume(stateId)是一个理想化的接口。在Spring AI Graph的当前版本中实现从持久化状态中恢复并继续执行特定节点需要更精细地管理GraphExecution实例。一种可行的实践是将Graph和StateStore封装在一个服务里当回调触发时该服务根据stateId加载状态并从humanApproval节点的下一个环节开始重新执行Graph。这可能需要你手动管理执行路径或者利用Graph的“子图”或“持久化执行器”特性。4.4 HITL实战中的经验与坑经验一状态设计的幂等性HITL回调可能因为网络问题被重复调用。你的回调接口和节点逻辑必须是幂等的。例如在handleApproval中可以先检查state.getCompetitorAnalysisApproved()是否已设置避免重复处理。经验二超时与异常处理人工审核可能无限期等待。必须在流程中设置超时机制。可以在HumanApprovalNode中启动一个定时器或者在Graph外部设置一个监控进程对于超时未处理的任务执行默认操作如自动批准、驳回或通知升级。经验三上下文恢复当Graph从持久化状态恢复执行时确保所有Bean如ChatClient的上下文是正确的。例如如果使用了对话记忆ChatMemory需要确保记忆也能随状态一起恢复。坑parallelResults的并发访问在HITL场景下generateReport节点可能在其他分支完成后、竞品分析批准前就被触发因为它的入边条件可能已部分满足。因此generateReport节点中的等待逻辑必须足够健壮要能处理“结果尚未准备好”的情况并优雅地返回当前状态允许Graph继续等待其他边被触发。5. 调试与监控让并行与HITL流程可视化当流程变得复杂后调试是个大问题。git graph插件能可视化代码提交图我们也需要工具来可视化Spring AI Graph的执行流。方案一日志染色为每个Graph执行实例生成一个唯一IDexecutionId并在该实例所有节点的日志中打印此ID。这样在日志系统中可以通过这个ID过滤出一次完整执行的所有步骤尽管它们是并行打印的。Component public class MarketAnalysisNode implements FunctionResearchState, ResearchState { Override public ResearchState apply(ResearchState state) { String execId state.getExecutionId(); // 需要在状态中增加此字段 log.info([{}] 开始市场分析..., execId); // ... 业务逻辑 log.info([{}] 市场分析完成。, execId); return state; } }方案二状态快照存储利用StateStore不仅为了恢复也为了调试。定期或在每个节点执行后将状态存储起来。开发一个简单的管理界面可以查看任意executionId的历史状态变化就像看一个时间线能清楚地看到parallelResultsMap是如何被逐步填充的。方案三集成Micrometer与ActuatorSpring AI Graph可能与Spring Boot Actuator和Micrometer集成暴露执行指标。你可以监控每个节点的平均执行时间。并行分支的数量分布。HITL节点的平均等待时间。 这些指标对于优化Graph性能和发现瓶颈至关重要。关于memory graph debugger和graph rag等热词它们指向了更高级的调试和优化技术。Memory Graph Debugger可能指用于诊断内存中Graph状态流转的工具。Graph RAG则是将Graph与检索增强生成结合让AI节点能动态检索外部知识。在我们的上下文中可以想象一个节点专门负责从向量数据库检索信息其执行效率会直接影响整个并行流程的耗时需要重点监控。构建一个支持并行和HITL的Spring AI Graph系统是一个从“线性思维”到“拓扑思维”的转变。你不再只是关心一个接一个的步骤而是要设计一张清晰、健壮、可观测的网。这张网要能并发处理任务以提升效率也要能在关键点暂停以融入人的智慧。
返回列表