一、什么是工作流
工作流的本质是:
用一张图组织多个任务,并明确任务之间的依赖关系、数据传递、条件分支、并行、汇合、循环和终止规则。
比如一个最基本的 RAG 工作流:
开始 ↓ 搜索知识库 ↓ 文档重排 ↓ 大模型生成答案 ↓ 回复用户转换成图结构就是:
Start → Search → Reranker → AI Chat → Reply其中:
- 每一个方框是一个节点
Node。 - 每一条箭头是一个边
Edge。 - 节点执行具体任务。
- 边决定执行顺序。
1、工作流有什么用
1. 把复杂任务拆成小任务
例如 RAG 不再是一个巨大的函数,而是:
问题处理 → 知识库检索 → 文档过滤 → 文档重排 → 提示词构造 → LLM 生成 → 答案输出每一步都能单独观察、测试和替换。
2. 固化业务顺序
如果企业规定:
必须先查询客户 → 再查询订单 → 再检查退款条件 → 最后才能创建退款工作流可以保证这个顺序。
3. 支持条件分支
例如:
查询订单 ↓ 订单是否存在? ├── 是 → 检查退款资格 └── 否 → 返回“订单不存在”4. 支持并行执行
例如:
用户问题 ├── 查询客户信息 ├── 查询订单信息 └── 查询物流信息 ↓ 汇总结果5. 支持观察和排错
工作流可以记录:
- 哪个节点开始执行。
- 哪个节点执行失败。
- 每个节点花了多长时间。
- 节点输入是什么。
- 节点输出是什么。
- 最终经过了哪个分支。
这对企业级 AI 应用很重要,因为 LLM 本身具有不确定性。
6. 用确定性流程约束概率性模型
LLM 的回答带有概率性,而工作流可以规定:
必须先检索 必须经过审核 必须符合条件 必须记录日志 必须得到人工批准因此:
工作流不是为了让 AI 更自由,而是为了让 AI 的行为更可控。
工作流(LangGraph):带状态的有向图,支持:
- 顺序执行
- 条件分支(判断检索资料够不够,走不同逻辑)
- 循环回流(资料质量差,重新改写 query 再检索)
- 全局共享数据(所有步骤共用一份状态)
- 断点保存、人工介入审核
2、核心组成三要素(LangGraph 标准工作流)
- State 全局状态:整个流程共享的数据容器,所有节点读写它;比如用户问题、检索文档、LLM 回答、重试次数、对话历史。
- Node 节点:流程里最小执行单元,对应一个业务动作(检索文档、打分、生成回答、改写查询词)。
- Edge 边:控制节点流转
- 普通边:固定跳转 A→B
- 条件边:根据 State 数据判断下一步去哪里(循环、分支核心)
3、 工作流能解决 RAG 什么痛点
- 检索到的文档相关性很低,直接丢给 LLM 会产生幻觉 → 工作流可以打分,不合格就重新检索
- 单次 query 检索覆盖不全 → 自动生成多条同义 query,多次检索融合结果
- 多轮对话上下文散乱 → 全局 State 统一存放对话历史
- 需要人工校验 AI 答案再输出 → 流程暂停,人工修改后继续执行
二、简单项目
1、前置依赖安装
pip install langgraph langchain-openai langchain-community langchain-core pydantic python-dotenv入门级工作流项目:自反思 RAG 工作流(完整可运行)
2、业务流程设计(工作流流转逻辑)
- START 接收用户问题
- 节点 1:向量检索知识库,拿到文档片段
- 节点 2:LLM 打分,判断文档是否和问题相关
- 分支 1:文档不相关(分数低)→ 改写用户 query,回到检索节点重新查(循环)
- 分支 2:文档相关 → 进入生成节点
- 节点 3:结合检索文档生成最终回答
- END 输出答案
流程图文字版: START → 检索节点 → 打分节点 打分不通过 → 改写 query → 检索节点(循环) 打分通过 → 生成回答 → END
步骤 1:环境文件 .env
OPENAI_API_KEY=sk-xxx EMBED_MODEL=BAAI/bge-small-zh-v1.5步骤 2:完整代码分步讲解
2.1 导入依赖、初始化全局模型、向量库
from typing import TypedDict, List from dotenv import load_dotenv import os # LangGraph核心 from langgraph.graph import StateGraph, START, END from langgraph.checkpoint.memory import MemorySaver # LangChain组件 from langchain_openai import ChatOpenAI, OpenAIEmbeddings from langchain_community.vectorstores import Chroma from langchain_core.documents import Document from langchain_core.prompts import ChatPromptTemplate # 加载环境变量 load_dotenv() llm = ChatOpenAI(model="gpt-3.5-turbo", api_key=os.getenv("OPENAI_API_KEY"), temperature=0) # 模拟本地向量库(你可替换Qdrant) # 测试知识库文档 test_docs = [ Document(page_content="FastAPI是Python高性能异步Web框架,常用于RAG后端开发"), Document(page_content="LangGraph用于构建带循环、分支的LLM工作流,解决普通RAG无法反思优化的问题"), Document(page_content="Qdrant是生产级向量数据库,支持异步客户端适配FastAPI") ] vector_store = Chroma.from_documents(test_docs, OpenAIEmbeddings()) retriever = vector_store.as_retriever(search_kwargs={"k": 2})2.2 定义工作流全局 State(核心:共享状态)
TypedDict 定义所有流程需要共用的数据字段
class RAGWorkflowState(TypedDict): question: str # 用户原始问题 rewrite_query: str # 改写后的查询词(循环检索用) docs: List[Document] # 检索到的知识库文档 answer: str # 最终AI回答 retry_times: int # 检索重试次数,防止无限死循环字段说明: 整个工作流所有节点都能读取 / 修改这 5 个字段,数据全局互通。
2.3 定义所有 Node 节点(每一个节点 = 独立业务步骤)
节点 1:检索文档
输入:状态里的改写 query / 原始问题
输出:更新 state 的 docs 字段
def retrieve_node(state: RAGWorkflowState): # 优先使用改写后的query,无改写则用原始问题 query = state.get("rewrite_query", state["question"]) print(f"【检索节点】使用查询词:{query}") docs = retriever.get_relevant_documents(query) return {"docs": docs}节点 2:打分节点,判断文档相关性(分支判断核心)
LLM 判断检索文档是否能回答用户问题,返回good/bad
grade_prompt = ChatPromptTemplate.from_messages([ ("system", "你是文档相关性打分器,只输出good或bad。如果文档内容可以回答用户问题输出good,完全无关输出bad"), ("human", "用户问题:{question}\n检索文档:{docs_content}") ]) def grade_docs_node(state: RAGWorkflowState): question = state["question"] docs = state["docs"] docs_content = "\n".join([d.page_content for d in docs]) chain = grade_prompt | llm res = chain.invoke({"question": question, "docs_content": docs_content}).content print(f"【打分节点】文档判定结果:{res}") return {"grade_result": res}这是LangChain 管道语法(LCEL,LangChain Expression Language),|是管道操作符,作用等价于 Linux 的管道cmd1 | cmd2:把前一个组件的输出,自动作为后一个组件的输入,串联成一条完整执行链路 =链式调用。
拆开grade_prompt | llm完整流程:
grade_prompt:提示词模板 接收传入变量question、docs_content,填充模板,生成完整发给大模型的字符串 Prompt;|管道:将拼接好的完整 Prompt 消息列表,自动传给后面的llm大模型;llm:大模型对象(如 OpenAI、Qwen、本地私有化 LLM) 接收上一步生成的提示词,调用模型推理,返回模型输出结果。
整行代码等价手动分步写法(不用管道,繁琐):
# 1. 填充模板生成消息 prompt_messages = grade_prompt.format_messages(question=question, docs_content=docs_content) # 2. 把消息丢给LLM推理 res_msg = llm.invoke(prompt_messages) # 3. 取文本 res = res_msg.contentgrade_prompt | llm一行替代上面多步,封装成一个可直接invoke()的链对象chain。
节点 3:改写查询词(文档不相关时循环使用)
优化原始问题,生成更精准的检索词
rewrite_prompt = ChatPromptTemplate.from_messages([ ("system", "根据用户原始问题,优化生成更精准的检索关键词,只输出优化后的查询词,不要多余文字"), ("human", "原始问题:{question}") ]) def rewrite_query_node(state: RAGWorkflowState): question = state["question"] retry = state["retry_times"] + 1 chain = rewrite_prompt | llm new_query = chain.invoke({"question": question}).content print(f"【改写节点】新检索词:{new_query},当前重试次数:{retry}") return {"rewrite_query": new_query, "retry_times": retry}节点 4:生成最终回答(文档合格后执行)
gen_prompt = ChatPromptTemplate.from_messages([ ("system", "仅根据提供的知识库文档回答问题,无相关内容直接说无资料"), ("human", "知识库文档:{docs_content}\n用户问题:{question}") ]) def generate_answer_node(state: RAGWorkflowState): question = state["question"] docs_content = "\n".join([d.page_content for d in state["docs"]]) chain = gen_prompt | llm ans = chain.invoke({"question": question, "docs_content": docs_content}).content return {"answer": ans}2.4 定义条件分支函数(控制循环流转)
接收 state,返回下一步节点名称 增加重试次数限制,最多循环 2 次,避免死循环
def route_grade(state: RAGWorkflowState): grade = state["grade_result"] retry = state["retry_times"] # 最多重试2次,强制进入生成节点,防止无限循环 if retry >= 2: print("【流转判断】达到最大重试次数,直接生成回答") return "generate" if grade == "good": return "generate" else: return "rewrite"2.5 组装完整工作流图(搭建节点与边)
# 1. 创建图构建器,绑定状态 graph_builder = StateGraph(RAGWorkflowState) # 2. 注册所有节点 graph_builder.add_node("retrieve", retrieve_node) graph_builder.add_node("grade", grade_docs_node) graph_builder.add_node("rewrite", rewrite_query_node) graph_builder.add_node("generate", generate_answer_node) # 3. 固定流转边 graph_builder.add_edge(START, "retrieve") graph_builder.add_edge("retrieve", "grade") graph_builder.add_edge("rewrite", "retrieve") # 改写后回到检索,实现循环 # 4. 条件分支边(打分后动态选择下一步) graph_builder.add_conditional_edges( source="grade", path=route_grade, path_map={ "generate": "generate", "rewrite": "rewrite" } ) # 5. 生成回答后结束流程 graph_builder.add_edge("generate", END) # 6. 编译工作流,开启断点持久化(记忆保存每一步状态) checkpointer = MemorySaver() rag_workflow = graph_builder.compile(checkpointer=checkpointer)2.6 调用工作流测试运行
# 会话id,区分不同用户的流程状态 config = {"thread_id": "user_001"} # 初始输入状态 input_state = { "question": "FastAPI如何搭配向量数据库做RAG", "rewrite_query": "", "docs": [], "answer": "", "retry_times": 0 } # 执行完整工作流 result = rag_workflow.invoke(input_state, config=config) # 输出最终结果 print("\n===== 工作流执行完成,最终回答 =====") print(result["answer"])三、代码运行效果讲解
场景 1:问题和知识库高度相关
流程: START → 检索 → 打分 (good) → 生成回答 → END 无循环,一次性执行完毕。
场景 2:问题模糊,初次检索文档无关
流程: START → 检索 → 打分 (bad) → 改写 query → 重新检索 → 打分 (good) → 生成回答 完成一次循环优化检索词。
限制保护
设置最大重试 2 次,即使多次检索文档都不相关,也会强制生成回答,杜绝工作流死循环(面试必问:如何防止 Agent / 工作流无限循环)。
四、工作流核心知识点(秋招面试考点)
1. State 全局状态
- 作用:跨节点共享数据,替代 LCEL 手动传递参数
- 缺点:状态税,大量字段频繁读写会增加耗时
- 优化:只保留流程必须的字段,不要冗余存储数据
2. Node 节点设计规范
一个节点只做一件事,单一职责:
- 禁止一个节点同时完成检索 + 打分 + 生成
- 方便单独替换、单元测试、流程微调
3. Edge 两种流转
- 普通边 add_edge:固定单向流转
- 条件边 add_conditional_edges:实现分支、循环、动态路由,是工作流核心能力
4. Checkpoint 断点持久化
MemorySaver 会保存每一步执行后的 State,通过 thread_id 区分会话; 业务价值:支持人工介入、流程中断恢复、多用户并发独立流程。
5. 工作流对比普通 LCEL RAG(面试高频问答)
- LCEL:线性单向,无状态,无法循环、分支;适合简单固定问答
- LangGraph 工作流:有全局状态、支持循环反思、条件判断、断点恢复;适合复杂智能 Agent、自优化 RAG、多工具调用