4.3 工作流编排


文档摘要

4.2 工作流编排 — 大模型应用开发从零到一 关键词短语 本节导读:深入理解大模型应用的工作流编排技术,掌握使用LangGraph等工具构建复杂多步骤任务的能力。 学习目标 理解工作流编排的基本概念和重要性 掌握LangGraph框架的使用方法 学会设计复杂的多步骤工作流 理解状态管理和消息传递机制 构建完整的智能工作流应用 核心概念 什么是工作流编排 工作流编排是指将多个独立的任务步骤按照逻辑顺序组织起来,形成一个完整的业务流程。在大模型应用中,工作流编排允许我们构建能够执行复杂多步骤任务的智能系统。

4.2 工作流编排 — 大模型应用开发从零到一 关键词短语

本节导读:深入理解大模型应用的工作流编排技术,掌握使用LangGraph等工具构建复杂多步骤任务的能力。

学习目标

  • 理解工作流编排的基本概念和重要性
  • 掌握LangGraph框架的使用方法
  • 学会设计复杂的多步骤工作流
  • 理解状态管理和消息传递机制
  • 构建完整的智能工作流应用

核心概念

什么是工作流编排

工作流编排是指将多个独立的任务步骤按照逻辑顺序组织起来,形成一个完整的业务流程。在大模型应用中,工作流编排允许我们构建能够执行复杂多步骤任务的智能系统。

工作流的核心要素

  • 节点(Nodes):执行具体任务的单元
  • 边(Edges):连接节点的逻辑关系
  • 状态(State):在流程中传递的数据
  • 条件(Conditions):控制流程分支的判断逻辑

工作流类型

1. 线性工作流

输入 → 处理步骤1 → 处理步骤2 → 输出

适用于简单的顺序任务,每一步完成后立即进入下一步。

2. 条件分支工作流

输入 → 条件判断 → ├── 分支A → 处理A → 输出 └── 分支B → 处理B → 输出

根据不同条件选择不同的处理路径,适用于需要决策的场景。

3. 循环工作流

输入 → 处理步骤 → 条件判断 → ├── 满足条件 → 输出 └── 不满足条件 → 返回处理步骤

需要重复执行某些步骤直到满足特定条件,适用于迭代优化任务。

4. 并行工作流

输入 → ├── 步骤A → 结果A ├── 步骤B → 结果B └── 步骤C → 结果C → 汇总处理 → 输出

多个步骤并行执行,最后汇总结果,适用于需要同时处理多个任务的高效场景。

环境准备 / 前置知识

Python环境配置

# 安装LangGraph和相关依赖 pip install langgraph langchain langchain-openai pip install python-dotenv loguru

基础依赖检查

import langgraph import langchain from langchain_openai import ChatOpenAI print(f"LangGraph版本: {langgraph.__version__}") print(f"LangChain版本: {langchain.__version__}")

模型配置

from langchain_openai import ChatOpenAI from dotenv import load_dotenv import os # 加载环境变量 load_dotenv() # 配置语言模型 llm = ChatOpenAI( model_name="gpt-4-turbo-preview", temperature=0.7, max_tokens=4000 )

分步实战

步骤 1:LangGraph基础概念

LangGraph是一个用于构建复杂工作流的图计算框架,特别适合大模型应用的编排需求。

#!/usr/bin/env python3 # -*- coding: utf-8 -*- """ LangGraph基础概念示例 """ from langgraph.graph import Graph, END from typing import Dict, Any, List from langchain_openai import ChatOpenAI from langchain.prompts import PromptTemplate class SimpleWorkflow: """简单工作流示例""" def __init__(self): self.llm = ChatOpenAI( model_name="gpt-4-turbo-preview", temperature=0.7 ) # 创建工作流图 self.workflow = Graph() # 添加节点 self.workflow.add_node("input", self.input_handler) self.workflow.add_node("process", self.process_handler) self.workflow.add_node("output", self.output_handler) # 添加边 self.workflow.add_edge("input", "process") self.workflow.add_edge("process", "output") self.workflow.add_edge("output", END) # 编译工作流 self.app = self.workflow.compile() def input_handler(self, state: Dict[str, Any]) -> Dict[str, Any]: """输入处理节点""" print("=== 输入处理 ===") print(f"原始输入: {state.get('input', '')}") # 预处理输入 processed_input = { **state, "processed_text": state.get('input', '').strip(), "timestamp": "2026-07-12T10:00:00" } return processed_input def process_handler(self, state: Dict[str, Any]) -> Dict[str, Any]: """处理节点""" print("=== 处理中 ===") input_text = state.get('processed_text', '') # 使用LLM进行处理 prompt = PromptTemplate( input_variables=["text"], template="请对以下文本进行分析和总结:\\n{text}" ) chain = prompt | self.llm result = chain.invoke({"text": input_text}) print(f"LLM处理结果: {result.content[:100]}...") return { **state, "llm_result": result.content, "processing_complete": True } def output_handler(self, state: Dict[str, Any]) -> Dict[str, Any]: """输出处理节点""" print("=== 输出生成 ===") # 格式化输出 output_data = { "original_input": state.get('input', ''), "processed_text": state.get('processed_text', ''), "llm_result": state.get('llm_result', ''), "timestamp": state.get('timestamp', ''), "status": "completed" } print("最终输出:") print(f"- 原始输入: {output_data['original_input']}") print(f"- 处理状态: {output_data['status']}") print(f"- 处理时间: {output_data['timestamp']}") return output_data def run(self, input_text: str) -> Dict[str, Any]: """运行工作流""" initial_state = {"input": input_text} return self.app.invoke(initial_state) # 使用示例 def basic_workflow_example(): """基础工作流示例""" workflow = SimpleWorkflow() # 测试不同类型的输入 test_inputs = [ "人工智能正在改变我们的生活方式", "大语言模型能够理解和生成人类语言", "机器学习算法从数据中学习模式" ] print("=== LangGraph基础工作流测试 ===") for i, text in enumerate(test_inputs, 1): print(f"\\n--- 测试案例 {i} ---") result = workflow.run(text) print(f"输入长度: {len(text)} 字符") print(f"输出状态: {result['status']}") print(f"处理时间: {result['timestamp']}") if __name__ == "__main__": basic_workflow_example()

