工作流学习,含简单上手项目
2026/7/23 23:36:50 网站建设 项目流程

一、什么是工作流

工作流的本质是:

用一张图组织多个任务,并明确任务之间的依赖关系、数据传递、条件分支、并行、汇合、循环和终止规则。

比如一个最基本的 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 标准工作流)

  1. State 全局状态:整个流程共享的数据容器,所有节点读写它;比如用户问题、检索文档、LLM 回答、重试次数、对话历史。
  2. Node 节点:流程里最小执行单元,对应一个业务动作(检索文档、打分、生成回答、改写查询词)。
  3. Edge 边:控制节点流转
    • 普通边:固定跳转 A→B
    • 条件边:根据 State 数据判断下一步去哪里(循环、分支核心)

3、 工作流能解决 RAG 什么痛点

  1. 检索到的文档相关性很低,直接丢给 LLM 会产生幻觉 → 工作流可以打分,不合格就重新检索
  2. 单次 query 检索覆盖不全 → 自动生成多条同义 query,多次检索融合结果
  3. 多轮对话上下文散乱 → 全局 State 统一存放对话历史
  4. 需要人工校验 AI 答案再输出 → 流程暂停,人工修改后继续执行

二、简单项目

1、前置依赖安装

pip install langgraph langchain-openai langchain-community langchain-core pydantic python-dotenv

入门级工作流项目:自反思 RAG 工作流(完整可运行)

2、业务流程设计(工作流流转逻辑)

  1. START 接收用户问题
  2. 节点 1:向量检索知识库,拿到文档片段
  3. 节点 2:LLM 打分,判断文档是否和问题相关
    • 分支 1:文档不相关(分数低)→ 改写用户 query,回到检索节点重新查(循环)
    • 分支 2:文档相关 → 进入生成节点
  4. 节点 3:结合检索文档生成最终回答
  5. 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完整流程:

  1. grade_prompt:提示词模板 接收传入变量questiondocs_content,填充模板,生成完整发给大模型的字符串 Prompt;
  2. |管道:将拼接好的完整 Prompt 消息列表,自动传给后面的llm大模型;
  3. 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.content

grade_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 全局状态

  1. 作用:跨节点共享数据,替代 LCEL 手动传递参数
  2. 缺点:状态税,大量字段频繁读写会增加耗时
  3. 优化:只保留流程必须的字段,不要冗余存储数据

2. Node 节点设计规范

一个节点只做一件事,单一职责:

  • 禁止一个节点同时完成检索 + 打分 + 生成
  • 方便单独替换、单元测试、流程微调

3. Edge 两种流转

  1. 普通边 add_edge:固定单向流转
  2. 条件边 add_conditional_edges:实现分支、循环、动态路由,是工作流核心能力

4. Checkpoint 断点持久化

MemorySaver 会保存每一步执行后的 State,通过 thread_id 区分会话; 业务价值:支持人工介入、流程中断恢复、多用户并发独立流程。

5. 工作流对比普通 LCEL RAG(面试高频问答)

  1. LCEL:线性单向,无状态,无法循环、分支;适合简单固定问答
  2. LangGraph 工作流:有全局状态、支持循环反思、条件判断、断点恢复;适合复杂智能 Agent、自优化 RAG、多工具调用

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询