Skip to Content
全部文章全栈、LLM智能体框架LangGraphGo工作流学习

智能体框架LangGraphGo工作流学习

发布时间: 2026-06-29

前言

最近不是流行FSE和FDE了吗,补一下智能体开发相关的知识,本来打算去玩玩dify或者n8n,后来经过仔细的了解后发现也就是拖拖拖,然后开发一些工具给它调用,就觉得反正都是学习,就学点最直接的,学LangChain吧,然后又了解到LangChain+LangGraph开发 智能体的玩法,于是就通过代码的形式初步窥探了一下这两个框架,LangChain的作用是封装好各家ai模型的的调用,不用自己来按照大模型要求的格式去调用了。本次主要是想聊聊LangGraph,它的核心概念非常的简单,但是把智能体编排通过图的概念这么一包装,就变得简单许多了(可能主要是理解容易堵和成本下降了)。

LangGraph基础概念

LangGraph像画一个流程图一样来定义了编排逻辑,分为 状态节点 ,节点就是一个逻辑节点,比如到了某个步骤该调工具查库还是调用大模型。边的作用是定义每个节点执行后,流转到下一个节点应该是谁。状态则在节点之间起数据共享作用。其中节点+边共同构成了 ,这就是LangGraph的核心流程。

节点

这就是一个节点的定义,本质上来说就是一个函数/方法,执行一段逻辑

