1. 项目概述

1.1 项目背景

随着客户服务需求的不断增长,传统的客服系统已经难以满足企业的需求。智能客服系统通过引入人工智能技术,能够24小时不间断地为客户提供服务,提高客户满意度,同时降低企业运营成本。

多智能体系统(MAS)在智能客服领域的应用具有显著优势:

  • 分工协作:不同智能体负责不同类型的任务,提高处理效率
  • 知识共享:智能体之间可以共享知识和经验,提升整体服务质量
  • 灵活性:系统可以根据业务需求灵活调整智能体的数量和职责
  • 可扩展性:易于添加新的智能体和功能,适应业务发展

1.2 项目目标

本项目旨在构建一个基于多智能体系统的智能客服平台,实现以下目标:

  1. 多渠道接入:支持网站、APP、微信、电话等多种渠道
  2. 智能分流:根据客户问题类型自动分流到相应的智能体
  3. 智能问答:基于知识库和历史数据,提供准确的回答
  4. 多轮对话:支持复杂问题的多轮对话处理
  5. 业务办理:支持常见业务的自动化办理
  6. 人工协作:在智能体无法解决问题时,无缝转接人工客服
  7. 数据分析:收集和分析客户反馈,持续优化系统

1.3 技术栈选择

类别 技术/框架 版本 用途
开发语言 Python 3.12+ 智能体开发
AI框架 LangChain 0.3.x LLM集成、工具链
LangGraph 0.2.x 多智能体协作管理
大语言模型 OpenAI GPT-4o - 核心对话能力
阿里云通义千问 - 中文语境优化
向量数据库 ChromaDB 0.4.x 知识库存储和检索
缓存 Redis 7.0+ 会话管理、临时数据存储
数据库 PostgreSQL 15.0+ 结构化数据存储
消息队列 RabbitMQ 3.12+ 任务分发、异步处理
Web框架 FastAPI 0.104+ API接口开发
容器化 Docker 20.10+ 环境隔离、部署一致性
编排 Kubernetes 1.26+ 容器编排、自动扩缩容
监控 Prometheus + Grafana 2.40+ / 9.0+ 系统监控、可视化

2. 系统架构设计

2.1 整体架构


2.2 智能体角色设计

智能体角色 职责 核心功能 技术实现
接入智能体 统一接入管理 多渠道消息接收与发送、格式转换、会话管理 FastAPI + WebSocket
意图识别智能体 客户意图分析 文本分类、意图识别、实体提取 LLM + 微调模型
分流智能体 智能任务分发 基于意图和上下文的任务分配、优先级管理 LangGraph路由节点
问答智能体 知识型问题解答 知识库检索、答案生成、多轮对话管理 RAG + LLM
业务办理智能体 业务流程处理 表单收集、业务系统调用、结果反馈 工具调用 + 状态管理
情绪管理智能体 客户情绪分析 情绪识别、情绪调节、危机干预 NLP模型 + 规则引擎
人工协作智能体 人工客服协作 人工转接、上下文传递、协作管理 WebSocket + 队列

2.3 数据流设计

  1. 客户请求流

    • 客户通过多渠道发送请求
    • 接入智能体接收并标准化请求
    • 意图识别智能体分析客户意图
    • 分流智能体将请求分发到对应智能体
    • 专业智能体处理请求并生成响应
    • 接入智能体将响应返回给客户
  2. 知识流

    • 知识库定期更新和维护
    • 向量数据库存储知识嵌入
    • 问答智能体检索相关知识
    • 智能体间共享处理经验
    • 系统从对话中学习新知识
  3. 监控流

    • 各智能体实时上报状态和性能指标
    • 监控系统收集和分析指标
    • 异常情况触发告警
    • 运维人员处理告警和故障

3. 核心功能实现

3.1 接入层实现

3.1.1 多渠道接入
# 接入智能体实现
from fastapi import FastAPI, WebSocket, WebSocketDisconnect
from fastapi.middleware.cors import CORSMiddleware
import json
import asyncio
from typing import Dict, Set

class ConnectionManager:
    def __init__(self):
        self.active_connections: Dict[str, Dict[str, WebSocket]] = {}
    
    async def connect(self, websocket: WebSocket, channel: str, user_id: str):
        await websocket.accept()
        if channel not in self.active_connections:
            self.active_connections[channel] = {}
        self.active_connections[channel][user_id] = websocket
    
    def disconnect(self, channel: str, user_id: str):
        if channel in self.active_connections and user_id in self.active_connections[channel]:
            del self.active_connections[channel][user_id]
    
    async def send_personal_message(self, message: str, channel: str, user_id: str):
        if channel in self.active_connections and user_id in self.active_connections[channel]:
            await self.active_connections[channel][user_id].send_text(message)
    
    async def broadcast(self, message: str, channel: str):
        if channel in self.active_connections:
            for connection in self.active_connections[channel].values():
                await connection.send_text(message)

app = FastAPI()
app.add_middleware(
    CORSMiddleware,
    allow_origins=["*"],
    allow_credentials=True,
    allow_methods=["*"],
    allow_headers=["*"],
)

manager = ConnectionManager()

@app.websocket("/ws/{channel}/{user_id}")
async def websocket_endpoint(websocket: WebSocket, channel: str, user_id: str):
    await manager.connect(websocket, channel, user_id)
    try:
        while True:
            data = await websocket.receive_text()
            # 处理接收到的消息
            message_data = json.loads(data)
            
            # 消息格式转换和标准化
            standardized_message = {
                "channel": channel,
                "user_id": user_id,
                "message_type": message_data.get("type", "text"),
                "content": message_data.get("content", ""),
                "timestamp": message_data.get("timestamp", time.time()),
                "metadata": message_data.get("metadata", {})
            }
            
            # 转发给意图识别智能体
            response = await process_message(standardized_message)
            
            # 返回响应给客户端
            await manager.send_personal_message(json.dumps(response), channel, user_id)
    except WebSocketDisconnect:
        manager.disconnect(channel, user_id)
        # 处理断开连接的逻辑

async def process_message(message):
    """处理消息并返回响应"""
    # 这里将消息转发给意图识别智能体
    # 实际实现中应使用消息队列或直接调用
    from agents.intent_agent import IntentRecognitionAgent
    
    intent_agent = IntentRecognitionAgent()
    result = await intent_agent.analyze_intent(message)
    
    return result

if __name__ == "__main__":
    import uvicorn
    uvicorn.run(app, host="0.0.0.0", port=8000)
3.1.2 会话管理
# 会话管理服务
import redis
import json
import time
from typing import Dict, Any, Optional