步骤 2:条件分支工作流

在实际应用中,我们经常需要根据不同条件选择不同的处理路径。

#!/usr/bin/env python3 # -*- coding: utf-8 -*- """ 条件分支工作流示例 """ from langgraph.graph import Graph, END from typing import Dict, Any, List from langchain_openai import ChatOpenAI from langchain.prompts import PromptTemplate class ConditionalWorkflow: """条件分支工作流""" def __init__(self): self.llm = ChatOpenAI( model_name="gpt-4-turbo-preview", temperature=0.3 ) self.workflow = Graph() # 添加节点 self.workflow.add_node("input", self.input_analyzer) self.workflow.add_node("text_analysis", self.text_analysis) self.workflow.add_node("code_analysis", self.code_analysis) self.workflow.add_node("summary", self.summary_generator) self.workflow.add_node("final_output", self.final_output) # 添加条件和边 self.workflow.add_edge("input", "text_analysis") self.workflow.add_edge("input", "code_analysis") # 条件分支:根据内容类型选择不同路径 self.workflow.add_conditional_edges( "text_analysis", self.route_to_summary, { "code": "code_analysis", "text": "summary", "mixed": "summary" } ) self.workflow.add_edge("code_analysis", "summary") self.workflow.add_edge("summary", "final_output") self.workflow.add_edge("final_output", END) self.app = self.workflow.compile() def input_analyzer(self, state: Dict[str, Any]) -> Dict[str, Any]: """输入分析节点""" input_text = state.get('input', '') # 简单的内容类型检测 has_code = "```" in input_text has_text = len(input_text) > 50 and not has_code content_type = "mixed" if has_code and not has_text: content_type = "code" elif has_text and not has_code: content_type = "text" return { **state, "content_type": content_type, "input_length": len(input_text), "has_code": has_code, "has_text": has_text } def route_to_summary(self, state: Dict[str, Any]) -> str: """路由到总结节点""" if state.get('content_type') == 'code': return "code" else: return "text" def run(self, input_text: str) -> Dict[str, Any]: """运行条件分支工作流""" initial_state = {"input": input_text} return self.app.invoke(initial_state) # 使用示例 def conditional_workflow_example(): """条件分支工作流示例""" workflow = ConditionalWorkflow() test_cases = [ { "name": "纯文本分析", "content": "人工智能技术正在快速发展,深度学习算法在图像识别、自然语言处理等领域取得了突破性进展。机器学习已经成为现代科技发展的重要驱动力。" }, { "name": "代码分析", "content": "def calculate_fibonacci(n): return n if n <= 1 else calculate_fibonacci(n-1) + calculate_fibonacci(n-2)" } ] print("=== 条件分支工作流测试 ===") for test_case in test_cases: print(f"\\n--- {test_case['name']} ---") result = workflow.run(test_case['content']) print(f"内容类型: {result['content_type']}") print(f"输入长度: {result['input_length']} 字符") if __name__ == "__main__": conditional_workflow_example()

步骤 3:复杂工作流状态管理

在复杂的工作流中,状态管理是关键。我们需要在多个节点之间传递和更新状态。

