智能体框架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。