workflow.AddNode("node_name", "某某功能节点",func MyNode(ctx context.Context, state MyState) (MyState, error) { // 执行逻辑 return newState, nil })

这是边的定义

// 普通边,a节点执行完执行b节点 workflow.AddEdge("nodeA", "nodeB") // 并行执行a和b节点 workflow.AddEdge("start", "branch_a") workflow.AddEdge("start", "branch_b") // 条件边,比如如果状态是调用工具,则进入tool节点,否则进入结束节点 workflow.AddConditionalEdge("agent", func(ctx context.Context, state InventoryState) string { // 这里就用到了状态,状态是在整个图中共享数据的 if state.NextAction == "call_tool" { return "tools" } return graph.END })

基本上核心概念就这几个,做过图功能的人就很好理解node、edge的关系,后面我直接用代码来深入学习一下LangGrahpGo的一些具体应用场景。

实际应用场景

工具调用+条件节点

这个是比较常见的一个场景,就是用户提了一个问题,使用大模型来识别用户的意图,然后提供一些工具给大模型调用,然后将结果返回给用户,这个例子比如用户要查询iphone的库存

// 关键代码 // 导入 import ( "github.com/tmc/langchaingo/llms/openai" ) func main(){ // openAIOptions中有baseurl、apikey、模型型号,本次使用deepseek来做学习支持 llm, err := openai.New(openAIOptions...) // 使用Gin注册一个chat路由 r.POST("/api/chat", func(c *gin.Context) { ... initialState := InventoryState{UserQuery: input.Question} ctx := context.WithValue(c.Request.Context(), requestIDContextKey{}, requestID) // 关键代码,把用户输入的问题传入我们的LangGraph,函数见下一段源码 finalState, err := inventoryAgent.Invoke(ctx, initialState) if err != nil { c.JSON(http.StatusInternalServerError, gin.H{"error": fmt.Sprintf("智能体运行出错: %v", err)}) return } // 返回结果 c.JSON(http.StatusOK, gin.H{ "status": "success", "data": gin.H{ "answer": finalState.FinalAnswer, "tool_called": finalState.ToolResult, }, }) }) }
package main import ( "github.com/smallnest/langgraphgo/graph" "github.com/tmc/langchaingo/llms" ) // 模拟的本地库存数据库查询函数。 func queryStockDatabase(productName string) string { if strings.Contains(strings.ToLower(productName), "iphone") { return "库存充足,目前北京仓还剩 52 台。" } return "该商品暂无库存。" } // 图流程:agent -> tools -> agent -> END func buildInventoryAgent(llm interface { Call(context.Context, string, ...llms.CallOption) (string, error) }) *graph.StateRunnable[InventoryState] { workflow := graph.NewStateGraph[InventoryState]() // 添加一个决策节点 workflow.AddNode("agent", "AI决策节点", func(ctx context.Context, state InventoryState) (InventoryState, error) { requestID := requestIDFromContext(ctx) nodeStart := time.Now() var prompt string // 这里通过状态中的一个变量来模拟调用工具的结果,第一次进来时空的 if state.ToolResult == "" { prompt = fmt.Sprintf(`你是库存助手。用户的提问是: "%s"。如果用户询问库存信息,请只回复: "NEED_TOOL: [商品名称]"。否则直接回答用户。`, state.UserQuery) } else { prompt = fmt.Sprintf(`用户提问: "%s"。库存查询结果: "%s"。请用自然语言回答用户的问题。`, state.UserQuery, state.ToolResult) } response, err := llm.Call(ctx, prompt) // 如果模型恢复内容中有NEED_TOOL就表示大模型要调用工具查询库存,此处是模拟调用工具 if strings.HasPrefix(response, "NEED_TOOL:") { state.NextAction = "call_tool" state.FinalAnswer = strings.TrimSpace(strings.TrimPrefix(response, "NEED_TOOL:")) } else { state.NextAction = "end" state.FinalAnswer = response } return state, nil }) // 第二个节点,工具调用节点 workflow.AddNode("tools", "工具调用节点", func(ctx context.Context, state InventoryState) (InventoryState, error) { requestID := requestIDFromContext(ctx) start := time.Now() // 调用模拟数据查询 dbResult := queryStockDatabase(state.FinalAnswer) state.ToolResult = dbResult state.NextAction = "" return state, nil }) // 从agent开始执行 workflow.SetEntryPoint("agent") // 条件边,意思是agent执行完,下一个节点需要通过这个函数来决定 workflow.AddConditionalEdge("agent", func(ctx context.Context, state InventoryState) string { // 第一次agent执行完,会把NextAction设置为call_tool,就表示要执行工具查询节点 if state.NextAction == "call_tool" { return "tools" } return graph.END }) // 工具查询节点执行完回到agent决策节点 workflow.AddEdge("tools", "agent") return mustCompile(workflow) }

这个编排图就是:agent -> tools -> agent -> END,它是一个循环决策的图,agent执行完调用工具,工具调用后又回agent

对于这个示例如果给请求提交:iPhone 有库存吗?

最终ai会回复类似于: 北京仓库还有52台,您可以前往购买。 为什么说 类似 因为模型回复除非你约束他的格式,否则它的回复不是一成不变的。

并行聚合

并行聚合一般用于处理一堆并行任务,汇总后再返回的场景,来看一个从多个角度评审的智能体怎么做并行审查。

func main(){ parallelReviewExample := buildParallelReviewExample() // 还是加一个路由 r.POST("/api/examples/parallel-review", func(c *gin.Context) { var input struct { Content string `json:"content" binding:"required"` } if err := c.ShouldBindJSON(&input); err != nil { c.JSON(http.StatusBadRequest, gin.H{"error": "参数无效,请输入 content"}) return } start := time.Now() finalState, err := parallelReviewExample.Invoke(c.Request.Context(), ParallelReviewState{Content: input.Content}) if err != nil { c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) return } c.JSON(http.StatusOK, gin.H{ "status": "success", "elapsed_ms": time.Since(start).Milliseconds(), "data": finalState, }) }) } // 编排图: parallel_review -> summarize -> End func buildParallelReviewExample() *graph.StateRunnable[ParallelReviewState] { workflow := graph.NewStateGraph[ParallelReviewState]() // 每个 reviewer 都接收同一个 state,但只返回自己负责的局部结果。 reviewers := map[string]func(context.Context, ParallelReviewState) (ParallelReviewState, error){ "risk": func(ctx context.Context, state ParallelReviewState) (ParallelReviewState, error) { time.Sleep(80 * time.Millisecond) return ParallelReviewState{Content: state.Content, Reviews: map[string]string{"risk": "风险评审:关注承诺过满、缺少边界条件和异常兜底。"}}, nil }, "copy": func(ctx context.Context, state ParallelReviewState) (ParallelReviewState, error) { time.Sleep(80 * time.Millisecond) return ParallelReviewState{Content: state.Content, Reviews: map[string]string{"copy": "文案评审:表达清楚,但可以补一个更具体的用户收益。"}}, nil }, "tech": func(ctx context.Context, state ParallelReviewState) (ParallelReviewState, error) { time.Sleep(80 * time.Millisecond) return ParallelReviewState{Content: state.Content, Reviews: map[string]string{"tech": "技术评审:建议把外部 API、超时和重试作为独立节点。"}}, nil }, } workflow.AddParallelNodes("parallel_review", reviewers, func(results []ParallelReviewState) ParallelReviewState { merged := ParallelReviewState{ Reviews: map[string]string{}, } for _, result := range results { if merged.Content == "" { merged.Content = result.Content } for key, value := range result.Reviews { merged.Reviews[key] = value } } return merged }) // 聚合节点在并行结果之后执行,适合做总结、排序、冲突处理或最终 LLM 汇总。 workflow.AddNode("summarize", "聚合评审意见", func(ctx context.Context, state ParallelReviewState) (ParallelReviewState, error) { keys := make([]string, 0, len(state.Reviews)) for key := range state.Reviews { keys = append(keys, key) } sort.Strings(keys) parts := make([]string, 0, len(keys)) for _, key := range keys { parts = append(parts, state.Reviews[key]) } state.Summary = strings.Join(parts, " ") return state, nil }) workflow.SetEntryPoint("parallel_review") workflow.AddEdge("parallel_review", "summarize") workflow.AddEdge("summarize", graph.END) return mustCompile(workflow) }

对于这个示例如果给请求提交:我们准备上线一个库存智能体,请帮我做上线前评审

最终得到的响应结构体:

{ "data": { "content": "我们准备上线一个库存智能体,请帮我做上线前评审。", "reviews": { "copy": "文案评审:表达清楚,但可以补一个更具体的用户收益。", "risk": "风险评审:关注承诺过满、缺少边界条件和异常兜底。", "tech": "技术评审:建议把外部 API、超时和重试作为独立节点。" }, "summary": "文案评审:表达清楚,但可以补一个更具体的用户收益。 风险评审:关注承诺过满、缺少边界条件和异常兜底。 技术评审:建议把外部 API、超时和重试作为独立节点。" }, "elapsed_ms": 81, "status": "success" }

中断/恢复(人工审批)

这个场景非常常见,我们在执行一些特殊操作时,需要用户二次确认才继续执行。

r.POST("/api/examples/approval", func(c *gin.Context) { var input struct { OrderID string `json:"order_id"` Amount int `json:"amount"` State ApprovalState `json:"state"` Approved bool `json:"approved"` ResumeFrom []string `json:"resume_from"` } if err := c.ShouldBindJSON(&input); err != nil { c.JSON(http.StatusBadRequest, gin.H{"error": "参数无效"}) return } initialState := input.State if initialState.OrderID == "" { initialState.OrderID = input.OrderID } if initialState.Amount == 0 { initialState.Amount = input.Amount } var ( state ApprovalState err error ) // 如果前端传了ResumeFrom则表示要恢复,否则就是第一次进入会话 if len(input.ResumeFrom) == 0 { state, err = approvalExample.Invoke(c.Request.Context(), initialState) } else { // 第二次请求拿到中断时的状态,继续执行 state, err = approvalExample.InvokeWithConfig(c.Request.Context(), initialState, &graph.Config{ ResumeFrom: input.ResumeFrom, ResumeValue: input.Approved, }) } var interrupt *graph.GraphInterrupt if err != nil && !errors.As(err, &interrupt) { c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) return } if interrupt != nil { // 第一次执行时会产生一个中断错误,这时候我们把中断的错误信息返回给前端 if interruptedState, ok := interrupt.State.(ApprovalState); ok { state = interruptedState } c.JSON(http.StatusOK, gin.H{ "status": "interrupted", "interrupt_node": interrupt.Node, "interrupt_value": interrupt.InterruptValue, "resume_from": interrupt.NextNodes, "data": state, }) return } c.JSON(http.StatusOK, gin.H{"status": "success", "data": state}) }) // 第一次请求: // prepare -> approval(Interrupt) 返回 interrupted // 第二次请求: // 从 approval 节点恢复,Interrupt 会拿到 ResumeValue,然后继续 finish -> END func buildApprovalExample() *graph.StateRunnable[ApprovalState] { workflow := graph.NewStateGraph[ApprovalState]() workflow.AddNode("prepare", "准备审批单", func(ctx context.Context, state ApprovalState) (ApprovalState, error) { if state.OrderID == "" { state.OrderID = "demo-order" } if state.Amount == 0 { state.Amount = 100 } return state, nil }) workflow.AddNode("approval", "人工确认节点", func(ctx context.Context, state ApprovalState) (ApprovalState, error) { // Interrupt 会让图暂停,并把提示信息放进 GraphInterrupt.InterruptValue。 // 如果本次是 resume 进来的,Interrupt 不会再报错,而是返回 ResumeValue。 resumeValue, err := graph.Interrupt(ctx, fmt.Sprintf("请确认订单 %s 是否允许支付 %d 元", state.OrderID, state.Amount)) if err != nil { return state, err } // 拿到用户的选择的数据 approved, _ := resumeValue.(bool) state.Approved = approved return state, nil }) workflow.AddNode("finish", "完成审批", func(ctx context.Context, state ApprovalState) (ApprovalState, error) { // 如果用户是同意则审批完成 if state.Approved { state.Message = fmt.Sprintf("订单 %s 已通过审批,可以继续支付。", state.OrderID) } else { state.Message = fmt.Sprintf("订单 %s 未通过审批,流程已停止。", state.OrderID) } return state, nil }) workflow.SetEntryPoint("prepare") workflow.AddEdge("prepare", "approval") workflow.AddEdge("approval", "finish") workflow.AddEdge("finish", graph.END) return mustCompile(workflow) }

来测试一下,第一次请求参数:

{"order_id":"A100","amount":299}

响应

{ "data": { "order_id": "A100", "amount": 299, "approved": false, "message": "" }, "interrupt_node": "approval", "interrupt_value": "请确认订单 A100 是否允许支付 299 元", "resume_from": [ "approval" ], "status": "interrupted" }

此时前端判断status是interrupted,中断,询问用户是否允许支付,用户选择允许后提交:

{"state":{"order_id":"A100","amount":299},"approved":true,"resume_from":["approval"]}

其中的resume_from透传第一次响应中resume_from,得到响应:

{ "data": { "order_id": "A100", "amount": 299, "approved": true, "message": "订单 A100 已通过审批,可以继续支付。" }, "status": "success" }

此时表示完成订单,这只是一个示例,实际上的LangGrayph不会把一些数据透传到前端,通常会把resume_from之类的状态存储到redis、mysql等,给前端签发一个id,等前端二次提交带上id,后端通过id去找对应的数据,避免前端传入不被允许的数据。

会话记忆

我们在调用大模型时,实际上每次调用都是一个新的会话,可能存在需要保留会话上下文的情况,本质上就是数据的持久化,LangGraph有一个checkpoint的概念,可以把数据存储到内存、redis、mysql、sqlite等。

比如上面的中断/恢复

package main import ( "context" "fmt" "log" "os" "github.com/tmc/langgraph/checkpointer/postgres" "github.com/tmc/langgraph/graph" "github.com/tmc/langgraph/model" "github.com/tmc/langgraph/types" ) // State 最简状态 type State struct { Messages []model.Message } func (s State) Copy() graph.State { ms := make([]model.Message, len(s.Messages)) copy(ms, s.Messages) return State{Messages: ms} } // 中断节点 func collectNode(ctx context.Context, s State) (State, error) { phone, err := types.Interrupt[string](ctx, "输入手机号") if err != nil { return s, err } s.Messages = append(s.Messages, model.NewHumanMessage("手机号:"+phone)) return s, nil } // 收尾节点 func chatNode(ctx context.Context, s State) (State, error) { s.Messages = append(s.Messages, model.NewAssistantMessage("信息已保存")) return s, nil } func main() { ctx := context.Background() // PG 检查点 cp, _ := postgres.New("postgres://postgres:123456@127.0.0.1:5432/langgraph_db?sslmode=disable") defer cp.Close() _ = cp.Setup(ctx) // 构建图 b := graph.NewStateGraph[State]() b.AddNode("collect", collectNode) b.AddNode("chat", chatNode) b.AddEdge(graph.StartNode, "collect") b.AddEdge("collect", "chat") b.AddEdge("chat", graph.EndNode) g := b.Compile(graph.WithCheckpointer(cp)) // 这里测试直接固定会话ID cfg := graph.Config{ThreadID: "t1"} // 执行 out, err := g.Invoke(ctx, State{Messages: []model.Message{model.NewHumanMessage("登记信息")}}, cfg) if e, ok := err.(types.InterruptError); ok { fmt.Println("中断提示:", e.Prompt) os.Exit(0) } if err != nil { log.Fatal(err) } // 打印结果 for _, m := range out.Messages { fmt.Printf("[%s] %s\n", m.Role, m.Content) } }

这个例子和上面的例子相比最大的区别就是,这个例子使用了pgsql来存储状态,自动建三张表langgraph postgres checkpointer,全部靠 thread_id 关联,流程走到 types.Interrupt,执行两个持久化动作: 把当前所有 State、中断标记、执行位置序列化存入 checkpoints。

最后编辑于

hi