class SessionManager:
    def __init__(self, redis_url: str = "redis://localhost:6379/0"):
        self.redis_client = redis.from_url(redis_url, decode_responses=True)
        self.session_ttl = 3600  # 会话过期时间(秒)
    
    def create_session(self, user_id: str, channel: str) -> str:
        """创建新会话"""
        session_id = f"session:{user_id}:{int(time.time())}"
        session_data = {
            "session_id": session_id,
            "user_id": user_id,
            "channel": channel,
            "created_at": time.time(),
            "last_activity": time.time(),
            "messages": [],
            "context": {}
        }
        
        self.redis_client.setex(
            session_id,
            self.session_ttl,
            json.dumps(session_data, ensure_ascii=False)
        )
        
        return session_id
    
    def get_session(self, session_id: str) -> Optional[Dict[str, Any]]:
        """获取会话信息"""
        session_data = self.redis_client.get(session_id)
        if not session_data:
            return None
        
        # 更新最后活动时间
        session = json.loads(session_data)
        session["last_activity"] = time.time()
        self.redis_client.setex(
            session_id,
            self.session_ttl,
            json.dumps(session, ensure_ascii=False)
        )
        
        return session
    
    def update_session(self, session_id: str, updates: Dict[str, Any]) -> bool:
        """更新会话信息"""
        session = self.get_session(session_id)
        if not session:
            return False
        
        # 更新会话数据
        session.update(updates)
        session["last_activity"] = time.time()
        
        self.redis_client.setex(
            session_id,
            self.session_ttl,
            json.dumps(session, ensure_ascii=False)
        )
        
        return True
    
    def add_message(self, session_id: str, message: Dict[str, Any]) -> bool:
        """添加消息到会话"""
        session = self.get_session(session_id)
        if not session:
            return False
        
        # 添加消息
        if "messages" not in session:
            session["messages"] = []
        
        session["messages"].append({
            "timestamp": time.time(),
            "role": message.get("role"),
            "content": message.get("content"),
            "message_type": message.get("type", "text")
        })
        
        # 限制消息数量,避免内存占用过大
        if len(session["messages"]) > 50:
            session["messages"] = session["messages"][-50:]
        
        return self.update_session(session_id, session)
    
    def set_context(self, session_id: str, context: Dict[str, Any]) -> bool:
        """设置会话上下文"""
        return self.update_session(session_id, {"context": context})
    
    def get_context(self, session_id: str) -> Dict[str, Any]:
        """获取会话上下文"""
        session = self.get_session(session_id)
        if not session:
            return {}
        
        return session.get("context", {})
    
    def end_session(self, session_id: str) -> bool:
        """结束会话"""
        return bool(self.redis_client.delete(session_id))

# 示例用法
session_manager = SessionManager()

def handle_user_message(user_id: str, channel: str, message: str):
    """处理用户消息"""
    # 查找或创建会话
    session_id = find_or_create_session(user_id, channel)
    
    # 添加用户消息
    session_manager.add_message(session_id, {
        "role": "user",
        "content": message,
        "type": "text"
    })
    
    # 处理消息...
    
    # 添加系统回复
    session_manager.add_message(session_id, {
        "role": "assistant",
        "content": "这是系统回复",
        "type": "text"
    })

def find_or_create_session(user_id: str, channel: str) -> str:
    """查找或创建会话"""
    # 实际实现中应根据业务逻辑查找现有会话
    # 这里简化处理,直接创建新会话
    return session_manager.create_session(user_id, channel)

3.2 智能体核心实现

3.2.1 意图识别智能体
# 意图识别智能体
from langchain_openai import ChatOpenAI
from langchain.prompts import PromptTemplate
from langchain_core.output_parsers import JsonOutputParser
from pydantic import BaseModel, Field
from typing import List, Optional

class IntentResult(BaseModel):
    intent: str = Field(..., description="客户意图类型")
    confidence: float = Field(..., description="意图识别置信度")
    entities: List[dict] = Field(default_factory=list, description="提取的实体")
    sentiment: str = Field(default="neutral", description="客户情绪")
    need_human: bool = Field(default=False, description="是否需要人工客服")

class IntentRecognitionAgent:
    def __init__(self):
        self.llm = ChatOpenAI(
            model="gpt-4o",
            temperature=0.1,
            max_tokens=500
        )
        
        self.prompt_template = PromptTemplate(
            template="""你是一个专业的智能客服意图识别系统,请分析以下客户消息,识别客户的意图和情绪。

客户消息:{message}

请按照以下格式输出JSON格式的结果:
{{
  "intent": "意图类型",
  "confidence": 置信度,
  "entities": [
    {{
      "type": "实体类型",
      "value": "实体值"
    }}
  ],
  "sentiment": "情绪类型",
  "need_human": 是否需要人工客服
}}

意图类型包括:
- 产品咨询
- 订单查询
- 故障报修
- 投诉建议
- 业务办理
- 其他

情绪类型包括:
- positive(积极)
- neutral(中性)
- negative(消极)

如果客户消息中包含以下内容,need_human 应设为 true:
- 明确要求转接人工客服
- 情绪非常激动或愤怒
- 涉及复杂的个人问题
- 智能系统无法处理的特殊情况

请确保输出格式正确,只包含JSON内容,不要添加其他文本。""",
            input_variables=["message"]
        )
        
        self.parser = JsonOutputParser(pydantic_object=IntentResult)
        self.chain = self.prompt_template | self.llm | self.parser
    
    async def analyze_intent(self, message: dict) -> dict:
        """分析客户意图"""
        content = message.get("content", "")
        
        try:
            result = await self.chain.ainvoke({"message": content})
            
            # 构建响应
            response = {
                "type": "intent_analysis",
                "intent": result.intent,
                "confidence": result.confidence,
                "entities": result.entities,
                "sentiment": result.sentiment,
                "need_human": result.need_human,
                "session_id": message.get("session_id"),
                "user_id": message.get("user_id")
            }
            
            return response
        except Exception as e:
            # 错误处理
            return {
                "type": "error",
                "message": f"意图识别失败: {str(e)}",
                "intent": "其他",
                "confidence": 0.5,
                "entities": [],
                "sentiment": "neutral",
                "need_human": False
            }

# 示例用法
async def test_intent_agent():
    agent = IntentRecognitionAgent()
    message = {
        "content": "你好,我想查询我昨天下单的订单状态,订单号是123456789",
        "user_id": "user123",
        "session_id": "session123"
    }
    result = await agent.analyze_intent(message)
    print(result)

if __name__ == "__main__":
    import asyncio
    asyncio.run(test_intent_agent())
3.2.2 分流智能体
# 分流智能体
from langgraph.graph import StateGraph, END
from langgraph.graph.state import CompiledStateGraph
from typing import Dict, Any, Optional

class RouterState(Dict[str, Any]):
    """分流智能体状态"""
    pass

def route_message(state: RouterState) -> str:
    """根据意图分流消息"""
    intent = state.get("intent", "")
    need_human = state.get("need_human", False)
    sentiment = state.get("sentiment", "neutral")
    
    # 优先处理需要人工客服的情况
    if need_human:
        return "human_agent"
    
    # 根据意图分流
    intent_route_map = {
        "产品咨询": "qa_agent",
        "订单查询": "qa_agent",
        "故障报修": "qa_agent",
        "投诉建议": "emotion_agent",
        "业务办理": "business_agent",
        "其他": "qa_agent"
    }
    
    # 情绪特别消极时,优先转到情绪管理智能体
    if sentiment == "negative" and intent not in ["业务办理"]:
        return "emotion_agent"
    
    return intent_route_map.get(intent, "qa_agent")

def handle_qa_agent(state: RouterState) -> RouterState:
    """问答智能体处理"""
    # 实际实现中应调用问答智能体
    from agents.qa_agent import QAAgent
    
    agent = QAAgent()
    result = agent.process_message(state)
    
    state["response"] = result
    state["agent"] = "qa_agent"
    return state

