Langchain搭建LLM应用程序之二 LangGraph+FastAPI
·
应用程序搭建流程图

实施步骤
创建FastAPI 应用
from fastapi import FastAPI
from app.config import settings
from app.routes import router
from app.utils.logging import setup_logging
# 配置日志
setup_logging()
# 创建 FastAPI 应用
app = FastAPI(
title=settings.APP_NAME,
version=settings.APP_VERSION,
docs_url="/docs",
redoc_url="/redoc"
)
定义路由
from typing import Dict, Any
import time
from fastapi import APIRouter, HTTPException, status, Depends, Query
from app.states import BaseResponse, ChatRequest, WorkflowType
from app.core.security import validate_api_key
from app.utils.logging import get_logger
from app.workflows import workflow_manager
logger = get_logger(__name__)
router = APIRouter()
@router.post("/workflow/{workflow_name}", response_model=BaseResponse)
async def execute_workflow(
workflow_name: WorkflowType,
request: Dict[str, Any],
api_key: str = Depends(validate_api_key)
):
"""
执行指定工作流
通过 URL 路径参数路由到正确的工作流
"""
try:
logger.info(f"执行工作流请求", extra={"workflow": workflow_name, "user": request.get("user_id")})
# 添加时间戳
request_data = {**request, "timestamp": time.time()}
# 执行工作流
result = workflow_manager.invoke_workflow(workflow_name.value, request_data)
if not result["success"]:
raise HTTPException(
status_code=status.HTTP_422_UNPROCESSABLE_ENTITY,
detail=result["error"]
)
return BaseResponse(**result)
except ValueError as e:
logger.warning(f"工作流未找到", extra={"workflow": workflow_name, "error": str(e)})
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail=f"工作流 '{workflow_name}' 不存在"
)
except Exception as e:
logger.error(f"工作流执行错误", extra={"workflow": workflow_name, "error": str(e)})
raise HTTPException(
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
detail=f"工作流执行失败: {str(e)}"
)
# 特定工作流的专用端点(可选)
@router.post("/chat", response_model=BaseResponse)
async def chat_workflow(request: ChatRequest, api_key: str = Depends(validate_api_key)):
"""聊天工作流专用端点"""
return await execute_workflow(WorkflowType.CHAT, request.dict(), api_key)
代码中定义了 /chat 路由,并且通过WorkflowType定义工作流名称,通过统一接口execute_workflow执行工作流(需要把要执行的工作流注册到WorkflowManager)
定义状态
from pydantic import BaseModel, Field
from typing import Any, Dict, Optional, List
from enum import Enum
class WorkflowType(str, Enum):
CHAT = "chat"
class BaseRequest(BaseModel):
"""基础请求模型"""
input_data: Dict[str, Any] = Field(..., description="输入数据")
parameters: Optional[Dict[str, Any]] = Field(default={}, description="处理参数")
user_id: Optional[str] = Field(default=None, description="用户ID")
class ChatRequest(BaseRequest):
"""聊天请求"""
message: str = Field(..., description="用户消息")
conversation_id: Optional[str] = Field(default=None, description="会话ID")
class BaseResponse(BaseModel):
"""基础响应模型"""
success: bool = Field(..., description="是否成功")
data: Optional[Dict[str, Any]] = Field(default=None, description="响应数据")
workflow: str = Field(..., description="使用的工作流")
execution_time: float = Field(..., description="执行时间(秒)")
error: Optional[str] = Field(default=None, description="错误信息")
class ChatState(BaseModel):
"""定义工作流状态"""
message: str = ""
cleaned_message: str = ""
message_length: int = 0
intent: str = ""
is_complex: bool = False
response: str = ""
response_type: str = ""
final_response: Dict = None
user_id: str = ""
conversation_id: str = ""
timestamp: float = 0
定义工作流
工作流管理器WorkflowMnager
import importlib
import json
import os
from typing import Dict, Type
from app.workflows.base import BaseWorkflow
from app.workflows.chat import ChatWorkflow
from app.utils.logging import get_logger
from app.utils import BASE_PATH
logger = get_logger(__name__)
default_graph_config_path = os.path.join(BASE_PATH, "config", "langgraph.json")
class WorkflowManager:
"""工作流管理器"""
def __init__(self, config_path:str = default_graph_config_path):
self._workflows: Dict[str, BaseWorkflow] = {}
self.config_path = config_path
self._register_workflows()
def _register_workflows(self):
"""注册所有工作流"""
with open(self.config_path, "r", encoding='utf-8') as f:
graph_config = json.load(f)
for config in graph_config.get("workflows",[]):
workflow_name = config.get("name")
class_path = config.get("class")
module_path, class_name = class_path.rsplit(".", 1)
module = importlib.import_module(module_path)
workflow_class = getattr(module, class_name)
workflow_inst = workflow_class()
self._workflows[workflow_name] = workflow_inst
def get_workflow(self, name: str) -> BaseWorkflow:
"""获取工作流实例"""
if name not in self._workflows:
raise ValueError(f"工作流 '{name}' 未注册")
return self._workflows[name]
def list_workflows(self) -> list:
"""列出所有可用工作流"""
return list(self._workflows.keys())
def invoke_workflow(self, name: str, input_data: Dict) -> Dict:
"""执行指定工作流"""
workflow = self.get_workflow(name)
return workflow.invoke(input_data)
# 全局工作流管理器实例
workflow_manager = WorkflowManager()
通过langgraph.json定义工作流,并动态加载到WorkflowManager类中,以工作流名称与路由定义的工作流名称对应
执行工作流
# 实际工作流
from typing import Literal, Dict, Any
from langgraph.graph import StateGraph, END
from pydantic import BaseModel
from app.workflows.base import BaseWorkflow
from app.llm_models import model_gemma3_1b
class ChatWorkflow(BaseWorkflow):
PROMPT= """### 要求
1. 你是一名乐于助人的助手,你精通python及大模型相关的技术
2. 你可以用专业的话术回答关于python和大模型相关的技术问题
### 问题
{question}
### 输出
1. 你的回答需要准确有依据,如果没有答案,请回答‘不知道’,切勿随意给出不正确或不相关的答案
2. 输出内容要求简明扼要,答案尽可能丰富,但总长度不能超过500
"""
def __init__(self):
super().__init__(name="chat")
self.llm = model_gemma3_1b
def _build_graph(self) -> StateGraph:
"""构建聊天工作流图"""
graph = StateGraph(State)
# 添加节点
graph.add_node("preprocess", self._preprocess_message)
graph.add_node("understand_intent", self._understand_intent)
graph.add_node("generate_response", self._generate_response)
graph.add_node("postprocess", self._postprocess_response)
# 设置工作流路径
graph.set_entry_point("preprocess")
graph.add_edge("preprocess", "understand_intent")
graph.add_conditional_edges(
"understand_intent",
self._route_based_on_intent,
{
"simple": "generate_response",
"complex": "generate_response",
"fallback": "generate_response"
}
)
graph.add_edge("generate_response", "postprocess")
graph.add_edge("postprocess", END)
return graph
def _preprocess_message(self, state: State) -> Dict[str, Any]:
"""预处理消息"""
message = state.message
logger.info("预处理消息", extra={"message": message})
# 简单的文本清理
cleaned_message = message.strip()
return {
"cleaned_message": cleaned_message,
"message_length": len(cleaned_message)
}
def _understand_intent(self, state: State) -> Dict[str, Any]:
"""理解用户意图"""
message = state.cleaned_message.lower()
# 简单的意图识别
if any(word in message for word in ["你好", "嗨", "hello"]):
intent = "greeting"
elif any(word in message for word in ["帮助", "help"]):
intent = "help"
elif any(word in message for word in ["谢谢", "感谢"]):
intent = "thanks"
else:
intent = "general"
logger.info("识别用户意图", extra={"intent": intent, "message": message})
return {
"intent": intent,
"is_complex": len(message) > 50 # 简单复杂度判断
}
def _route_based_on_intent(self, state: State) -> Literal["simple", "complex", "fallback"]:
"""基于意图路由"""
if state.is_complex:
return "complex"
elif state.intent in ["greeting", "thanks"]:
return "simple"
else:
return "fallback"
def _generate_response(self, state: State) -> Dict[str, Any]:
"""生成响应"""
intent = state.intent
message = self.PROMPT.format(question=state.cleaned_message)
# LLM调用
response = self.llm.invoke(message)
response_data = response.content
logger.info("生成响应", extra={"intent": intent, "response": response_data})
return {
"response": response_data,
"response_type": intent
}
def _postprocess_response(self, state: State) -> Dict[str, Any]:
"""后处理响应"""
response = state.response
# 添加响应元数据
processed_response = {
"content": response,
"timestamp": state.timestamp,
"conversation_id": state.conversation_id,
"user_id": state.user_id,
"workflow": "chat"
}
return {
"final_response": processed_response
}
prompt初始化好之后,引入llm model
LLM大模型
from langchain_community.chat_models import ChatOllama
model_gemma3_1b = ChatOllama(
model="gemma3:1b",
base_url= "http://127.0.0.1:11434", # 本地部署的Ollama
num_predict = 2048,
temperature = 0.5,
top_p = 0.9,
top_k = 40
)
本地部署Ollama,下载gemma模型,就可以直接调用,部署详情可参照Ollama官网
响应
大模型调用结果按照定义的response State返回即可
LangChain 高级组件Agent应用请看Langhain搭建LLM应用程序之三
更多推荐



所有评论(0)