应用程序搭建流程图

在这里插入图片描述

实施步骤
创建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应用程序之三

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