def handle_business_agent(state: RouterState) -> RouterState:
    """业务办理智能体处理"""
    # 实际实现中应调用业务办理智能体
    from agents.business_agent import BusinessAgent
    
    agent = BusinessAgent()
    result = agent.process_message(state)
    
    state["response"] = result
    state["agent"] = "business_agent"
    return state

def handle_emotion_agent(state: RouterState) -> RouterState:
    """情绪管理智能体处理"""
    # 实际实现中应调用情绪管理智能体
    from agents.emotion_agent import EmotionAgent
    
    agent = EmotionAgent()
    result = agent.process_message(state)
    
    state["response"] = result
    state["agent"] = "emotion_agent"
    return state

def handle_human_agent(state: RouterState) -> RouterState:
    """人工协作智能体处理"""
    # 实际实现中应调用人工协作智能体
    from agents.human_agent import HumanAgent
    
    agent = HumanAgent()
    result = agent.process_message(state)
    
    state["response"] = result
    state["agent"] = "human_agent"
    return state

def create_router_graph() -> CompiledStateGraph:
    """创建分流智能体图"""
    graph = StateGraph(RouterState)
    
    # 添加节点
    graph.add_node("route", route_message)
    graph.add_node("qa_agent", handle_qa_agent)
    graph.add_node("business_agent", handle_business_agent)
    graph.add_node("emotion_agent", handle_emotion_agent)
    graph.add_node("human_agent", handle_human_agent)
    
    # 添加边
    graph.add_conditional_edges(
        "route",
        lambda state: state,
        {
            "qa_agent": "qa_agent",
            "business_agent": "business_agent",
            "emotion_agent": "emotion_agent",
            "human_agent": "human_agent"
        }
    )
    
    # 所有智能体处理完成后结束
    graph.add_edge("qa_agent", END)
    graph.add_edge("business_agent", END)
    graph.add_edge("emotion_agent", END)
    graph.add_edge("human_agent", END)
    
    # 设置入口点
    graph.set_entry_point("route")
    
    return graph.compile()

# 示例用法
async def test_router_agent():
    """测试分流智能体"""
    graph = create_router_graph()
    
    # 测试订单查询
    state = {
        "content": "我想查询我的订单状态",
        "intent": "订单查询",
        "confidence": 0.95,
        "entities": [],
        "sentiment": "neutral",
        "need_human": False,
        "user_id": "user123",
        "session_id": "session123"
    }
    
    result = await graph.ainvoke(state)
    print("订单查询测试结果:", result)
    
    # 测试业务办理
    state = {
        "content": "我想办理会员升级",
        "intent": "业务办理",
        "confidence": 0.9,
        "entities": [],
        "sentiment": "neutral",
        "need_human": False,
        "user_id": "user123",
        "session_id": "session123"
    }
    
    result = await graph.ainvoke(state)
    print("业务办理测试结果:", result)
    
    # 测试需要人工客服的情况
    state = {
        "content": "我要投诉,转接人工客服",
        "intent": "投诉建议",
        "confidence": 0.95,
        "entities": [],
        "sentiment": "negative",
        "need_human": True,
        "user_id": "user123",
        "session_id": "session123"
    }
    
    result = await graph.ainvoke(state)
    print("人工客服测试结果:", result)

if __name__ == "__main__":
    import asyncio
    asyncio.run(test_router_agent())
3.2.3 问答智能体
# 问答智能体
from langchain_openai import ChatOpenAI, OpenAIEmbeddings
from langchain_chroma import Chroma
from langchain.prompts import PromptTemplate
from langchain.chains import RetrievalQA
from typing import Dict, Any

class QAAgent:
    def __init__(self):
        # 初始化嵌入模型
        self.embeddings = OpenAIEmbeddings(model="text-embedding-3-large")
        
        # 初始化向量数据库
        self.vector_db = Chroma(
            persist_directory="./vector_db",
            embedding_function=self.embeddings
        )
        
        # 初始化LLM
        self.llm = ChatOpenAI(
            model="gpt-4o",
            temperature=0.1
        )
        
        # 初始化检索器
        self.retriever = self.vector_db.as_retriever(
            search_kwargs={"k": 5}
        )
        
        # 初始化提示模板
        self.prompt_template = PromptTemplate(
            template="""你是一个专业的客服智能助手,请根据以下上下文信息回答客户的问题。

上下文信息:
{context}

客户问题:
{question}

请遵循以下要求:
1. 基于上下文信息回答问题,不要添加无关内容
2. 回答要准确、简洁、专业
3. 如果上下文信息不足以回答问题,请明确说明
4. 保持友好的语气和专业的态度

回答:""",
            input_variables=["context", "question"]
        )
        
        # 初始化问答链
        self.qa_chain = RetrievalQA.from_chain_type(
            llm=self.llm,
            chain_type="stuff",
            retriever=self.retriever,
            chain_type_kwargs={
                "prompt": self.prompt_template
            },
            return_source_documents=True
        )
    
    def process_message(self, state: Dict[str, Any]) -> Dict[str, Any]:
        """处理问答消息"""
        question = state.get("content", "")
        
        try:
            result = self.qa_chain.invoke({"query": question})
            
            # 构建响应
            response = {
                "type": "answer",
                "content": result["result"],
                "sources": [doc.metadata for doc in result.get("source_documents", [])],
                "confidence": self._calculate_confidence(result),
                "follow_up": self._generate_follow_up(question, result["result"])
            }
            
            return response
        except Exception as e:
            # 错误处理
            return {
                "type": "error",
                "content": "抱歉,我暂时无法回答这个问题,请稍后再试。",
                "error": str(e)
            }
    
    def _calculate_confidence(self, result: Dict[str, Any]) -> float:
        """计算回答的置信度"""
        # 实际实现中应基于检索结果和模型输出计算置信度
        # 这里简化处理,返回固定值
        return 0.85
    
    def _generate_follow_up(self, question: str, answer: str) -> list:
        """生成后续问题建议"""
        # 实际实现中应基于问题和回答生成后续问题建议
        # 这里简化处理,返回空列表
        return []
    
    def add_to_knowledge_base(self, documents: list):
        """添加文档到知识库"""
        # 实际实现中应处理文档的嵌入和存储
        self.vector_db.add_documents(documents)
        self.vector_db.persist()

# 示例用法
def test_qa_agent():
    """测试问答智能体"""
    agent = QAAgent()
    
    # 测试问题
    state = {
        "content": "如何修改密码?",
        "user_id": "user123",
        "session_id": "session123"
    }
    
    result = agent.process_message(state)
    print("问答智能体测试结果:")
    print(f"回答: {result['content']}")
    print(f"来源: {result.get('sources', [])}")
    print(f"置信度: {result.get('confidence', 0)}")

if __name__ == "__main__":
    test_qa_agent()
3.2.3 问答智能体
# 问答智能体
from langchain_openai import ChatOpenAI, OpenAIEmbeddings
from langchain_chroma import Chroma
from langchain.prompts import PromptTemplate
from langchain.chains import RetrievalQA
from typing import Dict, Any

