1. LangGraph 基础概念解析
LangGraph 作为 LangChain 生态系统的扩展组件,为构建复杂智能体工作流提供了一种全新的范式。与传统的对话式框架不同,LangGraph 采用状态机(State Machine)和有向图(Directed Graph)的模型来定义智能体的执行流程。这种设计理念使得构建具备循环、分支和状态保持能力的智能体变得直观且可控。
1.1 核心组件详解
LangGraph 的工作流由三个基本要素构成:
-
全局状态(State):这是贯穿整个工作流的共享数据容器,通常定义为 Python 的 TypedDict。它保存了智能体运行过程中需要追踪的所有信息,如对话历史、中间结果、当前步骤等。状态对象在节点间传递,每个节点都可以读取和修改其中的数据。
-
节点(Nodes):节点是工作流中的基本执行单元,每个节点都是一个 Python 函数,接收当前状态作为输入,返回更新后的状态。节点可以执行各种操作,如调用 LLM、执行工具、处理数据等。在我们的问答助手示例中,定义了三个核心节点:理解查询节点、搜索节点和回答节点。
-
边(Edges):边定义了节点之间的跳转逻辑。LangGraph 支持两种边类型:
- 常规边:固定指向下一个节点
- 条件边:根据状态动态决定跳转目标
1.2 状态机模型优势
状态机模型为智能体开发带来了几个关键优势:
- 显式流程控制:开发者可以清晰地看到整个工作流的执行路径
- 原生支持循环:通过条件边轻松实现"反思-修正"等循环逻辑
- 模块化设计:每个节点都是独立的函数,便于测试和维护
- 状态持久化:全局状态自动维护执行上下文,无需额外处理
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 问答助手实现详解
2.1 环境准备与初始化
在开始构建问答助手前,需要完成以下准备工作:
- 安装依赖包:
bash复制pip install langgraph langchain-openai tavily-python python-dotenv
- 配置环境变量:
创建.env文件并添加以下内容:
code复制LLM_API_KEY=your_openai_api_key
LLM_MODEL_ID=gpt-4
TAVILY_API_KEY=your_tavily_api_key
- 导入必要模块:
python复制import os
from dotenv import load_dotenv
from typing import TypedDict, Annotated
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage, AIMessage
from langgraph.graph import StateGraph
from langgraph.graph.message import add_messages
from tavily import TavilyClient
# 加载环境变量
load_dotenv()
2.2 状态定义与模型初始化
定义全局状态结构是构建 LangGraph 工作流的第一步。我们的问答助手需要跟踪以下信息:
python复制class SearchState(TypedDict):
messages: Annotated[list, add_messages] # 对话历史
user_query: str # 用户原始查询
search_query: str # 优化后的搜索关键词
search_results: str # 搜索结果
final_answer: str # 最终生成的答案
step: str # 当前执行步骤
初始化 LLM 和 Tavily 客户端:
python复制llm = ChatOpenAI(
model=os.getenv("LLM_MODEL_ID", "gpt-4"),
api_key=os.getenv("LLM_API_KEY"),
temperature=0.7
)
tavily_client = TavilyClient(api_key=os.getenv("TAVILY_API_KEY"))
2.3 核心节点实现
2.3.1 理解查询节点
这个节点负责分析用户意图并生成优化的搜索关键词:
python复制def understand_query_node(state: SearchState) -> dict:
"""步骤1:理解用户查询并生成搜索关键词"""
user_message = state["messages"][-1].content
understand_prompt = f"""分析用户的查询:"{user_message}"
请完成两个任务:
1. 简洁总结用户想要了解什么
2. 生成最适合搜索引擎的关键词(中英文均可,要精准)
格式:
理解:[用户需求总结]
搜索词:[最佳搜索关键词]"""
response = llm.invoke([HumanMessage(content=understand_prompt)])
response_text = response.content
# 解析LLM输出提取搜索词
search_query = user_message # 默认使用原始查询
if "搜索词:" in response_text:
search_query = response_text.split("搜索词:")[1].strip()
return {
"user_query": response_text,
"search_query": search_query,
"step": "understood",
"messages": [AIMessage(content=f"我理解您的需求:{response_text}")]
}
2.3.2 搜索节点
这个节点调用 Tavily API 执行实际搜索:
python复制def tavily_search_node(state: SearchState) -> dict:
"""步骤2:使用Tavily API进行真实搜索"""
search_query = state["search_query"]
try:
print(f"🔍 正在搜索: {search_query}")
response = tavily_client.search(
query=search_query,
search_depth="basic",
include_answer=True,
max_results=5
)
# 处理搜索结果
search_results = ""
if response.get("answer"):
search_results = f"综合答案:\n{response['answer']}\n\n"
if response.get("results"):
search_results += "相关信息:\n"
for i, result in enumerate(response["results"][:3], 1):
search_results += f"{i}. {result.get('title','')}\n{result.get('content','')}\n来源:{result.get('url','')}\n\n"
return {
"search_results": search_results or "未找到相关信息",
"step": "searched",
"messages": [AIMessage(content="✅ 搜索完成!正在整理答案...")]
}
except Exception as e:
return {
"search_results": f"搜索失败:{str(e)}",
"step": "search_failed",
"messages": [AIMessage(content="❌ 搜索遇到问题,我将基于已有知识为您回答")]
}
2.3.3 回答生成节点
基于搜索结果生成最终回答:
python复制def generate_answer_node(state: SearchState) -> dict:
"""步骤3:基于搜索结果生成最终答案"""
if state["step"] == "search_failed":
# 搜索失败时的回退策略
fallback_prompt = f"请基于您的知识回答:{state['user_query']}"
response = llm.invoke([HumanMessage(content=fallback_prompt)])
else:
# 基于搜索结果生成答案
answer_prompt = f"""基于以下信息回答问题:
用户问题:{state['user_query']}
搜索结果:
{state['search_results']}
要求:
1. 提供准确、完整的回答
2. 引用重要信息来源
3. 结构清晰,分点说明"""
response = llm.invoke([HumanMessage(content=answer_prompt)])
return {
"final_answer": response.content,
"step": "completed",
"messages": [AIMessage(content=response.content)]
}
2.4 工作流组装与执行
将各个节点连接成完整的工作流:
python复制def create_search_assistant():
workflow = StateGraph(SearchState)
# 添加节点
workflow.add_node("understand", understand_query_node)
workflow.add_node("search", tavily_search_node)
workflow.add_node("answer", generate_answer_node)
# 设置线性流程
workflow.add_edge(START, "understand")
workflow.add_edge("understand", "search")
workflow.add_edge("search", "answer")
workflow.add_edge("answer", END)
# 编译工作流
return workflow.compile()
3. 高级功能与优化
3.1 条件分支与循环
LangGraph 的强大之处在于支持条件分支和循环。例如,我们可以改进搜索节点,在搜索结果不理想时自动重试:
python复制def should_retry_search(state: SearchState) -> str:
"""判断是否需要重试搜索"""
if "未找到相关信息" in state["search_results"]:
return "retry"
return "continue"
# 修改工作流构建
workflow.add_conditional_edges(
"search",
should_retry_search,
{
"retry": "search", # 重试搜索
"continue": "answer" # 继续生成答案
}
)
3.2 状态检查点
LangGraph 支持状态持久化,可以保存和恢复执行状态:
python复制from langgraph.checkpoint.memory import InMemorySaver
def create_search_assistant():
workflow = StateGraph(SearchState)
# ... 添加节点和边 ...
# 添加状态检查点
memory = InMemorySaver()
return workflow.compile(checkpointer=memory)
3.3 多模态扩展
我们可以扩展工作流支持多模态处理,例如添加图像理解节点:
python复制from langchain_community.tools import ImageCaptionTool
def image_understand_node(state: SearchState) -> dict:
"""处理图像内容"""
image_url = extract_image_url(state["user_query"])
if not image_url:
return state
tool = ImageCaptionTool()
caption = tool.run(image_url)
return {
**state,
"image_caption": caption,
"messages": [AIMessage(content=f"识别到图像内容:{caption}")]
}
4. 生产环境部署建议
4.1 性能优化
- 异步执行:
python复制async def async_search_node(state: SearchState) -> dict:
"""异步搜索节点"""
# 实现异步搜索逻辑
pass
- 缓存机制:
python复制from functools import lru_cache
@lru_cache(maxsize=100)
def cached_search(query: str) -> dict:
"""带缓存的搜索"""
return tavily_client.search(query=query)
4.2 错误处理与监控
- 增强错误处理:
python复制def safe_llm_invoke(prompt):
"""带重试机制的LLM调用"""
max_retries = 3
for attempt in range(max_retries):
try:
return llm.invoke(prompt)
except Exception as e:
if attempt == max_retries - 1:
raise
time.sleep(2 ** attempt)
- 添加日志记录:
python复制import logging
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s'
)
logger = logging.getLogger(__name__)
def logged_search_node(state: SearchState) -> dict:
"""带日志记录的搜索节点"""
logger.info(f"开始搜索: {state['search_query']}")
# ... 搜索逻辑 ...
4.3 可扩展性设计
- 插件式架构:
python复制class NodePlugin:
"""节点插件基类"""
def process(self, state: dict) -> dict:
raise NotImplementedError
class SentimentPlugin(NodePlugin):
"""情感分析插件"""
def process(self, state):
# 分析用户查询情感
return state
- 配置化工作流:
python复制def build_workflow_from_config(config: dict) -> StateGraph:
"""根据配置文件构建工作流"""
workflow = StateGraph(SearchState)
for node in config["nodes"]:
workflow.add_node(node["name"], globals()[node["func"]])
# ... 添加边 ...
return workflow.compile()
5. 常见问题与解决方案
5.1 搜索相关问题
问题1:搜索结果质量不高
解决方案:
- 优化搜索关键词生成提示词
- 尝试不同的搜索深度参数
- 添加搜索结果过滤逻辑
python复制def optimize_search_query(query: str) -> str:
"""优化搜索查询"""
prompt = f"""原始查询:{query}
请生成3个更可能获得优质结果的搜索关键词,用逗号分隔:"""
response = llm.invoke([HumanMessage(content=prompt)])
return response.content.split(",")[0].strip()
问题2:API调用频率限制
解决方案:
- 实现请求限流
- 添加缓存层
- 使用指数退避重试
python复制from tenacity import retry, stop_after_attempt, wait_exponential
@retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=4, max=10))
def safe_tavily_search(query: str) -> dict:
"""带重试的搜索"""
return tavily_client.search(query=query)
5.2 LLM相关问题
问题1:回答不准确
解决方案:
- 改进提示词工程
- 添加事实核查步骤
- 设置温度参数为更低值
python复制def fact_check_answer(answer: str, sources: str) -> str:
"""事实核查"""
prompt = f"""请核查以下回答是否与提供的信息一致:
回答:{answer}
来源信息:{sources}
请指出任何不一致之处并提供修正建议:"""
return llm.invoke([HumanMessage(content=prompt)]).content
问题2:响应速度慢
解决方案:
- 使用流式响应
- 实现异步处理
- 考虑使用更轻量级的模型
python复制async def stream_answer(state: SearchState, websocket):
"""流式生成答案"""
prompt = build_answer_prompt(state)
chunks = []
async for chunk in llm.astream([HumanMessage(content=prompt)]):
await websocket.send_text(chunk.content)
chunks.append(chunk.content)
return "".join(chunks)
5.3 工作流调试技巧
- 状态可视化:
python复制def print_state(state: dict):
"""打印状态快照"""
print("=== 状态快照 ===")
for k, v in state.items():
if k != "messages":
print(f"{k}: {str(v)[:100]}{'...' if len(str(v))>100 else ''}")
print("最新消息:", state["messages"][-1].content)
- 单步调试模式:
python复制def debug_node(node_func, test_state: dict):
"""调试单个节点"""
print("=== 输入状态 ===")
print_state(test_state)
output = node_func(test_state)
print("\n=== 输出状态 ===")
print_state(output)
return output
- 测试用例生成:
python复制def generate_test_cases(num_cases=5) -> list:
"""生成测试用例"""
prompt = """请生成5个测试查询,涵盖不同类型的问题:
1. 事实性问题
2. 观点性问题
3. 技术性问题
4. 模糊查询
5. 多部分问题
格式:
1. [类型]: [查询]"""
response = llm.invoke([HumanMessage(content=prompt)]).content
return [line.split(": ")[1] for line in response.split("\n") if ": " in line]
6. 进阶应用场景
6.1 多智能体协作
LangGraph 可以协调多个智能体协同工作:
python复制def create_team_assistant():
workflow = StateGraph(SearchState)
# 添加不同角色的智能体节点
workflow.add_node("researcher", research_node)
workflow.add_node("analyst", analysis_node)
workflow.add_node("reviewer", review_node)
# 定义协作流程
workflow.add_edge(START, "researcher")
workflow.add_edge("researcher", "analyst")
workflow.add_edge("analyst", "reviewer")
# 添加循环审核机制
def needs_revision(state):
return "needs_revision" in state.get("feedback", "")
workflow.add_conditional_edges(
"reviewer",
needs_revision,
{
True: "analyst", # 需要修改
False: END # 审核通过
}
)
return workflow.compile()
6.2 长文档处理
扩展问答助手处理长文档:
python复制def process_document_node(state: SearchState) -> dict:
"""处理上传的文档"""
if not state.get("document"):
return state
from langchain_text_splitters import RecursiveCharacterTextSplitter
text_splitter = RecursiveCharacterTextSplitter(
chunk_size=1000,
chunk_overlap=200
)
chunks = text_splitter.split_text(state["document"])
return {
**state,
"document_chunks": chunks,
"messages": [AIMessage(content=f"已将文档分割为{len(chunks)}个段落")]
}
6.3 实时数据流处理
适应实时数据场景:
python复制async def handle_data_stream(websocket):
"""处理WebSocket数据流"""
app = create_search_assistant()
async for message in websocket:
state = {"messages": [HumanMessage(content=message)]}
async for output in app.astream(state):
for node_name, node_output in output.items():
if "messages" in node_output:
await websocket.send_text(node_output["messages"][-1].content)
在实际部署中,我发现为每个工作流实例维护独立的会话上下文非常重要。通过将 thread_id 绑定到用户会话,可以实现多轮对话的状态保持:
python复制async def handle_user_session(user_id: str, websocket):
"""处理用户会话"""
app = create_search_assistant()
config = {"configurable": {"thread_id": user_id}}
async for message in websocket:
state = {"messages": [HumanMessage(content=message)]}
async for output in app.astream(state, config=config):
# 处理输出...
另一个实用技巧是在状态中维护自定义的元数据,这可以帮助调试和监控:
python复制class EnhancedState(TypedDict):
messages: Annotated[list, add_messages]
metadata: dict # 自定义元数据字段
# ...其他字段...
def track_metrics_node(state: EnhancedState) -> EnhancedState:
"""跟踪性能指标"""
start_time = time.time()
# ...节点逻辑...
return {
**state,
"metadata": {
**state.get("metadata", {}),
"processing_time": time.time() - start_time,
"node_name": "track_metrics"
}
}
对于需要人工审核的场景,可以设计人工干预节点:
python复制def human_review_node(state: SearchState) -> dict:
"""等待人工审核"""
send_for_review(state["final_answer"])
return {
**state,
"status": "pending_review",
"messages": [AIMessage(content="回答已提交人工审核,请稍候...")]
}
def check_review_status(state: SearchState) -> str:
"""检查审核状态"""
if state.get("review_approved"):
return "approved"
return "rejected"