#!/usr/bin/env python3 # -*- coding: utf-8 -*- """ 复杂工作流状态管理示例 """ from langgraph.graph import Graph, END from langgraph.checkpoint.memory import MemorySaver from typing import Dict, Any, List, TypedDict from datetime import datetime from langchain_openai import ChatOpenAI class WorkflowState(TypedDict): """工作流状态定义""" input_text: str current_step: str processing_history: List[str] intermediate_results: Dict[str, Any] final_result: Dict[str, Any] metadata: Dict[str, Any] class ComplexWorkflow: """复杂工作流示例""" def __init__(self): self.llm = ChatOpenAI( model_name="gpt-4-turbo-preview", temperature=0.5 ) # 使用内存检查点保存工作流状态 memory = MemorySaver() self.workflow = Graph(checkpointer=memory) # 添加节点 self.workflow.add_node("initialize", self.initialize_workflow) self.workflow.add_node("preprocess", self.preprocess_text) self.workflow.add_node("analyze", self.content_analysis) self.workflow.add_node("finalize", self.finalize_result) # 添加工作流边 self.workflow.add_edge("initialize", "preprocess") self.workflow.add_edge("preprocess", "analyze") self.workflow.add_edge("analyze", "finalize") self.workflow.add_edge("finalize", END) self.app = self.workflow.compile() def initialize_workflow(self, state: WorkflowState) -> WorkflowState: """初始化工作流""" print("=== 工作流初始化 ===") return { **state, "current_step": "initialized", "processing_history": ["workflow_initialized"], "metadata": { "start_time": datetime.now().isoformat(), "workflow_version": "1.0" } } def preprocess_text(self, state: WorkflowState) -> WorkflowState: """文本预处理""" print("=== 文本预处理 ===") input_text = state["input_text"] # 文本清洗和标准化 cleaned_text = input_text.strip() word_count = len(cleaned_text.split()) preprocessing_result = { "cleaned_text": cleaned_text, "word_count": word_count, "preprocessing_complete": True } # 更新状态 updated_state = { **state, "current_step": "preprocessed", "processing_history": state["processing_history"] + ["preprocessing_completed"], "intermediate_results": { **state.get("intermediate_results", {}), "preprocessing": preprocessing_result } } print(f"预处理完成: {word_count} 词") return updated_state def run(self, input_text: str) -> WorkflowState: """运行复杂工作流""" initial_state: WorkflowState = { "input_text": input_text, "current_step": "initial", "processing_history": [], "intermediate_results": {}, "final_result": {}, "metadata": {} } return self.app.invoke(initial_state) # 使用示例 def complex_workflow_example(): """复杂工作流示例""" workflow = ComplexWorkflow() test_input = "在人工智能领域,深度学习算法已经成为重要的技术手段。" print("=== 复杂工作流状态管理测试 ===") # 运行工作流 result = workflow.run(test_input) # 输出结果摘要 print(f"\\n=== 工作流结果摘要 ===") print(f"处理步骤: {' → '.join(result['processing_history'])}") print(f"最终状态: {result['current_step']}") print(f"处理词数: {result['intermediate_results']['preprocessing']['word_count']}") if __name__ == "__main__": complex_workflow_example()

常见问题 FAQ

Q: LangGraph与其他工作流框架有什么区别?

A: LangGraph专门为大模型应用设计,支持状态管理、条件分支、循环等复杂逻辑,与LLM深度集成。相比传统的Airflow、Prefect等框架,LangGraph更适合处理基于LLM的智能任务。

Q: 如何处理工作流中的错误和异常?

A: 在每个节点中添加异常处理逻辑,使用try-catch捕获错误。可以设置重试机制、降级处理和错误记录,确保工作流的健壮性。

Q: 工作流状态如何持久化?

A: LangGraph支持多种检查点后端(内存、文件、数据库),可以保存工作流状态以便恢复和监控。在生产环境中建议使用数据库作为后端。

Q: 如何优化工作流性能?

A: 通过并行处理、缓存机制、结果复用等方法优化。合理的节点设计和状态管理可以显著提升工作流执行效率。

Q: 工作流监控和调试有哪些方法?

A: 使用LangGraph的日志记录、状态可视化、性能监控工具。可以通过检查点查看工作流执行历史,便于调试和优化。

本节小结

通过本节的学习,你已经掌握了:

✅ 工作流编排的基本概念和类型 ✅ LangGraph框架的使用方法
✅ 条件分支和状态管理技术 ✅ 复杂工作流的设计和实现
✅ 完整的智能工作流应用构建

这些技术将帮助你构建能够执行复杂多步骤任务的大模型应用,大幅提升系统的智能化程度和自动化水平。

延伸阅读

  • 官方文档:LangGraph框架使用指南(文字描述,不带链接)
  • 相关章节:本教程 4.3 节 Agent协作(文字描述,不带链接)
  • 技术博客:大模型工作流编排最佳实践(文字描述,不带链接)

关键词:大模型应用开发从零到一, 工作流编排, LangGraph, 状态管理, 条件分支, 智能工作流, 多步骤任务
难度:进阶
预计阅读:40 分钟


发布者: 作者: 内存溢出警告的小龙虾 转发
评论区 (0)
U