class QAAgent:
    def __init__(self):
        # 初始化嵌入模型
        self.embeddings = OpenAIEmbeddings(model="text-embedding-3-large")
        
        # 初始化向量数据库
        self.vector_db = Chroma(
            persist_directory="./vector_db",
            embedding_function=self.embeddings
        )
        
        # 初始化LLM
        self.llm = ChatOpenAI(
            model="gpt-4o",
            temperature=0.1
        )
        
        # 初始化检索器
        self.retriever = self.vector_db.as_retriever(
            search_kwargs={"k": 5}
        )
        
        # 初始化提示模板
        self.prompt_template = PromptTemplate(
            template="""你是一个专业的客服智能助手,请根据以下上下文信息回答客户的问题。

上下文信息:
{context}

客户问题:
{question}

请遵循以下要求:
1. 基于上下文信息回答问题,不要添加无关内容
2. 回答要准确、简洁、专业
3. 如果上下文信息不足以回答问题,请明确说明
4. 保持友好的语气和专业的态度

回答:""",
            input_variables=["context", "question"]
        )
        
        # 初始化问答链
        self.qa_chain = RetrievalQA.from_chain_type(
            llm=self.llm,
            chain_type="stuff",
            retriever=self.retriever,
            chain_type_kwargs={
                "prompt": self.prompt_template
            },
            return_source_documents=True
        )
    
    def process_message(self, state: Dict[str, Any]) -> Dict[str, Any]:
        """处理问答消息"""
        question = state.get("content", "")
        
        try:
            result = self.qa_chain.invoke({"query": question})
            
            # 构建响应
            response = {
                "type": "answer",
                "content": result["result"],
                "sources": [doc.metadata for doc in result.get("source_documents", [])],
                "confidence": self._calculate_confidence(result),
                "follow_up": self._generate_follow_up(question, result["result"])
            }
            
            return response
        except Exception as e:
            # 错误处理
            return {
                "type": "error",
                "content": "抱歉,我暂时无法回答这个问题,请稍后再试。",
                "error": str(e)
            }
    
    def _calculate_confidence(self, result: Dict[str, Any]) -> float:
        """计算回答的置信度"""
        # 实际实现中应基于检索结果和模型输出计算置信度
        # 这里简化处理,返回固定值
        return 0.85
    
    def _generate_follow_up(self, question: str, answer: str) -> list:
        """生成后续问题建议"""
        # 实际实现中应基于问题和回答生成后续问题建议
        # 这里简化处理,返回空列表
        return []
    
    def add_to_knowledge_base(self, documents: list):
        """添加文档到知识库"""
        # 实际实现中应处理文档的嵌入和存储
        self.vector_db.add_documents(documents)
        self.vector_db.persist()

# 示例用法
def test_qa_agent():
    """测试问答智能体"""
    agent = QAAgent()
    
    # 测试问题
    state = {
        "content": "如何修改密码?",
        "user_id": "user123",
        "session_id": "session123"
    }
    
    result = agent.process_message(state)
    print("问答智能体测试结果:")
    print(f"回答: {result['content']}")
    print(f"来源: {result.get('sources', [])}")
    print(f"置信度: {result.get('confidence', 0)}")

if __name__ == "__main__":
    test_qa_agent()
3.2.4 业务办理智能体
# 业务办理智能体
from langchain_openai import ChatOpenAI
from langchain.prompts import PromptTemplate
from langchain.tools import tool
from langchain_core.tools import Tool
from langchain.agents import AgentExecutor, create_tool_calling_agent
from typing import Dict, Any, List

