工作流学习,含简单上手项目

📅 2026/7/23 23:37:00
工作流学习,含简单上手项目
一、什么是工作流工作流的本质是用一张图组织多个任务并明确任务之间的依赖关系、数据传递、条件分支、并行、汇合、循环和终止规则。比如一个最基本的 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向量检索知识库拿到文档片段节点 2LLM 打分判断文档是否和问题相关分支 1文档不相关分数低→ 改写用户 query回到检索节点重新查循环分支 2文档相关 → 进入生成节点节点 3结合检索文档生成最终回答END 输出答案流程图文字版 START → 检索节点 → 打分节点 打分不通过 → 改写 query → 检索节点循环 打分通过 → 生成回答 → END步骤 1环境文件 .envOPENAI_API_KEYsk-xxx EMBED_MODELBAAI/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(modelgpt-3.5-turbo, api_keyos.getenv(OPENAI_API_KEY), temperature0) # 模拟本地向量库你可替换Qdrant # 测试知识库文档 test_docs [ Document(page_contentFastAPI是Python高性能异步Web框架常用于RAG后端开发), Document(page_contentLangGraph用于构建带循环、分支的LLM工作流解决普通RAG无法反思优化的问题), Document(page_contentQdrant是生产级向量数据库支持异步客户端适配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/badgrade_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 管道语法LCELLangChain 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(questionquestion, docs_contentdocs_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 rewrite2.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( sourcegrade, pathroute_grade, path_map{ generate: generate, rewrite: rewrite } ) # 5. 生成回答后结束流程 graph_builder.add_edge(generate, END) # 6. 编译工作流开启断点持久化记忆保存每一步状态 checkpointer MemorySaver() rag_workflow graph_builder.compile(checkpointercheckpointer)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, configconfig) # 输出最终结果 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、多工具调用