ARTICLE DETAIL

资讯详情

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

基于LangGraph与LangChain的workflow低代码平台架构与实战

基于LangGraph与LangChain的workflow低代码平台架构与实战 简介这是一套面向开发者与AI应用团队的低代码工作流平台源码用于快速搭建聊天机器人、RAG、Agent及Multi-Agent应用技术栈涵盖LangGraph、Langchain、FastAPI与NextJS适合具备一定Python与前端基础、希望缩短AI应用落地周期的中高级开发者。资源包共566个文件约2.12MB以219个ts与146个tsx构成前端交互与工作流编排界面111个py承载后端服务与Agent逻辑另有json、yml、dockerfile等配置部署文件及少量md说明文档整体结构清晰、模块划分明确。该版本新增MCP Node支持将MCP工具转换为LangChain工具供LangGraph Agent调用可连接多个MCP服务器并动态加载工具兼容stdio与SSE两种传输模式与现有工作流无缝集成。目前已有144人学习读者可借此理解低代码工作流编排、RAG与多智能体协作的实现思路并参考其目录组织与配置方式快速二次开发。1. 从一堆散装脚本到可视化编排workflow 低代码平台到底解决什么问题如果你做过三个以上的 LangChain 项目大概率经历过这种场面chain.py里塞了十几个Runnable改一个 prompt 要翻三遍文件产品经理问「能不能把检索那步换成别的向量库」你心里清楚能改但改完要重新跑一遍全链路回归。更别提多 Agent 场景A 调 B、B 调 C、C 又回调 A光靠读代码已经理不清执行顺序了。基于 workflow 工作流的低代码平台本质就是把这种「代码里硬编码的编排逻辑」抽出来变成一张可拖拽、可保存、可版本化的有向图。节点是 LLM 调用、检索器、工具、条件分支、循环体边是数据流向。后端用 LangGraph 做图执行引擎LangChain 提供模型和工具的标准接口FastAPI 暴露 HTTP 与流式接口NextJS 负责画布和调试面板。这套组合在 2024 年下半年之后逐渐成为主流选型不是因为它最优雅而是因为每一层都有成熟的替代方案出问题能换。适合谁一是手里有多个 RAG 或 Agent 需求、但不想每个都从零写编排的团队二是想把已有 LangChain 脚本快速产品化、又不想重写成另一套 DSL 的工程师。不适合谁如果你的业务只有一条固定链路、半年不改一次直接写 Python 更省事上低代码平台反而是负担。2. 拆解 workflow 编排的四层架构LangGraph 做引擎、LangChain 做接口、FastAPI 做通道、NextJS 做画布2.1 为什么用 LangGraph 而不是自己写 DAG 调度器自己写调度器听起来可控但很快会遇到三个硬骨头状态如何在节点间传递、条件边怎么表达、中断后怎么恢复。LangGraph 把这三件事都做成了原语。它的核心抽象是StateGraph你定义一个状态类型通常是TypedDict或 Pydantic 模型每个节点是一个接收状态、返回状态增量的函数边可以是普通边也可以是条件边。from typing import TypedDict, Annotated from langgraph.graph import StateGraph, END from langgraph.graph.message import add_messages class WorkflowState(TypedDict): # add_messages 是 reducer多个节点写 messages 时会自动合并而不是覆盖 messages: Annotated[list, add_messages] retrieved_docs: list next_action: str def retrieve_node(state: WorkflowState): # 实际项目里这里调向量库示例用占位 docs [doc_a, doc_b] return {retrieved_docs: docs} def generate_node(state: WorkflowState): # 调 LLM把检索结果拼进 prompt answer 基于检索结果的回答 return {messages: [{role: assistant, content: answer}]} def route_after_retrieve(state: WorkflowState) - str: # 条件边检索为空就走兜底否则走生成 return generate if state[retrieved_docs] else fallback builder StateGraph(WorkflowState) builder.add_node(retrieve, retrieve_node) builder.add_node(generate, generate_node) builder.add_node(fallback, lambda s: {messages: [{role: assistant, content: 未找到相关资料}]}) builder.set_entry_point(retrieve) builder.add_conditional_edges(retrieve, route_after_retrieve, {generate: generate, fallback: fallback}) builder.add_edge(generate, END) builder.add_edge(fallback, END) graph builder.compile()这段代码的关键在Annotated[list, add_messages]。很多人第一次用 LangGraph 会踩的坑是节点返回的字典直接覆盖了整个状态导致前面节点写入的字段丢失。reducer 就是解决这个问题的——它告诉图「这个字段是追加还是替换」。add_messages是内置的追加 reducer自定义字段可以写自己的 reducer 函数。参数上set_entry_point决定从哪个节点开始add_conditional_edges的第三个参数是路由映射表键是路由函数的返回值值是目标节点名。路由函数只读状态、不做副作用这是硬性约定否则并发执行时行为不可预测。2.2 LangChain 在这一层扮演什么角色别把它当框架当接口层LangChain 经常被吐槽抽象太厚但在低代码平台里它的价值恰恰是「统一接口」。你的平台要支持 OpenAI、Anthropic、本地 Ollama、通义千问如果每个都写适配代码维护成本会爆炸。LangChain 的BaseChatModel和BaseRetriever把这些差异抹平了平台只需要面向接口编程。from langchain_openai import ChatOpenAI from langchain_community.chat_models import ChatOllama from langchain_core.prompts import ChatPromptTemplate def build_llm(provider: str, model: str, **kwargs): # 平台配置里读 provider运行时决定实例化哪个模型 if provider openai: return ChatOpenAI(modelmodel, temperaturekwargs.get(temperature, 0.7)) elif provider ollama: # 本地部署场景base_url 指向本机 ollama 服务 return ChatOllama(modelmodel, base_urlkwargs.get(base_url, http://localhost:11434)) raise ValueError(funsupported provider: {provider}) prompt ChatPromptTemplate.from_messages([ (system, 你是一个严谨的助手只根据提供的资料回答。), (human, 资料{context}\n\n问题{question}) ]) chain prompt | build_llm(openai, gpt-4o-mini) result chain.invoke({context: 示例资料, question: 示例问题})这里有个选型细节值得说LangChain 的 LCEL 管道语法|操作符在简单链路上很好用但一旦涉及循环、分支、人工中断就应该交给 LangGraph不要试图用 LCEL 硬凑。我见过有人用RunnableBranch嵌套五层来实现多轮判断最后调试时连自己都看不懂。判断标准很简单如果你的链路需要「根据上一步结果决定下一步走哪」就用图不要用链。2.3 FastAPI 暴露工作流流式接口和会话状态怎么设计低代码平台的执行接口有两个硬需求一是流式输出用户要看到 token 一个个蹦出来二是会话隔离不同用户的对话状态不能串。FastAPI 的StreamingResponse配合 LangGraph 的astream_events能同时满足。from fastapi import FastAPI from fastapi.responses import StreamingResponse from pydantic import BaseModel import json app FastAPI() class RunRequest(BaseModel): workflow_id: str session_id: str input: dict app.post(/workflow/run) async def run_workflow(req: RunRequest): async def event_generator(): # config 里的 thread_id 是 LangGraph 做会话隔离的关键 config {configurable: {thread_id: req.session_id}} async for event in graph.astream_events(req.input, configconfig, versionv2): if event[event] on_chat_model_stream: chunk event[data][chunk].content if chunk: # SSE 格式前端 EventSource 直接消费 yield fdata: {json.dumps({token: chunk})}\n\n yield data: [DONE]\n\n return StreamingResponse(event_generator(), media_typetext/event-stream)thread_id是 LangGraph 检查点机制的核心参数。如果你用MemorySaver做检查点同一个thread_id的多次调用会共享状态这正是多轮对话需要的。但要注意MemorySaver存在进程内存里多 worker 部署时会丢状态生产环境要换成PostgresSaver或RedisSaver。这个坑我在压测时踩过——单 worker 跑得好好的一上uvicorn --workers 4会话就随机丢失。2.4 NextJS 画布与后端的数据契约前端画布最终要序列化成后端能执行的图定义。常见做法是前端维护一份 JSON schema节点数组加边数组保存时 POST 给 FastAPI后端把这份 JSON 翻译成StateGraph的构建调用。{ nodes: [ {id: n1, type: retriever, config: {top_k: 5, index: docs_v2}}, {id: n2, type: llm, config: {provider: openai, model: gpt-4o-mini}}, {id: n3, type: condition, config: {expression: len(docs) 0}} ], edges: [ {from: n1, to: n3}, {from: n3, to: n2, when: true}, {from: n3, to: end, when: false} ] }后端翻译层要做的事把type映射到节点构造函数把config透传成节点参数把edges里的when转成条件边路由。这里最容易出问题的是节点 ID 和 LangGraph 节点名的对应关系——前端允许用户重命名节点但重命名后边引用要同步更新否则图编译时报「节点不存在」。我的做法是前端内部始终用不可变的 UUID 做引用显示名单独存一个字段这样重命名不影响拓扑。3. 从零跑通一个最小 RAG workflow环境、代码、调试三步走3.1 环境准备与依赖版本锁定这套技术栈的依赖冲突是高频翻车点。LangChain 生态包更新极快langchain-core、langchain-community、langgraph三者版本必须匹配。我一般用uv或poetry锁版本不用裸pip install。# 用 uv 创建虚拟环境并锁定版本 uv venv .venv source .venv/bin/activate # Windows 用 .venv\Scripts\activate uv pip install \ fastapi0.115 \ uvicorn[standard]0.30 \ langgraph0.2.0 \ langchain-core0.3.0 \ langchain-openai0.2.0 \ langchain-community0.3.0 \ pydantic2.0版本号只写主次版本不写死到 patch因为 patch 版本更新频繁且通常兼容。但langgraph和langchain-core的大版本要一起升跨大版本升级时先跑一遍官方迁移指南里的破坏性变更清单。Windows 上如果遇到uvloop装不上那是正常的——uvloop不支持 Windowsuvicorn[standard]会自动降级到asyncio不影响功能。3.2 最小可运行 RAG 图的完整代码下面是一个能跑通的最小 RAG workflow包含检索、生成、条件兜底三个节点。import os from typing import TypedDict, Annotated from langgraph.graph import StateGraph, END from langgraph.graph.message import add_messages from langgraph.checkpoint.memory import MemorySaver from langchain_openai import ChatOpenAI, OpenAIEmbeddings from langchain_core.vectorstores import InMemoryVectorStore from langchain_core.documents import Document # 1. 准备向量库示例用内存版生产换 PGVector 或 Milvus embeddings OpenAIEmbeddings(modeltext-embedding-3-small) docs [ Document(page_contentLangGraph 是 LangChain 团队推出的图编排引擎。), Document(page_contentFastAPI 是基于 Starlette 的异步 Web 框架。), ] vectorstore InMemoryVectorStore.from_documents(docs, embeddings) retriever vectorstore.as_retriever(search_kwargs{k: 2}) llm ChatOpenAI(modelgpt-4o-mini, temperature0) class RAGState(TypedDict): question: str docs: list answer: str messages: Annotated[list, add_messages] def retrieve(state: RAGState): hits retriever.invoke(state[question]) return {docs: [d.page_content for d in hits]} def generate(state: RAGState): context \n.join(state[docs]) resp llm.invoke(f根据资料回答{context}\n\n问题{state[question]}) return {answer: resp.content, messages: [resp]} def fallback(state: RAGState): return {answer: 知识库中没有相关内容。, messages: []} def has_docs(state: RAGState) - str: return generate if state[docs] else fallback builder StateGraph(RAGState) builder.add_node(retrieve, retrieve) builder.add_node(generate, generate) builder.add_node(fallback, fallback) builder.set_entry_point(retrieve) builder.add_conditional_edges(retrieve, has_docs, {generate: generate, fallback: fallback}) builder.add_edge(generate, END) builder.add_edge(fallback, END) # checkpointer 让同一 thread_id 的对话保留历史 graph builder.compile(checkpointerMemorySaver()) if __name__ __main__: config {configurable: {thread_id: demo-1}} result graph.invoke({question: LangGraph 是什么}, configconfig) print(result[answer])逻辑说明retrieve节点只负责取文档不调 LLMgenerate节点把文档拼进 prompthas_docs是纯路由函数不产生副作用。参数上search_kwargs{k: 2}控制召回数量RAG 场景一般 3 到 5 比较稳太大反而引入噪声。temperature0是 RAG 的默认选择需要创意生成时才调高。3.3 用 LangGraph Studio 或日志定位执行卡点图跑不通时最有效的调试手段是看每个节点的输入输出。LangGraph 支持astream_events逐事件输出也可以接 LangSmith 做可视化追踪。本地开发我一般先加一段打印。async def debug_run(question: str): config {configurable: {thread_id: debug-1}} async for event in graph.astream_events({question: question}, configconfig, versionv2): kind event[event] name event.get(name, ) if kind on_chain_start: print(f[进入节点] {name}) elif kind on_chain_end: print(f[离开节点] {name} 输出: {event[data].get(output)}) elif kind on_chat_model_stream: print(event[data][chunk].content, end, flushTrue)如果发现某个节点一直不执行先检查条件边的路由函数返回值是否和映射表的键完全一致——字符串大小写、空格都会导致路由失败而 LangGraph 默认不会报错只会静默走到END。这个「静默失败」是新手最容易卡住的地方加日志能立刻定位。4. 多 Agent 与 RAG 混合编排的避坑清单五个真实踩坑记录4.1 状态字段被覆盖导致检索结果丢失现象检索节点明明返回了文档生成节点拿到的docs却是空列表。原因状态类型里docs字段没有配 reducerLangGraph 默认行为是「后写入的覆盖先写入的」。如果中间有节点返回了不含docs的字典某些版本下会触发整体状态替换。解决给需要累积的字段加Annotated[list, operator.add]或自定义 reducer。只读字段如question不需要 reducer但任何会被多个节点写入的字段都必须显式声明合并策略。4.2 多 worker 部署后会话状态随机丢失现象本地单进程跑多轮对话正常部署到uvicorn --workers 4后第二轮对话经常「失忆」。原因MemorySaver把检查点存在进程内存里请求被负载均衡到另一个 worker 时读不到之前的检查点。解决生产环境换持久化 checkpointer。Postgres 方案用langgraph-checkpoint-postgresRedis 方案用社区维护的 saver。切换后thread_id语义不变但要注意数据库连接池大小要匹配 worker 数量否则高并发下会排队等连接。4.3 流式接口在 Nginx 后面被缓冲现象本地curl能看到 token 逐个输出部署到 Nginx 反代后变成一次性返回。原因Nginx 默认开启proxy_buffering会把上游的流式响应攒起来再发。解决在 location 块里加proxy_buffering off;和proxy_cache off;同时确认proxy_set_header Connection ;配合proxy_http_version 1.1。FastAPI 侧StreamingResponse的media_type必须是text/event-stream少了这个头某些客户端不会按流处理。4.4 条件边路由函数抛异常导致整图静默结束现象某个分支该走却没走图直接结束没有任何报错。原因路由函数内部访问了不存在的状态字段抛了KeyError但 LangGraph 在某些版本下会吞掉路由函数的异常。解决路由函数里所有字段访问用.get()加默认值并在函数入口加 try/except 打日志。更稳妥的做法是路由逻辑尽量简单复杂判断放到前置节点里算好路由函数只读一个布尔字段。4.5 前端画布保存的图定义与后端节点注册不一致现象画布上连线正常保存后执行报「节点 n3 未注册」。原因前端允许用户删除节点但没清理关联边或者节点类型下拉框新增了类型但后端映射表没同步。解决保存前在前端做一次图校验孤立边、悬空节点、环检测后端翻译层对未知type直接抛 400 并返回具体节点 ID。环检测尤其重要——LangGraph 支持循环但无出口的环会让执行永不结束必须在前端就拦住。5. 把 workflow 平台做扎实的两个进阶技巧检查点持久化与节点级超时5.1 用 PostgresSaver 替换 MemorySaver 的完整迁移步骤会话持久化是低代码平台从 demo 走向可用的分水岭。迁移本身不复杂但有几个参数决定成败。from langgraph.checkpoint.postgres import PostgresSaver from psycopg_pool import ConnectionPool # 连接池大小建议 worker 数 × 每 worker 并发请求数 × 1.5 pool ConnectionPool( conninfopostgresql://user:passlocalhost:5432/workflow, min_size2, max_size20, kwargs{autocommit: True, prepare_threshold: 0} ) # 首次运行需要建表之后注释掉 setup with pool.connection() as conn: checkpointer PostgresSaver(conn) checkpointer.setup() graph builder.compile(checkpointercheckpointer)prepare_threshold0这个参数容易被忽略。psycopg3 默认会做 prepared statement 缓存但在 PgBouncer 的 transaction 模式下会导致「prepared statement already exists」错误。设成 0 关闭缓存代价是轻微的性能损失换来的是连接池兼容性。表结构由setup()自动创建包含checkpoints和checkpoint_writes两张表前者存状态快照后者存待写入的增量。清理历史检查点用checkpointer.delete_thread(thread_id)别直接DELETE FROM否则可能留下孤儿写入记录。5.2 给每个节点加超时避免一个慢检索拖垮整条链路LangGraph 本身没有内置节点级超时但可以用asyncio.wait_for包一层。这个技巧在多 Agent 场景尤其重要——某个 Agent 调外部 API 卡住整条图会一直挂着。import asyncio from functools import wraps def with_timeout(seconds: float): def decorator(fn): wraps(fn) async def wrapper(state, *args, **kwargs): try: return await asyncio.wait_for(fn(state, *args, **kwargs), timeoutseconds) except asyncio.TimeoutError: # 超时后返回降级结果让图继续走而不是整体失败 return {error: fnode {fn.__name__} timeout after {seconds}s} return wrapper return decorator with_timeout(5.0) async def retrieve_node(state): # 异步检索超过 5 秒走降级 hits await retriever.ainvoke(state[question]) return {docs: [d.page_content for d in hits]}超时时间怎么定我的经验是检索节点 3 到 5 秒LLM 生成节点 30 到 60 秒取决于模型和输出长度工具调用节点按外部 API 的 P99 延迟乘 2。降级策略要提前想好——检索超时可以返回空列表走兜底分支生成超时可以返回「服务繁忙请重试」工具超时可以标记该工具不可用并继续。关键是不要让超时变成异常往上抛否则整图中断用户体验比降级还差。最后说个我自己的习惯每次给图加新节点先在本地用astream_events跑一遍把每个节点的输入输出打印出来存成日志文件。上线后如果用户报「结果不对」第一件事就是拿thread_id去查检查点表看状态在哪一步偏了。这个习惯帮我省了无数次「靠猜」的时间。希望帮到你。本文还有配套的精品资源点击获取
返回列表