class BusinessAgent:
    def __init__(self):
        self.llm = ChatOpenAI(
            model="gpt-4o",
            temperature=0.1
        )
        
        # 定义工具
        self.tools = [
            self._create_user_info_tool(),
            self._create_order_tool(),
            self._create_service_tool()
        ]
        
        # 初始化提示模板
        self.prompt = PromptTemplate(
            template="""你是一个业务办理智能助手,负责帮助客户办理各种业务。

当前对话:
{chat_history}
客户请求:
{input}
请根据客户的请求,决定是否需要使用工具来获取更多信息,或者直接回答客户的问题。
可用工具:
{tools}
工具使用格式:
{{
  "tool_call": {{
    "name": "工具名称",
    "arguments": {{
      "参数名": "参数值"
    }}
  }}
}}

直接回答格式:

{{
  "direct_answer": "你的回答"
}}


请确保输出格式正确,只包含JSON内容。""",
        input_variables=["chat_history", "input", "tools"]
        
        
        # 初始化智能体
        self.agent = create_tool_calling_agent(
            llm=self.llm,
            tools=self.tools,
            prompt=self.prompt
        )
        
        self.agent_executor = AgentExecutor(
            agent=self.agent,
            tools=self.tools,
            verbose=True
        )
    
    def _create_user_info_tool(self) -> Tool:
        """创建用户信息工具"""
        @tool
        def get_user_info(user_id: str) -> Dict[str, Any]:
            """获取用户信息
            
            Args:
                user_id: 用户ID
                
            Returns:
                用户信息字典,包含姓名、手机号、会员等级等信息
            """
            # 实际实现中应从数据库或API获取用户信息
            return {
                "user_id": user_id,
                "name": "张三",
                "phone": "13800138000",
                "member_level": "黄金会员",
                "points": 1000
            }
        
        return get_user_info
    
    def _create_order_tool(self) -> Tool:
        """创建订单工具"""
        @tool
        def get_order_info(order_id: str) -> Dict[str, Any]:
            """获取订单信息
            
            Args:
                order_id: 订单ID
                
            Returns:
                订单信息字典,包含订单状态、商品信息、物流信息等
            """
            # 实际实现中应从数据库或API获取订单信息
            return {
                "order_id": order_id,
                "status": "已发货",
                "items": [{
                    "name": "商品1",
                    "quantity": 1,
                    "price": 100
                }],
                "shipping_info": {
                    "company": "顺丰速运",
                    "tracking_number": "SF1234567890"
                },
                "total_amount": 100
            }
        
        return get_order_info
    
    def _create_service_tool(self) -> Tool:
        """创建服务工具"""
        @tool
        def process_service_request(service_type: str, user_id: str, details: Dict[str, Any]) -> Dict[str, Any]:
            """处理服务请求
            
            Args:
                service_type: 服务类型,如 "密码重置"、"会员升级" 等
                user_id: 用户ID
                details: 服务详情
                
            Returns:
                服务处理结果
            """
            # 实际实现中应调用相应的服务API
            return {
                "status": "success",
                "message": f"{service_type} 处理成功",
                "service_id": f"SVC{int(time.time())}",
                "user_id": user_id
            }
        
        return process_service_request
    
    def process_message(self, state: Dict[str, Any]) -> Dict[str, Any]:
        """处理业务办理消息"""
        input_text = state.get("content", "")
        chat_history = state.get("chat_history", "")
        
        try:
            result = self.agent_executor.invoke({
                "input": input_text,
                "chat_history": chat_history
            })
            
            # 构建响应
            response = {
                "type": "service_response",
                "content": result["output"],
                "status": "completed" if "成功" in result["output"] else "pending"
            }
            
            return response
        except Exception as e:
            # 错误处理
            return {
                "type": "error",
                "content": "抱歉,业务办理过程中出现错误,请稍后再试。",
                "error": str(e)
            }

# 示例用法
import time
def test_business_agent():
    """测试业务办理智能体"""
    agent = BusinessAgent()
    
    # 测试订单查询
    state = {
        "content": "查询订单号123456789的状态",
        "user_id": "user123",
        "session_id": "session123"
    }
    
    result = agent.process_message(state)
    print("业务办理智能体测试结果:")
    print(f"响应: {result['content']}")
    print(f"状态: {result['status']}")

if __name__ == "__main__":
    test_business_agent()

3.3 知识库管理

# 知识库管理
from langchain_community.document_loaders import TextLoader, PDFLoader, Docx2txtLoader
from langchain.text_splitter import RecursiveCharacterTextSplitter
from langchain_openai import OpenAIEmbeddings
from langchain_chroma import Chroma
from typing import List, Dict, Any
import os

class KnowledgeBaseManager:
    def __init__(self, persist_directory="./vector_db"):
        self.persist_directory = persist_directory
        self.embeddings = OpenAIEmbeddings(model="text-embedding-3-large")
        self.text_splitter = RecursiveCharacterTextSplitter(
            chunk_size=1000,
            chunk_overlap=200
        )
        self.vector_db = Chroma(
            persist_directory=persist_directory,
            embedding_function=self.embeddings
        )
    
    def add_document(self, file_path: str, metadata: Dict[str, Any] = None) -> int:
        """添加文档到知识库"""
        try:
            # 根据文件类型选择加载器
            if file_path.endswith('.txt'):
                loader = TextLoader(file_path)
            elif file_path.endswith('.pdf'):
                loader = PDFLoader(file_path)
            elif file_path.endswith('.docx'):
                loader = Docx2txtLoader(file_path)
            else:
                raise ValueError(f"不支持的文件类型: {file_path}")
            
            # 加载文档
            documents = loader.load()
            
            # 分割文档
            splits = self.text_splitter.split_documents(documents)
            
            # 添加元数据
            if metadata:
                for split in splits:
                    split.metadata.update(metadata)
            
            # 添加到向量数据库
            self.vector_db.add_documents(splits)
            self.vector_db.persist()
            
            return len(splits)
        except Exception as e:
            print(f"添加文档失败: {str(e)}")
            return 0
    
    def add_text(self, text: str, metadata: Dict[str, Any] = None) -> bool:
        """添加文本到知识库"""
        try:
            # 创建文档对象
            from langchain_core.documents import Document
            
            document = Document(page_content=text, metadata=metadata or {})
            
            # 分割文本
            splits = self.text_splitter.split_documents([document])
            
            # 添加到向量数据库
            self.vector_db.add_documents(splits)
            self.vector_db.persist()
            
            return True
        except Exception as e:
            print(f"添加文本失败: {str(e)}")
            return False
    
    def search(self, query: str, k: int = 5) -> List[Dict[str, Any]]:
        """搜索知识库"""
        try:
            results = self.vector_db.similarity_search_with_score(query, k=k)
            
            # 格式化结果
            formatted_results = []
            for doc, score in results:
                formatted_results.append({
                    "content": doc.page_content,
                    "metadata": doc.metadata,
                    "score": score
                })
            
            return formatted_results
        except Exception as e:
            print(f"搜索失败: {str(e)}")
            return []
    
    def delete_document(self, document_id: str) -> bool:
        """删除文档"""
        try:
            # 实际实现中应根据文档ID删除
            # 这里简化处理
            return True
        except Exception as e:
            print(f"删除文档失败: {str(e)}")
            return False
    
    def clear_knowledge_base(self) -> bool:
        """清空知识库"""
        try:
            # 删除向量数据库文件
            import shutil
            if os.path.exists(self.persist_directory):
                shutil.rmtree(self.persist_directory)
            
            # 重新初始化
            self.vector_db = Chroma(
                persist_directory=self.persist_directory,
                embedding_function=self.embeddings
            )
            
            return True
        except Exception as e:
            print(f"清空知识库失败: {str(e)}")
            return False
    
    def get_statistics(self) -> Dict[str, Any]:
        """获取知识库统计信息"""
        try:
            # 实际实现中应获取更详细的统计信息
            return {
                "persist_directory": self.persist_directory,
                "document_count": self.vector_db._collection.count()
            }
        except Exception as e:
            print(f"获取统计信息失败: {str(e)}")
            return {}

# 示例用法
def test_knowledge_base():
    """测试知识库管理"""
    manager = KnowledgeBaseManager()
    
    # 添加文档
    print("添加文档...")
    count = manager.add_document("./docs/faq.txt", {
        "source": "faq",
        "category": "客户服务",
        "update_date": "2024-01-01"
    })
    print(f"添加了 {count} 个文档片段")
    
    # 添加文本
    print("添加文本...")
    success = manager.add_text(
        "如何修改密码?\n1. 登录账户\n2. 进入个人中心\n3. 点击密码修改\n4. 按照提示操作",
        {
            "source": "manual",
            "category": "账户管理",
            "update_date": "2024-01-01"
        }
    )
    print(f"添加文本成功: {success}")
    
    # 搜索
    print("搜索知识库...")
    results = manager.search("修改密码")
    print(f"搜索结果: {len(results)} 条")
    for i, result in enumerate(results):
        print(f"结果 {i+1}: 相似度 {result['score']:.4f}")
        print(f"内容: {result['content'][:100]}...")
        print(f"来源: {result['metadata'].get('source')}")
    
    # 获取统计信息
    print("获取统计信息...")
    stats = manager.get_statistics()
    print(f"知识库统计: {stats}")

if __name__ == "__main__":
    test_knowledge_base()

3.4 人工协作实现

# 人工协作智能体
import asyncio
import json
from typing import Dict, Any, List

class HumanAgent:
    def __init__(self):
        self.available_agents = []  # 可用的人工客服
        self.queue = []  # 等待队列
        self.active_sessions = {}  # 活跃会话
    
    def process_message(self, state: Dict[str, Any]) -> Dict[str, Any]:
        """处理需要人工客服的消息"""
        session_id = state.get("session_id")
        user_id = state.get("user_id")
        content = state.get("content", "")
        
        # 检查是否已有活跃的人工会话
        if session_id in self.active_sessions:
            # 已有活跃会话,直接转发消息
            agent_id = self.active_sessions[session_id]
            self._forward_to_agent(session_id, agent_id, content)
            
            return {
                "type": "human_agent",
                "content": "您的消息已转发给人工客服,请稍候...",
                "status": "forwarded",
                "agent_id": agent_id
            }
        else:
            # 没有活跃会话,尝试分配客服
            agent_id = self._assign_agent()
            
            if agent_id:
                # 分配成功,创建会话
                self.active_sessions[session_id] = agent_id
                self._forward_to_agent(session_id, agent_id, content)
                
                return {
                    "type": "human_agent",
                    "content": "正在为您转接人工客服,请稍候...",
                    "status": "assigned",
                    "agent_id": agent_id
                }
            else:
                # 没有可用客服,加入等待队列
                self.queue.append({
                    "session_id": session_id,
                    "user_id": user_id,
                    "content": content,
                    "timestamp": time.time()
                })
                
                return {
                    "type": "human_agent",
                    "content": "当前人工客服繁忙,请耐心等待,我们会尽快为您服务。",
                    "status": "queued",
                    "queue_position": len(self.queue)
                }
    
    def _assign_agent(self) -> str:
        """分配人工客服"""
        # 实际实现中应基于客服的工作量、技能等因素分配
        # 这里简化处理,返回第一个可用客服
        if self.available_agents:
            return self.available_agents[0]
        return "agent_001"  # 模拟返回一个客服ID
    
    def _forward_to_agent(self, session_id: str, agent_id: str, message: str):
        """转发消息给人工客服"""
        # 实际实现中应通过WebSocket或消息队列转发消息
        print(f"转发消息到客服 {agent_id}: {message}")
    
    def agent_respond(self, session_id: str, agent_id: str, response: str):
        """人工客服回复"""
        # 实际实现中应将回复发送给用户
        print(f"客服 {agent_id} 回复: {response}")
        
        # 检查会话是否存在
        if session_id in self.active_sessions:
            # 可以在这里添加回复的处理逻辑
            pass
    
    def agent_available(self, agent_id: str):
        """客服可用"""
        if agent_id not in self.available_agents:
            self.available_agents.append(agent_id)
        
        # 检查是否有等待的会话
        if self.queue:
            # 分配第一个等待的会话
            session_info = self.queue.pop(0)
            self.active_sessions[session_info["session_id"]] = agent_id
            self._forward_to_agent(
                session_info["session_id"],
                agent_id,
                session_info["content"]
            )
    
    def agent_unavailable(self, agent_id: str):
        """客服不可用"""
        if agent_id in self.available_agents:
            self.available_agents.remove(agent_id)
    
    def end_session(self, session_id: str):
        """结束会话"""
        if session_id in self.active_sessions:
            agent_id = self.active_sessions[session_id]
            del self.active_sessions[session_id]
            
            # 可以在这里添加会话结束的处理逻辑
            print(f"会话 {session_id} 已结束,客服 {agent_id} 已释放")

# 示例用法
import time
def test_human_agent():
    """测试人工协作智能体"""
    agent = HumanAgent()
    
    # 模拟客服上线
    agent.agent_available("agent_001")
    agent.agent_available("agent_002")
    
    # 测试消息
    state = {
        "content": "我要投诉,转接人工客服",
        "user_id": "user123",
        "session_id": "session123"
    }
    
    result = agent.process_message(state)
    print("人工协作智能体测试结果:")
    print(f"响应: {result['content']}")
    print(f"状态: {result['status']}")
    print(f"客服ID: {result.get('agent_id')}")
    
    # 模拟客服回复
    time.sleep(1)
    agent.agent_respond("session123", "agent_001", "您好,我是客服小李,请问有什么可以帮助您的?")
    
    # 模拟结束会话
    time.sleep(2)
    agent.end_session("session123")

if __name__ == "__main__":
    test_human_agent()

4. 系统集成与部署

4.1 容器化部署

4.1.1 Docker Compose配置
version: '3.8'

services:
  # 智能客服API服务
  api-service:
    build: ./api-service
    restart: unless-stopped
    ports:
      - "8000:8000"
    volumes:
      - ./api-service:/app
    environment:
      - OPENAI_API_KEY=${OPENAI_API_KEY}
      - REDIS_URL=redis://redis:6379
      - DATABASE_URL=postgresql://admin:password@postgres:5432/example_db
      - CHROMA_DB_PATH=/app/vector_db
    depends_on:
      - redis
      - postgres
    networks:
      - agent-net

  # Redis缓存
  redis:
    image: redis:7
    restart: unless-stopped
    volumes:
      - redis-data:/data
    networks:
      - agent-net

  # PostgreSQL数据库
  postgres:
    image: postgres:15
    restart: unless-stopped
    volumes:
      - postgres-data:/var/lib/postgresql/data
    environment:
      - POSTGRES_USER=admin
      - POSTGRES_PASSWORD=password
      - POSTGRES_DB=example_db
    networks:
      - agent-net

  # 向量数据库
  chroma:
    build: ./chroma
    restart: unless-stopped
    volumes:
      - chroma-data:/app/chroma_db
    networks:
      - agent-net

  # 监控服务
  prometheus:
    image: prom/prometheus:latest
    restart: unless-stopped
    volumes:
      - ./prometheus/prometheus.yml:/etc/prometheus/prometheus.yml
    ports:
      - "9090:9090"
    networks:
      - agent-net

  # 可视化服务
  grafana:
    image: grafana/grafana:latest
    restart: unless-stopped
    ports:
      - "3000:3000"
    volumes:
      - grafana-data:/var/lib/grafana
    networks:
      - agent-net

volumes:
  redis-data:
  postgres-data:
  chroma-data:
  grafana-data:

networks:
  agent-net:
    driver: bridge
4.1.2 Kubernetes部署
# 智能客服API服务部署
apiVersion: apps/v1
kind: Deployment
metadata:
  name: api-service
  namespace: smart-customer-service
spec:
  replicas: 3
  selector:
    matchLabels:
      app: api-service
  template:
    metadata:
      labels:
        app: api-service
    spec:
      containers:
      - name: api-service
        image: your-registry/api-service:latest
        ports:
        - containerPort: 8000
        env:
        - name: OPENAI_API_KEY
          valueFrom:
            secretKeyRef:
              name: api-keys
              key: openai-api-key
        - name: REDIS_URL
          value: redis://redis:6379
        - name: DATABASE_URL
          value: postgresql://admin:password@postgres:5432/example_db
        - name: CHROMA_DB_PATH
          value: /app/vector_db
        resources:
          requests:
            cpu: "500m"
            memory: "1Gi"
          limits:
            cpu: "1"
            memory: "2Gi"

---
# 智能客服API服务Service
apiVersion: v1
kind: Service
metadata:
  name: api-service
  namespace: smart-customer-service
spec:
  selector:
    app: api-service
  ports:
  - port: 8000
    targetPort: 8000
  type: LoadBalancer

---
# Redis部署
apiVersion: apps/v1
kind: Deployment
metadata:
  name: redis
  namespace: smart-customer-service
spec:
  replicas: 1
  selector:
    matchLabels:
      app: redis
  template:
    metadata:
      labels:
        app: redis
    spec:
      containers:
      - name: redis
        image: redis:7
        ports:
        - containerPort: 6379
        volumeMounts:
        - name: redis-data
          mountPath: /data
  volumes:
  - name: redis-data
    persistentVolumeClaim:
      claimName: redis-pvc

---
# Redis Service
apiVersion: v1
kind: Service
metadata:
  name: redis
  namespace: smart-customer-service
spec:
  selector:
    app: redis
  ports:
  - port: 6379
    targetPort: 6379

---
# PostgreSQL部署
apiVersion: apps/v1
kind: Deployment
metadata:
  name: postgres
  namespace: smart-customer-service
spec:
  replicas: 1
  selector:
    matchLabels:
      app: postgres
  template:
    metadata:
      labels:
        app: postgres
    spec:
      containers:
      - name: postgres
        image: postgres:15
        ports:
        - containerPort: 5432
        env:
        - name: POSTGRES_USER
          value: admin
        - name: POSTGRES_PASSWORD
          value: password
        - name: POSTGRES_DB
          value: example_db
        volumeMounts:
        - name: postgres-data
          mountPath: /var/lib/postgresql/data
  volumes:
  - name: postgres-data
    persistentVolumeClaim:
      claimName: postgres-pvc

---
# PostgreSQL Service
apiVersion: v1
kind: Service
metadata:
  name: postgres
  namespace: smart-customer-service
spec:
  selector:
    app: postgres
  ports:
  - port: 5432
    targetPort: 5432

4.2 系统集成

4.2.1 与业务系统集成
# 业务系统集成服务
import requests
import json
from typing import Dict, Any, Optional

class BusinessSystemIntegration:
    def __init__(self):
        self.services = {
            "order": {
                "base_url": "http://order-service:8000",
                "endpoints": {
                    "get_order": "/api/orders/{order_id}",
                    "create_order": "/api/orders",
                    "cancel_order": "/api/orders/{order_id}/cancel"
                }
            },
            "user": {
                "base_url": "http://user-service:8000",
                "endpoints": {
                    "get_user": "/api/users/{user_id}",
                    "update_user": "/api/users/{user_id}",
                    "reset_password": "/api/users/{user_id}/reset-password"
                }
            },
            "payment": {
                "base_url": "http://payment-service:8000",
                "endpoints": {
                    "create_payment": "/api/payments",
                    "get_payment": "/api/payments/{payment_id}"
                }
            }
        }
    
    def get_order_info(self, order_id: str) -> Optional[Dict[str, Any]]:
        """获取订单信息"""
        try:
            service = self.services["order"]
            endpoint = service["endpoints"]["get_order"].format(order_id=order_id)
            url = f"{service['base_url']}{endpoint}"
            
            response = requests.get(url, timeout=10)
            response.raise_for_status()
            
            return response.json()
        except Exception as e:
            print(f"获取订单信息失败: {str(e)}")
            return None
    
    def cancel_order(self, order_id: str, user_id: str) -> Optional[Dict[str, Any]]:
        """取消订单"""
        try:
            service = self.services["order"]
            endpoint = service["endpoints"]["cancel_order"].format(order_id=order_id)
            url = f"{service['base_url']}{endpoint}"
            
            response = requests.post(url, json={"user_id": user_id}, timeout=10)
            response.raise_for_status()
            
            return response.json()
        except Exception as e:
            print(f"取消订单失败: {str(e)}")
            return None
    
    def get_user_info(self, user_id: str) -> Optional[Dict[str, Any]]:
        """获取用户信息"""
        try:
            service = self.services["user"]
            endpoint = service["endpoints"]["get_user"].format(user_id=user_id)
            url = f"{service['base_url']}{endpoint}"
            
            response = requests.get(url, timeout=10)
            response.raise_for_status()
            
            return response.json()
        except Exception as e:
            print(f"获取用户信息失败: {str(e)}")
            return None
    
    def reset_password(self, user_id: str, phone: str) -> Optional[Dict[str, Any]]:
        """重置密码"""
        try:
            service = self.services["user"]
            endpoint = service["endpoints"]["reset_password"].format(user_id=user_id)
            url = f"{service['base_url']}{endpoint}"
            
            response = requests.post(url, json={"phone": phone}, timeout=10)
            response.raise_for_status()
            
            return response.json()
        except Exception as e:
            print(f"重置密码失败: {str(e)}")
            return None
    
    def create_payment(self, order_id: str, user_id: str, amount: float) -> Optional[Dict[str, Any]]:
        """创建支付"""
        try:
            service = self.services["payment"]
            endpoint = service["endpoints"]["create_payment"]
            url = f"{service['base_url']}{endpoint}"
            
            payload = {
                "order_id": order_id,
                "user_id": user_id,
                "amount": amount,
                "payment_method": "online"
            }
            
            response = requests.post(url, json=payload, timeout=10)
            response.raise_for_status()
            
            return response.json()
        except Exception as e:
            print(f"创建支付失败: {str(e)}")
            return None

# 示例用法
def test_business_integration():
    """测试业务系统集成"""
    integration = BusinessSystemIntegration()
    
    # 测试获取订单信息
    order_info = integration.get_order_info("123456789")
    print("订单信息:", order_info)
    
    # 测试获取用户信息
    user_info = integration.get_user_info("user123")
    print("用户信息:", user_info)

if __name__ == "__main__":
    test_business_integration()
4.2.2 与第三方服务集成
# 第三方服务集成
import requests
import json
from typing import Dict, Any, Optional

class ThirdPartyIntegration:
    def __init__(self):
        self.services = {
            "sms": {
                "base_url": "http://sms-service:8000",
                "endpoints": {
                    "send_sms": "/api/sms/send"
                }
            },
            "email": {
                "base_url": "http://email-service:8000",
                "endpoints": {
                    "send_email": "/api/email/send"
                }
            },
            "notification": {
                "base_url": "http://notification-service:8000",
                "endpoints": {
                    "send_notification": "/api/notifications/send"
                }
            }
        }
    
    def send_sms(self, phone: str, message: str) -> bool:
        """发送短信"""
        try:
            service = self.services["sms"]
            endpoint = service["endpoints"]["send_sms"]
            url = f"{service['base_url']}{endpoint}"
            
            payload = {
                "phone": phone,
                "message": message,
                "type": "verification"
            }
            
            response = requests.post(url, json=payload, timeout=10)
            response.raise_for_status()
            
            return True
        except Exception as e:
            print(f"发送短信失败: {str(e)}")
            return False
    
    def send_email(self, email: str, subject: str, content: str) -> bool:
        """发送邮件"""
        try:
            service = self.services["email"]
            endpoint = service["endpoints"]["send_email"]
            url = f"{service['base_url']}{endpoint}"
            
            payload = {
                "email": email,
                "subject": subject,
                "content": content,
                "type": "notification"
            }
            
            response = requests.post(url, json=payload, timeout=10)
            response.raise_for_status()
            
            return True
        except Exception as e:
            print(f"发送邮件失败: {str(e)}")
            return False
    
    def send_notification(self, user_id: str, title: str, content: str, channel: str) -> bool:
        """发送通知"""
        try:
            service = self.services["notification"]
            endpoint = service["endpoints"]["send_notification"]
            url = f"{service['base_url']}{endpoint}"
            
            payload = {
                "user_id": user_id,
                "title": title,
                "content": content,
                "channel": channel
            }
            
            response = requests.post(url, json=payload, timeout=10)
            response.raise_for_status()
            
            return True
        except Exception as e:
            print(f"发送通知失败: {str(e)}")
            return False

# 示例用法
def test_third_party_integration():
    """测试第三方服务集成"""
    integration = ThirdPartyIntegration()
    
    # 测试发送短信
    success = integration.send_sms("13800138000", "您的验证码是123456")
    print(f"发送短信成功: {success}")
    
    # 测试发送邮件
    success = integration.send_email("user@example.com", "测试邮件", "这是一封测试邮件")
    print(f"发送邮件成功: {success}")

if __name__ == "__main__":
    test_third_party_integration()

5. 系统监控与运维

5.1 监控系统

# 监控系统集成
import prometheus_client
from prometheus_client import Counter, Gauge, Histogram, Summary
from fastapi import FastAPI, Request
import time

# 定义指标
REQUEST_COUNT = Counter('request_count', 'Total request count', ['endpoint', 'method', 'status'])
REQUEST_LATENCY = Histogram('request_latency_seconds', 'Request latency', ['endpoint'])
AGENT_RESPONSE_TIME = Gauge('agent_response_time_seconds', 'Agent response time', ['agent_type'])
ACTIVE_SESSIONS = Gauge('active_sessions', 'Number of active sessions')
ERROR_COUNT = Counter('error_count', 'Error count', ['error_type'])

app = FastAPI()

# 启动Prometheus指标服务
prometheus_client.start_http_server(8000)

@app.middleware("http")
async def metrics_middleware(request: Request, call_next):
    """请求指标中间件"""
    start_time = time.time()
    
    # 处理请求
    response = await call_next(request)
    
    # 记录指标
    endpoint = request.url.path
    method = request.method
    status = response.status_code
    
    REQUEST_COUNT.labels(endpoint=endpoint, method=method, status=status).inc()
    REQUEST_LATENCY.labels(endpoint=endpoint).observe(time.time() - start_time)
    
    return response

def record_agent_response_time(agent_type: str, response_time: float):
    """记录智能体响应时间"""
    AGENT_RESPONSE_TIME.labels(agent_type=agent_type).set(response_time)

def record_error(error_type: str):
    """记录错误"""
    ERROR_COUNT.labels(error_type=error_type).inc()

def update_active_sessions(count: int):
    """更新活跃会话数"""
    ACTIVE_SESSIONS.set(count)

# 示例用法
def process_user_message(message: dict):
    """处理用户消息"""
    start_time = time.time()
    
    try:
        # 处理消息...
        
        # 记录智能体响应时间
        response_time = time.time() - start_time
        record_agent_response_time("qa_agent", response_time)
        
        # 更新活跃会话数
        update_active_sessions(10)
        
    except Exception as e:
        # 记录错误
        record_error("processing_error")
        raise

if __name__ == "__main__":
    import uvicorn
    uvicorn.run(app, host="0.0.0.0", port=8001)

5.2 日志管理

# 日志管理
import logging
import logging.config
import json
import os
from datetime import datetime

# 日志配置
LOG_CONFIG = {
    "version": 1,
    "disable_existing_loggers": False,
    "formatters": {
        "standard": {
            "format": "%(asctime)s [%(levelname)s] %(name)s: %(message)s"
        },
        "json": {
            "()": "logging.Formatter",
            "format": json.dumps({
                "timestamp": "%(asctime)s",
                "level": "%(levelname)s",
                "logger": "%(name)s",
                "message": "%(message)s",
                "module": "%(module)s",
                "function": "%(funcName)s",
                "line": "%(lineno)d"
            })
        }
    },
    "handlers": {
        "console": {
            "class": "logging.StreamHandler",
            "formatter": "standard",
            "level": "INFO"
        },
        "file": {
            "class": "logging.handlers.RotatingFileHandler",
            "filename": "logs/agent-system.log",
            "maxBytes": 10485760,  # 10MB
            "backupCount": 5,
            "formatter": "json",
            "level": "INFO"
        },
        "error": {
            "class": "logging.handlers.RotatingFileHandler",
            "filename": "logs/agent-system-error.log",
            "maxBytes": 10485760,  # 10MB
            "backupCount": 5,
            "formatter": "json",
            "level": "ERROR"
        }
    },
    "loggers": {
        "": {
            "handlers": ["console", "file", "error"],
            "level": "INFO"
        },
        "agent": {
            "handlers": ["console", "file", "error"],
            "level": "DEBUG"
        }
    }
}

# 创建日志目录
os.makedirs('logs', exist_ok=True)

# 配置日志
logging.config.dictConfig(LOG_CONFIG)

# 获取日志器
logger = logging.getLogger("agent")

def log_user_message(user_id: str, session_id: str, message: str):
    """记录用户消息"""
    logger.info(
        f"User message",
        extra={
            "user_id": user_id,
            "session_id": session_id,
            "message": message,
            "event_type": "user_message"
        }
    )

def log_agent_response(agent_type: str, session_id: str, response: str):
    """记录智能体响应"""
    logger.info(
        f"Agent response",
        extra={
            "agent_type": agent_type,
            "session_id": session_id,
            "response": response,
            "event_type": "agent_response"
        }
    )

def log_error(error_type: str, session_id: str, error: str):
    """记录错误"""
    logger.error(
        f"Error occurred",
        extra={
            "error_type": error_type,
            "session_id": session_id,
            "error": error,
            "event_type": "error"
        }
    )

# 示例用法
def handle_message(user_id: str, session_id: str, message: str):
    """处理消息"""
    try:
        # 记录用户消息
        log_user_message(user_id, session_id, message)
        
        # 处理消息...
        response = "这是智能体的回复"
        
        # 记录智能体响应
        log_agent_response("qa_agent", session_id, response)
        
        return response
    except Exception as e:
        # 记录错误
        log_error("processing_error", session_id, str(e))
        raise

if __name__ == "__main__":
    # 测试日志
    handle_message("user123", "session123", "测试消息")

6. 案例分析与最佳实践

6.1 成功案例

案例:某大型电商平台智能客服系统

背景:该电商平台日均处理超过100万条客户咨询,传统客服系统难以应对高峰期的咨询量。

解决方案

  1. 多智能体架构:部署了7个不同角色的智能体,各司其职
  2. 知识库构建:构建了包含100万+条知识的企业级知识库
  3. 多渠道接入:支持网站、APP、微信、电话等多种渠道
  4. 智能分流:基于意图识别的智能任务分发
  5. 人工协作:智能体无法解决的问题无缝转接人工客服

效果

  • 客服效率提升:智能客服处理了85%的常规咨询
  • 响应时间缩短:平均响应时间从3分钟缩短到30秒
  • 客户满意度提升:客户满意度从82%提升到95%
  • 运营成本降低:客服运营成本降低40%

6.2 最佳实践

  1. 智能体设计最佳实践

    • 职责单一:每个智能体专注于特定功能
    • 接口标准化:统一智能体间的通信接口
    • 状态管理:合理管理智能体状态
    • 错误处理:完善的错误处理和恢复机制
  2. 知识库管理最佳实践

    • 知识分类:按业务领域和使用频率分类
    • 定期更新:定期更新和维护知识库
    • 质量控制:确保知识库内容的准确性
    • 使用分析:分析知识库使用情况,优化内容
  3. 系统部署最佳实践

    • 容器化部署:使用Docker和Kubernetes
    • 环境隔离:开发、测试、生产环境隔离
    • 自动化部署:CI/CD流水线自动化部署
    • 监控告警:完善的监控和告警机制
  4. 性能优化最佳实践

    • 缓存策略:合理使用缓存,减少重复计算
    • 异步处理:非关键操作使用异步处理
    • 批量处理:批量处理相似请求
    • 资源调度:智能调度计算资源
Logo

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

更多推荐