目 录CONTENT

文章目录

实战:从零搭建企业级智能客服系统

PySuper
2025-09-13 / 0 评论 / 0 点赞 / 2 阅读 / 0 字
温馨提示:
所有牛逼的人都有一段苦逼的岁月。 但是你只要像SB一样去坚持,终将牛逼!!! ✊✊✊

一、项目背景与问题定义

1.1 业务背景

拥有日均 50万+ 的客服咨询量,峰值时段可达 2万/分钟

原有的客服系统存在以下痛点:

指标

现状

目标

人工客服数量

800人

降至200人

平均响应时间

45秒

<5秒(AI)

单次咨询成本

¥8.5

<¥2.0

问题解决率

72%

>85%

7x24覆盖

仅大夜班

全天候

1.2 核心问题

经过深入调研,我们发现传统客服系统存在三大核心问题:

1. 人工成本高昂

  • 人力成本占客服中心总成本的 60-70%

  • 客服培训周期长(平均2-3个月)

  • 人员流动率高(年流失率约 35%

  • 节假日/大促期间需要大量临时人员

2. 响应速度慢

  • 人工客服同时接待能力有限(通常1:5-1:10)

  • 高峰期排队严重,用户体验差

  • 知识更新滞后,答案不一致

3. 知识管理混乱

  • 产品知识分散在多个系统

  • 跨部门协作困难

  • 知识迭代慢,无法实时更新

1.3 为什么选择 RAG + 大模型

在技术选型阶段,我们对比了三种方案:

┌─────────────────────────────────────────────────────────────────┐
│                      技术方案对比                                  │
├─────────────────┬─────────────────┬─────────────────┬───────────┤
│     方案        │     方案A        │     方案B        │   方案C    │
│                 │  (纯规则匹配)     │ (大模型微调)     │ (RAG+LLM)  │
├─────────────────┼─────────────────┼─────────────────┼───────────┤
│ 实施难度        │      低         │      高         │    中     │
│ 知识更新成本    │      高         │      高         │    低     │
│ 回答准确性      │      中         │      中-高      │    高     │
│ 可控性          │      高         │      低         │    中     │
│ 部署成本        │      低         │      高         │    中     │
│ 幻觉问题        │      无         │      少         │    少     │
│ 冷启动速度      │      快         │      慢         │    快     │
└─────────────────┴─────────────────┴─────────────────┴───────────┘

最终选择 方案C(RAG + 大模型),原因如下:

  • 知识更新无需重新训练,成本低

  • 回答基于真实知识库,可溯源

  • 支持复杂多轮对话

  • 响应速度快,用户体验好


二、系统架构设计

2.1 整体架构图

                                    ┌──────────────────────────────────────────────────────────────┐
                                    │                        用户层                                    │
                                    │   ┌─────────┐  ┌─────────┐  ┌─────────┐  ┌─────────────────┐   │
                                    │   │  Web    │  │  APP    │  │ 微信    │  │   企业微信      │   │
                                    │   │  Portal │  │  SDK    │  │小程序   │  │   客服插件      │   │
                                    │   └────┬────┘  └────┬────┘  └────┬────┘  └────────┬────────┘   │
                                    └────────┼───────────┼───────────┼───────────┼────────────────────┘
                                            │           │           │           │
                                            ▼           ▼           ▼           ▼
┌──────────────────────────────────────────────────────────────────────────────────────────┐
│                              接入层 (Load Balancer + API Gateway)                          │
│    ┌────────────────────────────────────────────────────────────────────────────────┐      │
│    │                        Kong / Nginx + SSL Offload                             │      │
│    └────────────────────────────────────────────────────────────────────────────────┘      │
└────────────────────────────────────────────┬────────────────────────────────────────────────┘
                                             │
                                             ▼
┌──────────────────────────────────────────────────────────────────────────────────────────┐
│                                     网关层 (Django + DRF)                                  │
│  ┌──────────────┐  ┌──────────────┐  ┌──────────────┐  ┌──────────────┐  ┌──────────────┐   │
│  │  会话管理    │  │  意图识别    │  │   知识检索   │  │   意图路由   │  │   转人工     │   │
│  │   Service   │  │   Service   │  │   Service   │  │   Service   │  │   Service   │   │
│  └──────┬───────┘  └──────┬───────┘  └──────┬───────┘  └──────┬───────┘  └──────┬───────┘   │
│         │                 │                 │                 │                 │           │
│         └─────────────────┴────────┬────────┴─────────────────┴─────────────────┘           │
│                                    │                                                           │
│                                    ▼                                                           │
│  ┌────────────────────────────────────────────────────────────────────────────────────────┐    │
│  │                              消息队列 (Redis Streams + RabbitMQ)                        │    │
│  └────────────────────────────────────────────────────────────────────────────────────────┘    │
└────────────────────────────────────────────┬────────────────────────────────────────────────┘
                                             │
         ┌───────────────────────────────────┼───────────────────────────────────┐
         │                                   │                                   │
         ▼                                   ▼                                   ▼
┌─────────────────────┐        ┌─────────────────────┐        ┌─────────────────────┐
│    AI 处理层        │        │    数据处理层        │        │    人工客服层       │
│  ┌───────────────┐  │        │  ┌───────────────┐  │        │  ┌───────────────┐  │
│  │  意图识别模型  │  │        │  │  知识库构建   │  │        │  │  客服工作台   │  │
│  │  (BCE + LLM)  │  │        │  │  (ETL Pipeline)│ │        │  │  (WebSocket)  │  │
│  └───────────────┘  │        │  └───────────────┘  │        │  └───────────────┘  │
│  ┌───────────────┐  │        │  ┌───────────────┐  │        │  ┌───────────────┐  │
│  │   RAG 检索    │  │        │  │   离线分析    │  │        │  │  工单系统     │  │
│  │  (Embedding) │  │        │  │   (Spark)    │  │        │  │              │  │
│  └───────────────┘  │        │  └───────────────┘  │        │  └───────────────┘  │
│  ┌───────────────┐  │        │  ┌───────────────┐  │        │                     │
│  │   LLM 生成    │  │        │  │   监控报警    │  │        │                     │
│  │ (Qwen/GLM)   │  │        │  │  (Prometheus) │  │        │                     │
│  └───────────────┘  │        │  └───────────────┘  │        │                     │
└─────────────────────┘        └─────────────────────┘        └─────────────────────┘
         │
         ▼
┌──────────────────────────────────────────────────────────────────────────────────────────┐
│                                  知识库层                                                 │
│  ┌──────────────────┐  ┌──────────────────┐  ┌──────────────────┐  ┌──────────────────┐   │
│  │    Milvus        │  │    Elasticsearch │  │     MySQL        │  │   Redis          │   │
│  │  (向量数据库)     │  │   (全文检索)     │  │   (业务数据)     │  │   (缓存)         │   │
│  └──────────────────┘  └──────────────────┘  └──────────────────┘  └──────────────────┘   │
└──────────────────────────────────────────────────────────────────────────────────────────┘

2.2 核心流程

用户提问 ──┬──> 意图识别 ──> 意图路由
           │      │
           │      ├──> 知识问答 (FAQ) ──> RAG检索 ──> LLM生成
           │      │
           │      ├──> 闲聊 (Chat) ──> LLM生成
           │      │
           │      ├──> 任务型 (Task) ──> 意图槽位填充 ──> 动作执行
           │      │
           │      └──> 未知/无法处理 ──> 转人工 ──> 客服工作台
           │
           └──> 多轮对话上下文管理

三、核心模块设计与实现

3.1 项目结构

customer_service/
├── config/                    # 配置文件
│   ├── settings/             # Django设置
│   │   ├── base.py
│   │   ├── development.py
│   │   └── production.py
│   ├── settings.yaml         # 主配置
│   └── rag_config.yaml       # RAG配置
├── apps/
│   ├── conversation/          # 对话管理
│   │   ├── models.py
│   │   ├── views.py
│   │   ├── services/
│   │   │   ├── dialogue_manager.py
│   │   │   ├── intent_classifier.py
│   │   │   └── human_handoff.py
│   │   └── urls.py
│   ├── knowledge/             # 知识库管理
│   │   ├── models.py
│   │   ├── services/
│   │   │   ├── rag_engine.py
│   │   │   ├── embedding_service.py
│   │   │   └── reranker.py
│   │   └── management/
│   │       └── commands/
│   ├── agent/                # AI Agent
│   │   ├── llm_client.py
│   │   └── prompt_templates.py
│   └── routing/              # 路由管理
├── common/                   # 公共组件
│   ├── middleware.py
│   ├── cache.py
│   └── utils.py
├── docker/
│   ├── docker-compose.yml
│   └── Dockerfile
├── k8s/                      # Kubernetes部署
│   ├── deployment.yaml
│   ├── service.yaml
│   ├── configmap.yaml
│   └── hpa.yaml
├── requirements.txt
└── manage.py

3.2 Django 核心配置

config/settings.yaml - 主配置文件

# 企业智能客服系统配置文件
# 版本: 1.0.0
# 更新时间: 2025-09-15

app:
  name: "智能客服系统"
  version: "2.0.0"
  debug: false
  secret_key: "${SECRET_KEY}"

server:
  host: "0.0.0.0"
  port: 8000
  workers: 4
  timeout: 120

# 数据库配置
database:
  engine: "django.db.backends.postgresql"
  name: "customer_service"
  host: "${DB_HOST}"
  port: 5432
  user: "${DB_USER}"
  password: "${DB_PASSWORD}"
  pool_size: 20
  max_overflow: 10

# Redis 配置
redis:
  host: "${REDIS_HOST}"
  port: 6379
  db: 0
  password: "${REDIS_PASSWORD}"
  max_connections: 100
  session_ttl: 86400  # 24小时

# LLM 模型配置
llm:
  provider: "dashscope"  # 阿里云百炼
  model: "qwen-max"
  api_key: "${DASHSCOPE_API_KEY}"
  base_url: "https://dashscope.aliyuncs.com/compatible-mode/v1"
  temperature: 0.7
  max_tokens: 2000
  top_p: 0.95
  timeout: 30

# 向量数据库配置
vector_db:
  type: "milvus"
  host: "${MILVUS_HOST}"
  port: 19530
  collection_name: "knowledge_base"
  dimension: 1536
  metric_type: "COSINE"

# Embedding 模型配置
embedding:
  provider: "dashscope"
  model: "text-embedding-v3"
  dimension: 1536
  batch_size: 100

# 重排序模型配置
reranker:
  model_name: "BAAI/bge-reranker-base"
  top_k: 10  # RAG检索返回的候选数

# 意图识别配置
intent:
  model_type: "bert_classifier"
  model_path: "/models/intent_classifier"
  labels:
    - "faq"           # 知识问答
    - "chitchat"      # 闲聊
    - "task"          # 任务型
    - "complaint"     # 投诉
    - "transfer"      # 转人工
  threshold: 0.75     # 意图置信度阈值
  fallback_threshold: 0.5  # 低于此值转人工

# 对话配置
conversation:
  max_history_turns: 10      # 最大历史轮次
  session_timeout: 1800     # 30分钟无交互则会话结束
  max_response_time: 5      # 最大响应时间(秒)

# 人工客服配置
human_handoff:
  auto_transfer_keywords:   # 触发转人工的关键词
    - "投诉"
    - "退款"
    - "人工"
    - "升级"
    - "经理"
  escalation_threshold: 3   # AI连续3次回答不好则转人工
  queue_name: "customer_service_queue"
  max_wait_time: 300       # 等待超过5分钟自动转人工

# 限流配置
rate_limit:
  enabled: true
  requests_per_minute: 60
  burst: 10

# 日志配置
logging:
  level: "INFO"
  format: "%(asctime)s - %(name)s - %(levelname)s - %(message)s"
  file: "/var/log/customer_service.log"
  max_bytes: 104857600  # 100MB
  backup_count: 10

3.3 RAG 检索引擎实现

apps/knowledge/services/rag_engine.py - RAG核心引擎

"""
RAG 检索引擎
基于向量检索 + 重排序的双阶段检索方案

核心流程:
1. 用户 query -> Embedding
2. 向量数据库检索 -> Top-K 候选文档
3. 重排序模型 -> 最终排序结果
4. 构建上下文 -> LLM 生成回答

Author: AI Team
Date: 2025-09-15
"""

import time
import hashlib
from typing import List, Dict, Optional, Tuple
from dataclasses import dataclass
from functools import lru_cache
import logging

import torch
import numpy as np
from pymilvus import Collection, connections, utility
from sentence_transformers import SentenceTransformer
from transformers import AutoModelForSequenceClassification, AutoTokenizer
import redis
from django.conf import settings

logger = logging.getLogger(__name__)


@dataclass
class Document:
    """文档数据结构"""
    id: str                      # 文档唯一ID
    content: str                 # 文档内容
    metadata: Dict[str, any]     # 元信息(来源、类别、更新时间等)
    score: float = 0.0          # 检索相关性分数
    vector: Optional[np.ndarray] = None  # 向量表示


@dataclass
class SearchResult:
    """检索结果"""
    documents: List[Document]     # 文档列表
    query: str                   # 原始查询
    total_time: float           # 检索耗时(ms)
    retrieval_type: str         # 检索类型


class RAGEngine:
    """
    RAG 检索引擎
    
    支持功能:
    - 双阶段检索:向量检索 + 重排序
    - 混合检索:向量 + 关键词
    - 缓存优化:热点query结果缓存
    - 批量检索:支持批量处理
    """
    
    def __init__(self):
        """初始化RAG引擎"""
        self._initialized = False
        self._init_clients()
        self._init_models()
        self._init_cache()
        
    def _init_clients(self):
        """初始化客户端连接"""
        # Milvus 连接
        connections.connect(
            alias="default",
            host=settings.VECTOR_DB.get("host"),
            port=settings.VECTOR_DB.get("port")
        )
        self.collection = Collection(settings.VECTOR_DB.get("collection_name"))
        self.collection.load()
        
        # Redis 连接(用于缓存)
        self.redis_client = redis.Redis(
            host=settings.REDIS.get("host"),
            port=settings.REDIS.get("port"),
            db=settings.REDIS.get("db", 1),  # RAG使用独立的DB
            password=settings.REDIS.get("password"),
            decode_responses=False  # 存储向量需要二进制
        )
        
        logger.info("RAG引擎客户端初始化完成")
        
    def _init_models(self):
        """初始化模型"""
        # Embedding 模型
        self.embedding_model = SentenceTransformer(
            'paraphrase-multilingual-MiniLM-L12-v2'  # 多语言支持
        )
        
        # 重排序模型
        self.reranker_tokenizer = AutoTokenizer.from_pretrained(
            settings.RERANKER.get("model_name")
        )
        self.reranker_model = AutoModelForSequenceClassification.from_pretrained(
            settings.RERANKER.get("model_name")
        )
        self.reranker_model.eval()
        
        # 设备选择
        self.device = "cuda" if torch.cuda.is_available() else "cpu"
        self.reranker_model.to(self.device)
        
        logger.info(f"RAG模型加载完成,设备: {self.device}")
        
    def _init_cache(self):
        """初始化缓存策略"""
        self.cache_ttl = 3600  # 缓存1小时
        self.cache_enabled = True
        
    def _get_cache_key(self, query: str) -> str:
        """生成缓存键"""
        # 使用query的hash作为缓存key
        query_hash = hashlib.md5(query.encode()).hexdigest()
        return f"rag:cache:{query_hash}"
    
    @lru_cache(maxsize=1000)
    def _encode_query(self, query: str) -> np.ndarray:
        """Query编码(带缓存)"""
        embedding = self.embedding_model.encode(
            query,
            convert_to_numpy=True,
            normalize_embeddings=True  # L2归一化
        )
        return embedding.astype(np.float32)
    
    def search(
        self,
        query: str,
        top_k: int = 5,
        use_reranker: bool = True,
        filter_conditions: Optional[Dict] = None,
        enable_cache: bool = True
    ) -> SearchResult:
        """
        文档检索
        
        Args:
            query: 用户查询
            top_k: 返回的文档数量
            use_reranker: 是否使用重排序
            filter_conditions: 过滤条件(如按类别、按时间范围)
            enable_cache: 是否启用缓存
            
        Returns:
            SearchResult: 检索结果
        """
        start_time = time.time()
        
        # 1. 检查缓存
        if enable_cache and self.cache_enabled:
            cache_key = self._get_cache_key(query)
            cached_result = self._get_from_cache(cache_key)
            if cached_result:
                logger.debug(f"缓存命中: {query[:50]}")
                cached_result.total_time = (time.time() - start_time) * 1000
                return cached_result
        
        # 2. Query向量化
        query_vector = self._encode_query(query)
        
        # 3. 向量检索 (Milvus ANN检索)
        search_params = {
            "metric_type": "COSINE",
            "params": {"nprobe": 10}
        }
        
        # 构建过滤表达式
        expr = self._build_filter_expr(filter_conditions) if filter_conditions else None
        
        # 检索候选集(取top_k的10倍用于重排序)
        candidate_k = top_k * 10
        results = self.collection.search(
            data=[query_vector.tolist()],
            anns_field="embedding",
            param=search_params,
            limit=candidate_k,
            expr=expr,
            output_fields=["id", "content", "metadata", "category", "source"]
        )
        
        # 4. 转换为Document对象
        candidates = []
        for hit in results[0]:
            doc = Document(
                id=hit.entity.get("id"),
                content=hit.entity.get("content"),
                metadata={
                    "category": hit.entity.get("category"),
                    "source": hit.entity.get("source"),
                    "updated_at": hit.entity.get("updated_at")
                },
                score=hit.distance
            )
            candidates.append(doc)
        
        # 5. 重排序(如启用)
        if use_reranker and candidates:
            candidates = self._rerank(query, candidates, top_k)
        
        # 6. 构建返回结果
        search_result = SearchResult(
            documents=candidates[:top_k],
            query=query,
            total_time=(time.time() - start_time) * 1000,
            retrieval_type="vector+reranker" if use_reranker else "vector"
        )
        
        # 7. 更新缓存
        if enable_cache and self.cache_enabled:
            self._save_to_cache(cache_key, search_result)
        
        logger.info(
            f"检索完成: query='{query[:30]}...', "
            f"candidates={len(candidates)}, "
            f"time={search_result.total_time:.2f}ms"
        )
        
        return search_result
    
    def _rerank(
        self,
        query: str,
        candidates: List[Document],
        top_k: int
    ) -> List[Document]:
        """
        使用交叉编码器进行重排序
        
        重排序原理:
        - 将 query 和 document 拼接作为输入
        - 使用预训练的 Cross-Encoder 模型预测相关性分数
        - 按分数排序返回
        """
        # 准备输入数据
        pairs = [(query, doc.content) for doc in candidates]
        
        # Tokenize
        inputs = self.reranker_tokenizer(
            pairs,
            padding=True,
            truncation=True,
            max_length=512,
            return_tensors="pt"
        ).to(self.device)
        
        # 推理
        with torch.no_grad():
            outputs = self.reranker_model(**inputs)
            scores = outputs.logits.squeeze(-1).cpu().numpy()
        
        # 排序
        sorted_indices = np.argsort(scores)[::-1]
        
        # 更新文档分数并返回
        reranked = []
        for idx in sorted_indices[:top_k]:
            doc = candidates[idx]
            doc.score = float(scores[idx])
            reranked.append(doc)
        
        logger.debug(f"重排序完成: {len(candidates)} -> {len(reranked)}")
        return reranked
    
    def _build_filter_expr(self, conditions: Dict) -> str:
        """构建Milvus过滤表达式"""
        expr_parts = []
        
        if "category" in conditions:
            categories = conditions["category"]
            if isinstance(categories, str):
                categories = [categories]
            cat_expr = " || ".join([f'category == "{c}"' for c in categories])
            expr_parts.append(f"({cat_expr})")
        
        if "source" in conditions:
            expr_parts.append(f'source == "{conditions["source"]}"')
        
        if "date_range" in conditions:
            start, end = conditions["date_range"]
            expr_parts.append(f"updated_at >= {start} && updated_at <= {end}")
        
        return " && ".join(expr_parts) if expr_parts else None
    
    def batch_search(
        self,
        queries: List[str],
        top_k: int = 5,
        **kwargs
    ) -> List[SearchResult]:
        """批量检索"""
        return [self.search(q, top_k, **kwargs) for q in queries]
    
    def _get_from_cache(self, cache_key: str) -> Optional[SearchResult]:
        """从缓存获取"""
        try:
            import pickle
            cached = self.redis_client.get(cache_key)
            if cached:
                return pickle.loads(cached)
        except Exception as e:
            logger.warning(f"缓存读取失败: {e}")
        return None
    
    def _save_to_cache(self, cache_key: str, result: SearchResult):
        """保存到缓存"""
        try:
            import pickle
            self.redis_client.setex(
                cache_key,
                self.cache_ttl,
                pickle.dumps(result)
            )
        except Exception as e:
            logger.warning(f"缓存写入失败: {e}")
    
    def add_documents(self, documents: List[Dict]):
        """
        添加文档到知识库
        
        Args:
            documents: 文档列表,每项包含 content, metadata
        """
        from pymilvus import CollectionSchema, FieldSchema, DataType
        
        entities = []
        embeddings = []
        
        for doc in documents:
            # 向量化
            embedding = self.embedding_model.encode(
                [doc["content"]],
                convert_to_numpy=True,
                normalize_embeddings=True
            )[0].tolist()
            embeddings.append(embedding)
            
            # 构建实体
            entities.append({
                "id": doc.get("id", hashlib.md5(doc["content"].encode()).hexdigest()),
                "content": doc["content"],
                "metadata": doc.get("metadata", {}),
                "category": doc.get("category", "general"),
                "source": doc.get("source", "manual"),
                "updated_at": int(time.time())
            })
        
        # 批量插入
        self.collection.insert([entities, embeddings])
        self.collection.flush()
        
        logger.info(f"成功添加 {len(documents)} 篇文档到知识库")
    
    def warm_up(self, popular_queries: List[str]):
        """
        预热缓存
        
        针对热点query提前进行检索和缓存
        """
        logger.info(f"开始预热缓存,预热 {len(popular_queries)} 个查询...")
        for query in popular_queries:
            self.search(query, enable_cache=True)
        logger.info("缓存预热完成")
    
    def health_check(self) -> Dict:
        """健康检查"""
        health = {
            "status": "healthy",
            "milvus": False,
            "redis": False,
            "embedding_model": False,
            "reranker_model": False
        }
        
        try:
            # 检查Milvus
            health["milvus"] = utility.get_connection_status("default")
        except:
            health["status"] = "degraded"
        
        try:
            # 检查Redis
            self.redis_client.ping()
            health["redis"] = True
        except:
            health["status"] = "degraded"
        
        # 检查模型
        health["embedding_model"] = self.embedding_model is not None
        health["reranker_model"] = self.reranker_model is not None
        
        return health
    
    def __del__(self):
        """清理资源"""
        try:
            connections.disconnect("default")
            self.redis_client.close()
        except:
            pass

3.4 对话管理器实现

apps/conversation/services/dialogue_manager.py - 多轮对话管理

"""
对话管理器
负责多轮对话的上下文管理、历史记录和状态跟踪

核心功能:
1. 会话生命周期管理
2. 上下文窗口维护(滑动窗口策略)
3. 对话状态机
4. 跨会话上下文复用

Author: AI Team
Date: 2025-09-15
"""

import json
import time
import uuid
from typing import Dict, List, Optional, Any
from dataclasses import dataclass, field, asdict
from enum import Enum
from collections import deque
import logging

import redis
from django.conf import settings
import jsonpickle

logger = logging.getLogger(__name__)


class DialogueState(Enum):
    """对话状态枚举"""
    INIT = "init"                    # 初始状态
    INTENT_DETECTED = "intent_detected"  # 意图已识别
    KNOWLEDGE_QA = "knowledge_qa"    # 知识问答中
    TASK_EXECUTION = "task_execution"  # 任务执行中
    WAITING_HUMAN = "waiting_human"  # 等待人工客服
    HUMAN_IN_CHARGE = "human_in_charge"  # 人工接管
    COMPLETED = "completed"          # 对话完成
    EXPIRED = "expired"              # 会话过期


class MessageRole(Enum):
    """消息角色"""
    USER = "user"
    ASSISTANT = "assistant"
    SYSTEM = "system"
    HUMAN_AGENT = "human_agent"     # 人工客服


@dataclass
class Message:
    """对话消息"""
    role: MessageRole
    content: str
    timestamp: float = field(default_factory=time.time)
    message_id: str = field(default_factory=lambda: str(uuid.uuid4()))
    metadata: Dict[str, Any] = field(default_factory=dict)
    
    def to_dict(self) -> Dict:
        """转换为字典"""
        return {
            "role": self.role.value,
            "content": self.content,
            "timestamp": self.timestamp,
            "message_id": self.message_id,
            "metadata": self.metadata
        }
    
    @classmethod
    def from_dict(cls, data: Dict) -> "Message":
        """从字典创建"""
        return cls(
            role=MessageRole(data["role"]),
            content=data["content"],
            timestamp=data.get("timestamp", time.time()),
            message_id=data.get("message_id", str(uuid.uuid4())),
            metadata=data.get("metadata", {})
        )


@dataclass
class DialogueContext:
    """
    对话上下文
    
    包含:
    - 会话基本信息
    - 历史消息
    - 意图状态
    - 槽位信息
    - 用户画像
    """
    session_id: str
    user_id: Optional[str] = None
    state: DialogueState = DialogueState.INIT
    created_at: float = field(default_factory=time.time)
    updated_at: float = field(default_factory=time.time)
    messages: List[Message] = field(default_factory=list)
    intent: Optional[str] = None
    intent_confidence: float = 0.0
    slots: Dict[str, Any] = field(default_factory=dict)
    human_escalation_count: int = 0  # 连续无法回答次数
    last_ai_response: Optional[str] = None
    metadata: Dict[str, Any] = field(default_factory=dict)
    
    def to_dict(self) -> Dict:
        """序列化"""
        return {
            "session_id": self.session_id,
            "user_id": self.user_id,
            "state": self.state.value,
            "created_at": self.created_at,
            "updated_at": self.updated_at,
            "messages": [m.to_dict() for m in self.messages],
            "intent": self.intent,
            "intent_confidence": self.intent_confidence,
            "slots": self.slots,
            "human_escalation_count": self.human_escalation_count,
            "last_ai_response": self.last_ai_response,
            "metadata": self.metadata
        }
    
    @classmethod
    def from_dict(cls, data: Dict) -> "DialogueContext":
        """反序列化"""
        messages = [Message.from_dict(m) for m in data.get("messages", [])]
        return cls(
            session_id=data["session_id"],
            user_id=data.get("user_id"),
            state=DialogueState(data.get("state", "init")),
            created_at=data.get("created_at", time.time()),
            updated_at=data.get("updated_at", time.time()),
            messages=messages,
            intent=data.get("intent"),
            intent_confidence=data.get("intent_confidence", 0.0),
            slots=data.get("slots", {}),
            human_escalation_count=data.get("human_escalation_count", 0),
            last_ai_response=data.get("last_ai_response"),
            metadata=data.get("metadata", {})
        )


class DialogueManager:
    """
    对话管理器
    
    核心职责:
    1. 创建和管理会话
    2. 维护对话历史(滑动窗口)
    3. 管理对话状态转换
    4. 提供上下文检索
    """
    
    # Redis key 前缀
    SESSION_PREFIX = "dialogue:session:"
    SESSION_INDEX_PREFIX = "dialogue:user_sessions:"
    
    def __init__(self):
        """初始化对话管理器"""
        self.redis_client = redis.Redis(
            host=settings.REDIS.get("host"),
            port=settings.REDIS.get("port"),
            db=settings.REDIS.get("db", 0),
            password=settings.REDIS.get("password"),
            decode_responses=False  # 需要存储序列化对象
        )
        
        # 配置
        self.max_history_turns = settings.CONVERSATION.get("max_history_turns", 10)
        self.session_timeout = settings.CONVERSATION.get("session_timeout", 1800)
        self.max_context_tokens = 4000  # 约等于2000个中文字符
        
        logger.info("对话管理器初始化完成")
    
    def create_session(
        self,
        user_id: Optional[str] = None,
        metadata: Optional[Dict] = None
    ) -> DialogueContext:
        """
        创建新会话
        
        Args:
            user_id: 用户ID(可选)
            metadata: 额外元信息
            
        Returns:
            DialogueContext: 新建的会话上下文
        """
        session_id = str(uuid.uuid4())
        
        context = DialogueContext(
            session_id=session_id,
            user_id=user_id,
            state=DialogueState.INIT,
            metadata=metadata or {}
        )
        
        # 持久化到Redis
        self._save_context(context)
        
        # 更新用户会话索引
        if user_id:
            self._add_to_user_index(user_id, session_id)
        
        logger.info(f"创建新会话: {session_id}, user_id: {user_id}")
        
        return context
    
    def get_session(self, session_id: str) -> Optional[DialogueContext]:
        """
        获取会话上下文
        
        Args:
            session_id: 会话ID
            
        Returns:
            DialogueContext 或 None(会话不存在或已过期)
        """
        key = f"{self.SESSION_PREFIX}{session_id}"
        
        try:
            data = self.redis_client.get(key)
            if not data:
                return None
            
            context = jsonpickle.decode(data)
            
            # 检查是否过期
            if time.time() - context.updated_at > self.session_timeout:
                logger.info(f"会话已过期: {session_id}")
                self.delete_session(session_id)
                return None
            
            return context
            
        except Exception as e:
            logger.error(f"获取会话失败: {session_id}, error: {e}")
            return None
    
    def update_session(self, context: DialogueContext) -> bool:
        """
        更新会话上下文
        
        Args:
            context: 会话上下文
            
        Returns:
            bool: 更新是否成功
        """
        context.updated_at = time.time()
        return self._save_context(context)
    
    def add_message(
        self,
        session_id: str,
        role: MessageRole,
        content: str,
        metadata: Optional[Dict] = None
    ) -> Optional[DialogueContext]:
        """
        添加消息到会话
        
        Args:
            session_id: 会话ID
            role: 消息角色
            content: 消息内容
            metadata: 额外元信息
            
        Returns:
            更新后的上下文,或None(会话不存在)
        """
        context = self.get_session(session_id)
        if not context:
            logger.warning(f"会话不存在,无法添加消息: {session_id}")
            return None
        
        # 创建消息
        message = Message(
            role=role,
            content=content,
            metadata=metadata or {}
        )
        
        # 添加到历史
        context.messages.append(message)
        
        # 滑动窗口:保持最近的max_history_turns轮对话
        # 每轮 = user + assistant 各一条
        max_messages = self.max_history_turns * 2
        if len(context.messages) > max_messages:
            context.messages = context.messages[-max_messages:]
        
        # 特殊处理:记录AI的最后回复(用于判断是否需要转人工)
        if role == MessageRole.ASSISTANT:
            context.last_ai_response = content
        
        # 更新状态
        context.updated_at = time.time()
        
        # 保存
        self._save_context(context)
        
        return context
    
    def get_conversation_history(
        self,
        session_id: str,
        max_turns: Optional[int] = None,
        include_metadata: bool = False
    ) -> List[Dict]:
        """
        获取对话历史
        
        Args:
            session_id: 会话ID
            max_turns: 最大返回轮数(每轮包含user+assistant)
            include_metadata: 是否包含元信息
            
        Returns:
            对话历史列表
        """
        context = self.get_session(session_id)
        if not context:
            return []
        
        messages = context.messages
        
        # 限制返回的轮数
        if max_turns:
            messages = messages[-(max_turns * 2):]
        
        if include_metadata:
            return [m.to_dict() for m in messages]
        else:
            return [{"role": m.role.value, "content": m.content} for m in messages]
    
    def get_context_for_llm(
        self,
        session_id: str,
        system_prompt: Optional[str] = None
    ) -> List[Dict]:
        """
        获取适合LLM处理的对话上下文
        
        格式化为消息列表,可直接用于API调用
        
        Args:
            session_id: 会话ID
            system_prompt: 系统提示词
            
        Returns:
            格式化后的消息列表
        """
        context = self.get_session(session_id)
        if not context:
            return []
        
        messages = []
        
        # 系统提示词
        if system_prompt:
            messages.append({
                "role": "system",
                "content": system_prompt
            })
        
        # 对话历史
        for msg in context.messages[-self.max_history_turns * 2:]:
            messages.append({
                "role": msg.role.value,
                "content": msg.content
            })
        
        return messages
    
    def update_state(
        self,
        session_id: str,
        state: DialogueState,
        **kwargs
    ) -> Optional[DialogueContext]:
        """
        更新对话状态
        
        Args:
            session_id: 会话ID
            state: 新状态
            **kwargs: 其他更新的字段
        """
        context = self.get_session(session_id)
        if not context:
            return None
        
        context.state = state
        context.updated_at = time.time()
        
        # 更新其他字段
        for key, value in kwargs.items():
            if hasattr(context, key):
                setattr(context, key, value)
        
        self._save_context(context)
        
        logger.info(f"会话状态更新: {session_id} -> {state.value}")
        
        return context
    
    def increment_escalation(self, session_id: str) -> int:
        """
        增加人工转接计数
        
        Returns:
            更新后的计数
        """
        context = self.get_session(session_id)
        if context:
            context.human_escalation_count += 1
            self._save_context(context)
        return context.human_escalation_count if context else 0
    
    def check_escalation_needed(self, session_id: str) -> bool:
        """
        检查是否需要转人工
        
        条件:连续3次AI回答不满意
        """
        context = self.get_session(session_id)
        if not context:
            return False
        
        threshold = settings.HUMAN_HANDOFF.get("escalation_threshold", 3)
        return context.human_escalation_count >= threshold
    
    def transfer_to_human(
        self,
        session_id: str,
        reason: str = "user_request",
        priority: int = 1
    ) -> Dict:
        """
        执行转人工操作
        
        Args:
            session_id: 会话ID
            reason: 转人工原因
            priority: 优先级(1-5,5最高)
            
        Returns:
            转接结果信息
        """
        context = self.get_session(session_id)
        if not context:
            return {"success": False, "error": "会话不存在"}
        
        # 更新状态
        context.state = DialogueState.HUMAN_IN_CHARGE
        context.metadata["transfer_reason"] = reason
        context.metadata["transfer_time"] = time.time()
        context.metadata["priority"] = priority
        self._save_context(context)
        
        # 获取等待中的客服
        # 这里可以接入实际的客服队列系统
        agent_id = self._allocate_agent(priority)
        
        logger.info(
            f"会话转人工: {session_id}, reason={reason}, "
            f"priority={priority}, agent={agent_id}"
        )
        
        return {
            "success": True,
            "session_id": session_id,
            "agent_id": agent_id,
            "estimated_wait_time": 60  # 预估等待时间(秒)
        }
    
    def _allocate_agent(self, priority: int) -> Optional[str]:
        """
        分配客服
        
        实际实现中应该查询客服队列系统
        """
        # 简化实现:返回模拟的agent_id
        return f"agent_{priority}_{int(time.time())}"
    
    def end_session(
        self,
        session_id: str,
        summary: Optional[str] = None
    ) -> bool:
        """
        结束会话
        
        Args:
            session_id: 会话ID
            summary: 对话摘要(用于存储日志)
        """
        context = self.get_session(session_id)
        if not context:
            return False
        
        # 更新状态
        context.state = DialogueState.COMPLETED
        context.metadata["end_time"] = time.time()
        if summary:
            context.metadata["summary"] = summary
        
        # 保存最终状态(可以转移到持久化存储)
        self._save_context(context)
        
        # 从用户索引中移除
        if context.user_id:
            self._remove_from_user_index(context.user_id, session_id)
        
        logger.info(f"会话结束: {session_id}")
        
        return True
    
    def delete_session(self, session_id: str) -> bool:
        """删除会话"""
        key = f"{self.SESSION_PREFIX}{session_id}"
        self.redis_client.delete(key)
        return True
    
    def _save_context(self, context: DialogueContext) -> bool:
        """保存上下文到Redis"""
        key = f"{self.SESSION_PREFIX}{context.session_id}"
        
        try:
            # 设置过期时间
            ttl = self.session_timeout * 2  # 双倍过期时间作为缓冲
            
            # 序列化并存储
            data = jsonpickle.encode(context)
            self.redis_client.setex(key, ttl, data)
            
            return True
            
        except Exception as e:
            logger.error(f"保存会话失败: {context.session_id}, error: {e}")
            return False
    
    def _add_to_user_index(self, user_id: str, session_id: str):
        """添加到用户会话索引"""
        key = f"{self.SESSION_INDEX_PREFIX}{user_id}"
        self.redis_client.sadd(key, session_id)
    
    def _remove_from_user_index(self, user_id: str, session_id: str):
        """从用户会话索引移除"""
        key = f"{self.SESSION_INDEX_PREFIX}{user_id}"
        self.redis_client.srem(key, session_id)
    
    def get_user_sessions(self, user_id: str) -> List[str]:
        """获取用户的所有会话ID"""
        key = f"{self.SESSION_INDEX_PREFIX}{user_id}"
        return list(self.redis_client.smembers(key))
    
    def get_session_stats(self) -> Dict:
        """获取会话统计信息"""
        info = self.redis_client.info()
        
        # 统计各状态的会话数(需要扫描)
        # 这里简化为返回Redis统计信息
        return {
            "total_keys": info.get("db0", {}).get("keys", 0),
            "used_memory": info.get("used_memory_human"),
            "connected_clients": info.get("connected_clients")
        }

3.5 意图识别与路由

apps/conversation/services/intent_classifier.py - 意图分类器

"""
意图识别与路由
使用多策略融合的意图识别方案

策略:
1. 规则匹配(快速、精确)
2. 语义相似度(通用、鲁棒)
3. LLM分类(复杂场景兜底)

Author: AI Team
Date: 2025-09-15
"""

import time
import re
from typing import Dict, List, Optional, Tuple
from dataclasses import dataclass
from enum import Enum
import logging

import torch
from transformers import BertTokenizer, BertModel
import numpy as np
from django.conf import settings

logger = logging.getLogger(__name__)


class Intent(Enum):
    """意图枚举"""
    FAQ = "faq"               # 知识问答
    CHITCHAT = "chitchat"    # 闲聊
    TASK = "task"            # 任务型
    COMPLAINT = "complaint"  # 投诉
    TRANSFER = "transfer"    # 转人工
    
    # 子意图(知识问答)
    PRODUCT_INQUIRY = "product_inquiry"      # 产品咨询
    ORDER_INQUIRY = "order_inquiry"          # 订单咨询
    RETURN_REFUND = "return_refund"          # 退换货
    PAYMENT_ISSUE = "payment_issue"          # 支付问题
    SHIPPING_INQUIRY = "shipping_inquiry"    # 物流咨询


@dataclass
class IntentResult:
    """意图识别结果"""
    intent: Intent
    confidence: float
    sub_intent: Optional[str] = None
    sub_confidence: Optional[float] = None
    slots: Dict[str, str] = None
    reasoning: str = ""
    method: str = "hybrid"  # 识别方法: rule/similarity/model


class IntentClassifier:
    """
    多策略意图识别器
    
    采用分层策略:
    Level 1: 规则匹配(关键词、正则)
    Level 2: 语义相似度(Embedding + 余弦相似度)
    Level 3: BERT分类模型(复杂/模糊场景)
    """
    
    def __init__(self):
        """初始化意图识别器"""
        self._init_rules()
        self._init_model()
        self._init_keywords()
        
    def _init_rules(self):
        """初始化规则配置"""
        # 转人工关键词
        self.transfer_keywords = [
            "人工", "客服", "人工客服", "转人工", "真人",
            "投诉", "升级", "经理", "主管", "老板",
            "解决不了", "不行", "太差", "非常不满"
        ]
        
        # 投诉关键词
        self.complaint_keywords = [
            "投诉", "举报", "差评", "不满", "失望",
            "愤怒", "垃圾", "骗子", "欺诈", "虚假"
        ]
        
        # 闲聊模式
        self.chitchat_patterns = [
            r"今天天气",
            r"你好|您好|嗨|hi|hello",
            r"你是谁|你叫什么",
            r"再见|拜拜|下次见",
            r"谢谢|感谢",
            r"辛苦了|加油"
        ]
        
        # 任务型关键词
        self.task_keywords = {
            "order_status": ["订单状态", "订单查询", "查订单", "发货了没"],
            "tracking": ["物流", "快递", "到哪了", "什么时候到"],
            "refund": ["退款", "退货", "取消订单"],
            "modify_address": ["改地址", "修改地址", "地址错了"],
            "invoice": ["发票", "开票", "增票"]
        }
        
        logger.info("规则配置加载完成")
    
    def _init_keywords(self):
        """初始化关键词向量(用于快速匹配)"""
        # FAQ 子意图关键词
        self.faq_keywords = {
            "product_inquiry": [
                "怎么用", "使用方法", "功能", "规格", "参数",
                "材质", "尺寸", "颜色", "款式", "好不好"
            ],
            "order_inquiry": [
                "什么时候", "多久", "发货", "到货",
                "状态", "单号", "订单号"
            ],
            "return_refund": [
                "退货", "退款", "换货", "七天无理由",
                "售后", "维修", "保修", "质保"
            ],
            "payment_issue": [
                "支付", "付款", "银行卡", "微信", "支付宝",
                "优惠券", "红包", "积分", "抵扣"
            ],
            "shipping_inquiry": [
                "快递", "物流", "配送", "送货", "自提",
                "网点", "站点", "派送", "签收"
            ]
        }
        
        # 编译正则表达式
        self.chitchat_regex = [
            re.compile(p, re.IGNORECASE) 
            for p in self.chitchat_patterns
        ]
        
        # 编译关键词匹配
        self.transfer_regex = re.compile(
            "|".join([re.escape(k) for k in self.transfer_keywords]),
            re.IGNORECASE
        )
        self.complaint_regex = re.compile(
            "|".join([re.escape(k) for k in self.complaint_keywords]),
            re.IGNORECASE
        )
    
    def _init_model(self):
        """初始化分类模型"""
        # 使用预训练的中文BERT模型
        self.tokenizer = BertTokenizer.from_pretrained("bert-base-chinese")
        self.model = BertModel.from_pretrained("bert-base-chinese")
        self.model.eval()
        
        # 意图标签映射
        self.intent_labels = {
            0: Intent.FAQ,
            1: Intent.CHITCHAT,
            2: Intent.TASK,
            3: Intent.COMPLAINT,
            4: Intent.TRANSFER
        }
        
        # 加载微调的分类头(如果有)
        # self.classifier.load_state_dict(...)
        
        self.device = "cuda" if torch.cuda.is_available() else "cpu"
        self.model.to(self.device)
        
        logger.info(f"意图分类模型加载完成,设备: {self.device}")
    
    def classify(self, query: str, context: Optional[Dict] = None) -> IntentResult:
        """
        意图分类主入口
        
        Args:
            query: 用户query
            context: 上下文信息(可选)
            
        Returns:
            IntentResult: 识别结果
        """
        start_time = time.time()
        query = query.strip()
        
        # Level 1: 规则匹配
        rule_result = self._match_rules(query)
        if rule_result and rule_result.confidence >= 0.9:
            logger.debug(f"规则匹配命中: {rule_result.intent.value}")
            rule_result.reasoning = f"规则匹配 ({time.time() - start_time:.3f}s)"
            return rule_result
        
        # Level 2: 关键词+语义相似度
        keyword_result = self._match_keywords(query)
        
        # Level 3: 模型分类(作为兜底)
        model_result = self._predict_with_model(query)
        
        # 融合结果
        final_result = self._fuse_results(
            query, 
            [rule_result, keyword_result, model_result],
            context
        )
        
        final_result.reasoning = f"融合策略 ({time.time() - start_time:.3f}s)"
        
        logger.info(
            f"意图识别: '{query[:30]}...' -> {final_result.intent.value} "
            f"(conf={final_result.confidence:.3f}, method={final_result.method})"
        )
        
        return final_result
    
    def _match_rules(self, query: str) -> Optional[IntentResult]:
        """Level 1: 规则匹配"""
        # 检查转人工
        if self.transfer_regex.search(query):
            return IntentResult(
                intent=Intent.TRANSFER,
                confidence=0.95,
                method="rule"
            )
        
        # 检查投诉
        if self.complaint_regex.search(query):
            return IntentResult(
                intent=Intent.COMPLAINT,
                confidence=0.92,
                method="rule"
            )
        
        # 检查闲聊
        for pattern in self.chitchat_regex:
            if pattern.search(query):
                return IntentResult(
                    intent=Intent.CHITCHAT,
                    confidence=0.88,
                    method="rule"
                )
        
        return None
    
    def _match_keywords(self, query: str) -> IntentResult:
        """Level 2: 关键词匹配 + 简单语义"""
        query_lower = query.lower()
        
        # 检查任务型关键词
        for task_type, keywords in self.task_keywords.items():
            for keyword in keywords:
                if keyword in query_lower:
                    return IntentResult(
                        intent=Intent.TASK,
                        confidence=0.85,
                        sub_intent=task_type,
                        method="keyword"
                    )
        
        # 检查FAQ子意图
        for faq_type, keywords in self.faq_keywords.items():
            match_count = sum(1 for k in keywords if k in query_lower)
            if match_count >= 2:
                return IntentResult(
                    intent=Intent.FAQ,
                    confidence=0.75,
                    sub_intent=faq_type,
                    method="keyword"
                )
        
        # 默认返回FAQ(最常见)
        return IntentResult(
            intent=Intent.FAQ,
            confidence=0.5,
            method="keyword"
        )
    
    @torch.no_grad()
    def _predict_with_model(self, query: str) -> IntentResult:
        """Level 3: BERT模型预测"""
        # Tokenize
        inputs = self.tokenizer(
            query,
            return_tensors="pt",
            padding=True,
            truncation=True,
            max_length=128
        ).to(self.device)
        
        # 编码
        outputs = self.model(**inputs)
        
        # 取[CLS]token的表示作为句子向量
        cls_embedding = outputs.last_hidden_state[:, 0, :]
        
        # 简化的分类逻辑(实际应该接分类头)
        # 这里用向量范数模拟置信度
        confidence = float(torch.norm(cls_embedding).item()) / 10
        
        return IntentResult(
            intent=Intent.FAQ,  # 简化处理
            confidence=min(confidence, 0.8),
            method="model"
        )
    
    def _fuse_results(
        self,
        query: str,
        results: List[IntentResult],
        context: Optional[Dict]
    ) -> IntentResult:
        """
        融合多个策略的结果
        
        融合策略:
        1. 规则匹配 > 关键词 > 模型
        2. 考虑置信度和上下文
        """
        # 过滤无效结果
        valid_results = [r for r in results if r is not None]
        
        if not valid_results:
            return IntentResult(
                intent=Intent.FAQ,
                confidence=0.5,
                method="default"
            )
        
        # 按置信度排序
        valid_results.sort(key=lambda x: x.confidence, reverse=True)
        
        # 选择最佳结果
        best = valid_results[0]
        
        # 上下文修正
        if context:
            # 如果上轮是闲聊,继续闲聊
            prev_intent = context.get("last_intent")
            if prev_intent == Intent.CHITCHAT and best.intent == Intent.FAQ:
                if best.confidence < 0.7:
                    best.intent = Intent.CHITCHAT
                    best.confidence = 0.6
                    best.method = "context_adjusted"
        
        return best
    
    def extract_slots(self, query: str, intent: Intent) -> Dict[str, str]:
        """
        槽位提取
        
        根据意图类型提取关键信息
        """
        slots = {}
        
        if intent == Intent.TASK:
            # 订单号提取
            order_pattern = r"订单[号]?[::]?\s*(\d{10,20})"
            match = re.search(order_pattern, query)
            if match:
                slots["order_id"] = match.group(1)
            
            # 手机号提取
            phone_pattern = r"1[3-9]\d{9}"
            match = re.search(phone_pattern, query)
            if match:
                slots["phone"] = match.group(0)
            
            # 时间提取
            time_patterns = [
                (r"(\d{4})年(\d{1,2})月(\d{1,2})日", "date"),
                (r"(\d{1,2})月(\d{1,2})日", "month_date"),
                (r"今天|明天|后天", "relative_date")
            ]
            for pattern, time_type in time_patterns:
                match = re.search(pattern, query)
                if match:
                    slots["time"] = match.group(0)
                    slots["time_type"] = time_type
                    break
        
        return slots


class IntentRouter:
    """
    意图路由器
    
    根据识别结果,将请求路由到不同的处理模块
    """
    
    def __init__(self):
        """初始化路由"""
        self.classifier = IntentClassifier()
        self.rag_engine = None  # 延迟初始化
        self.llm_client = None   # 延迟初始化
        
    def set_rag_engine(self, rag_engine):
        """设置RAG引擎"""
        self.rag_engine = rag_engine
    
    def set_llm_client(self, llm_client):
        """设置LLM客户端"""
        self.llm_client = llm_client
    
    async def route(self, query: str, session_id: str, context: Dict) -> Dict:
        """
        路由请求
        
        Args:
            query: 用户query
            session_id: 会话ID
            context: 上下文
            
        Returns:
            路由结果,包含响应和元信息
        """
        # 1. 意图识别
        intent_result = self.classifier.classify(query, context)
        
        # 2. 根据意图路由
        if intent_result.intent == Intent.TRANSFER:
            return await self._handle_transfer(query, session_id, intent_result)
        
        elif intent_result.intent == Intent.COMPLAINT:
            return await self._handle_complaint(query, session_id, intent_result)
        
        elif intent_result.intent == Intent.CHITCHAT:
            return await self._handle_chitchat(query, session_id)
        
        elif intent_result.intent == Intent.TASK:
            return await self._handle_task(query, session_id, intent_result)
        
        else:  # FAQ
            return await self._handle_faq(query, session_id, intent_result)
    
    async def _handle_faq(
        self,
        query: str,
        session_id: str,
        intent_result: IntentResult
    ) -> Dict:
        """处理知识问答"""
        if not self.rag_engine:
            return {"error": "RAG引擎未初始化"}
        
        # RAG检索
        search_result = self.rag_engine.search(query, top_k=5)
        
        # 提取相关文档
        contexts = [doc.content for doc in search_result.documents]
        context_text = "\n\n".join(contexts)
        
        # 构建Prompt
        prompt = f"""基于以下知识库内容回答用户问题。如果知识库中没有相关信息,请说明"抱歉,我无法从知识库中找到相关答案,建议您联系人工客服获取帮助"。

知识库内容:
{context_text}

用户问题:{query}

请给出准确、友好的回答。"""
        
        # LLM生成
        response = await self.llm_client.chat(prompt)
        
        return {
            "response": response,
            "intent": intent_result.intent.value,
            "confidence": intent_result.confidence,
            "source": "knowledge_base",
            "references": [{"content": doc.content[:100]} for doc in search_result.documents[:3]]
        }
    
    async def _handle_chitchat(self, query: str, session_id: str) -> Dict:
        """处理闲聊"""
        prompt = f"""你是一个友善的客服助手,正在与用户进行轻松的对话。请用简洁、友好的方式回复。

用户:{query}"""
        
        response = await self.llm_client.chat(prompt)
        
        return {
            "response": response,
            "intent": "chitchat",
            "source": "llm"
        }
    
    async def _handle_complaint(
        self,
        query: str,
        session_id: str,
        intent_result: IntentResult
    ) -> Dict:
        """处理投诉"""
        # 投诉需要更多同理心,先安抚情绪
        prompt = f"""用户表达了不满或投诉。请先用同理心回应,表达理解和歉意,然后再尝试提供帮助。

用户反馈:{query}

请先安抚用户情绪,然后询问具体情况。"""
        
        response = await self.llm_client.chat(prompt)
        
        # 标记需要关注
        escalation_needed = True
        
        return {
            "response": response,
            "intent": "complaint",
            "escalation_needed": escalation_needed,
            "source": "llm"
        }
    
    async def _handle_task(
        self,
        query: str,
        session_id: str,
        intent_result: IntentResult
    ) -> Dict:
        """处理任务型请求"""
        # 槽位提取
        slots = self.classifier.extract_slots(query, intent_result.intent)
        
        # 这里可以调用具体业务系统
        # 例如订单系统、物流系统等
        
        prompt = f"""用户有一个任务请求,请帮助用户完成。

任务类型:{intent_result.sub_intent}
提取的槽位信息:{slots}
用户请求:{query}

请根据槽位信息,判断是否需要补充信息,并给出回复引导用户完成操作。"""
        
        response = await self.llm_client.chat(prompt)
        
        return {
            "response": response,
            "intent": "task",
            "sub_intent": intent_result.sub_intent,
            "slots": slots,
            "source": "llm"
        }
    
    async def _handle_transfer(
        self,
        query: str,
        session_id: str,
        intent_result: IntentResult
    ) -> Dict:
        """处理转人工"""
        # 构建转人工提示
        response = "好的,我将为您转接人工客服,请稍候..."
        
        return {
            "response": response,
            "intent": "transfer",
            "action": "handoff",
            "source": "system"
        }

3.6 Django API 实现

apps/conversation/views.py - API视图

"""
客服对话 API 接口
支持实时对话、流式响应、WebSocket长连接

Author: AI Team
Date: 2025-09-15
"""

import json
import asyncio
import time
from typing import Optional
from dataclasses import dataclass

from django.http import JsonResponse, StreamingHttpResponse
from django.views import View
from django.views.decorators.csrf import csrf_exempt
from django.utils.decorators import method_decorator
from rest_framework.views import APIView
from rest_framework.response import Response
from rest_framework import status
from rest_framework.decorators import api_view
import logging

from .services.dialogue_manager import DialogueManager, DialogueState, MessageRole
from .services.intent_classifier import IntentRouter, Intent
from apps.knowledge.services.rag_engine import RAGEngine
from apps.agent.llm_client import LLMClient
from common.middleware import RateLimitMiddleware
from common.utils import generate_trace_id

logger = logging.getLogger(__name__)


@dataclass
class ChatRequest:
    """聊天请求"""
    query: str
    session_id: Optional[str] = None
    user_id: Optional[str] = None
    stream: bool = False
    metadata: dict = None


class ChatView(APIView):
    """
    对话接口
    
    POST /api/v1/chat
    {
        "query": "我想查一下订单状态",
        "session_id": "可选,用于多轮对话",
        "user_id": "可选,用户标识",
        "stream": false,
        "metadata": {}
    }
    """
    
    def __init__(self, **kwargs):
        super().__init__(**kwargs)
        self.dialogue_manager = DialogueManager()
        self.intent_router = IntentRouter()
        
        # 延迟初始化重量级组件
        self._rag_engine = None
        self._llm_client = None
    
    @property
    def rag_engine(self):
        if self._rag_engine is None:
            self._rag_engine = RAGEngine()
        return self._rag_engine
    
    @property
    def llm_client(self):
        if self._llm_client is None:
            self._llm_client = LLMClient()
        return self._llm_client
    
    def post(self, request):
        """
        处理对话请求
        """
        trace_id = generate_trace_id()
        
        try:
            # 1. 解析请求
            data = request.data
            chat_request = ChatRequest(
                query=data.get("query", "").strip(),
                session_id=data.get("session_id"),
                user_id=data.get("user_id"),
                stream=data.get("stream", False),
                metadata=data.get("metadata", {})
            )
            
            if not chat_request.query:
                return Response(
                    {"error": "query不能为空"},
                    status=status.HTTP_400_BAD_REQUEST
                )
            
            # 2. 获取或创建会话
            if chat_request.session_id:
                context = self.dialogue_manager.get_session(chat_request.session_id)
                if not context:
                    # 会话不存在或已过期,创建新会话
                    context = self.dialogue_manager.create_session(
                        user_id=chat_request.user_id,
                        metadata=chat_request.metadata
                    )
            else:
                context = self.dialogue_manager.create_session(
                    user_id=chat_request.user_id,
                    metadata=chat_request.metadata
                )
            
            # 3. 添加用户消息
            self.dialogue_manager.add_message(
                session_id=context.session_id,
                role=MessageRole.USER,
                content=chat_request.query,
                metadata={"trace_id": trace_id}
            )
            
            # 4. 设置组件引用
            self.intent_router.set_rag_engine(self.rag_engine)
            self.intent_router.set_llm_client(self.llm_client)
            
            # 5. 意图识别与路由
            context_dict = {
                "last_intent": context.intent,
                "slots": context.slots
            }
            
            # 同步执行路由(实际生产应该用异步)
            loop = asyncio.new_event_loop()
            asyncio.set_event_loop(loop)
            result = loop.run_until_complete(
                self.intent_router.route(
                    chat_request.query,
                    context.session_id,
                    context_dict
                )
            )
            loop.close()
            
            # 6. 添加AI回复到会话
            self.dialogue_manager.add_message(
                session_id=context.session_id,
                role=MessageRole.ASSISTANT,
                content=result.get("response", ""),
                metadata={
                    "trace_id": trace_id,
                    "intent": result.get("intent"),
                    "confidence": result.get("confidence")
                }
            )
            
            # 7. 更新会话状态
            self.dialogue_manager.update_session(context)
            
            # 8. 处理转人工
            if result.get("intent") == "transfer":
                handoff_result = self.dialogue_manager.transfer_to_human(
                    context.session_id,
                    reason="user_request"
                )
                result["handoff"] = handoff_result
            
            # 9. 返回响应
            return Response({
                "code": 0,
                "message": "success",
                "data": {
                    "session_id": context.session_id,
                    "response": result.get("response"),
                    "intent": result.get("intent"),
                    "confidence": result.get("confidence"),
                    "trace_id": trace_id,
                    "metadata": {
                        "response_time": result.get("response_time", 0),
                        "source": result.get("source")
                    }
                }
            })
            
        except Exception as e:
            logger.error(f"对话处理失败: {str(e)}", exc_info=True)
            return Response(
                {"error": f"服务器内部错误: {str(e)}"},
                status=status.HTTP_500_INTERNAL_SERVER_ERROR
            )


class ChatStreamView(APIView):
    """
    流式对话接口
    
    使用Server-Sent Events (SSE) 实现流式响应
    """
    
    def __init__(self, **kwargs):
        super().__init__(**kwargs)
        self.dialogue_manager = DialogueManager()
    
    def post(self, request):
        """流式响应"""
        data = request.data
        query = data.get("query", "").strip()
        session_id = data.get("session_id")
        
        if not query:
            return Response(
                {"error": "query不能为空"},
                status=status.HTTP_400_BAD_REQUEST
            )
        
        def event_stream():
            """SSE事件流"""
            try:
                # 初始化对话
                if not session_id:
                    context = self.dialogue_manager.create_session()
                    session_id = context.session_id
                    yield f"data: {json.dumps({'type': 'session', 'session_id': session_id})}\n\n"
                
                # 流式生成响应
                # 这里简化处理,实际应该使用真实的流式LLM调用
                for chunk in self._generate_stream_response(query, session_id):
                    yield f"data: {json.dumps({'type': 'content', 'content': chunk})}\n\n"
                
                # 结束信号
                yield f"data: {json.dumps({'type': 'done'})}\n\n"
                
            except Exception as e:
                yield f"data: {json.dumps({'type': 'error', 'error': str(e)})}\n\n"
        
        response = StreamingHttpResponse(
            event_stream(),
            content_type='text/event-stream'
        )
        response['Cache-Control'] = 'no-cache'
        response['X-Accel-Buffering'] = 'no'
        
        return response
    
    def _generate_stream_response(self, query: str, session_id: str):
        """生成流式响应(模拟)"""
        # 实际应该使用LLM的流式API
        llm_client = LLMClient()
        
        for chunk in llm_client.stream_chat(query):
            yield chunk


class SessionView(APIView):
    """
    会话管理接口
    """
    
    def __init__(self, **kwargs):
        super().__init__(**kwargs)
        self.dialogue_manager = DialogueManager()
    
    def get(self, request, session_id: str):
        """获取会话详情"""
        context = self.dialogue_manager.get_session(session_id)
        
        if not context:
            return Response(
                {"error": "会话不存在或已过期"},
                status=status.HTTP_404_NOT_FOUND
            )
        
        return Response({
            "code": 0,
            "data": {
                "session_id": context.session_id,
                "user_id": context.user_id,
                "state": context.state.value,
                "created_at": context.created_at,
                "updated_at": context.updated_at,
                "messages": [m.to_dict() for m in context.messages[-20:]],
                "intent": context.intent,
                "slots": context.slots
            }
        })
    
    def delete(self, request, session_id: str):
        """结束会话"""
        summary = request.data.get("summary")
        
        success = self.dialogue_manager.end_session(session_id, summary)
        
        if not success:
            return Response(
                {"error": "会话不存在"},
                status=status.HTTP_404_NOT_FOUND
            )
        
        return Response({"code": 0, "message": "会话已结束"})
    
    def get_history(self, request, session_id: str):
        """获取对话历史"""
        max_turns = int(request.query_params.get("max_turns", 10))
        
        history = self.dialogue_manager.get_conversation_history(
            session_id,
            max_turns=max_turns
        )
        
        return Response({
            "code": 0,
            "data": {
                "session_id": session_id,
                "messages": history
            }
        })


class HumanHandoffView(APIView):
    """
    人工客服转接接口
    """
    
    def __init__(self, **kwargs):
        super().__init__(**kwargs)
        self.dialogue_manager = DialogueManager()
    
    def post(self, request, session_id: str):
        """请求转人工"""
        reason = request.data.get("reason", "user_request")
        priority = request.data.get("priority", 1)
        
        # 增加转接计数
        self.dialogue_manager.increment_escalation(session_id)
        
        # 执行转接
        result = self.dialogue_manager.transfer_to_human(
            session_id,
            reason=reason,
            priority=priority
        )
        
        return Response({
            "code": 0,
            "data": result
        })
    
    def get(self, request, session_id: str):
        """获取转接状态"""
        context = self.dialogue_manager.get_session(session_id)
        
        if not context:
            return Response(
                {"error": "会话不存在"},
                status=status.HTTP_404_NOT_FOUND
            )
        
        return Response({
            "code": 0,
            "data": {
                "state": context.state.value,
                "is_human_in_charge": context.state == DialogueState.HUMAN_IN_CHARGE,
                "escalation_count": context.human_escalation_count,
                "transfer_reason": context.metadata.get("transfer_reason")
            }
        })


class HealthCheckView(APIView):
    """
    健康检查接口
    """
    
    def get(self, request):
        """检查系统健康状态"""
        from django.db import connection
        
        health = {
            "status": "healthy",
            "timestamp": time.time(),
            "components": {}
        }
        
        # 检查数据库
        try:
            with connection.cursor() as cursor:
                cursor.execute("SELECT 1")
            health["components"]["database"] = "healthy"
        except Exception as e:
            health["components"]["database"] = f"unhealthy: {str(e)}"
            health["status"] = "degraded"
        
        # 检查Redis
        try:
            from django.core.cache import cache
            cache.set("health_check", "ok", 1)
            if cache.get("health_check") == "ok":
                health["components"]["redis"] = "healthy"
            else:
                health["components"]["redis"] = "unhealthy"
                health["status"] = "degraded"
        except Exception as e:
            health["components"]["redis"] = f"unhealthy: {str(e)}"
            health["status"] = "degraded"
        
        # 检查RAG引擎
        try:
            rag_engine = RAGEngine()
            rag_health = rag_engine.health_check()
            health["components"]["rag_engine"] = rag_health
            if rag_health["status"] != "healthy":
                health["status"] = "degraded"
        except Exception as e:
            health["components"]["rag_engine"] = f"unhealthy: {str(e)}"
            health["status"] = "degraded"
        
        status_code = 200 if health["status"] == "healthy" else 503
        
        return Response(health, status=status_code)

四、知识库构建

4.1 知识库配置

apps/knowledge/management/commands/build_knowledge_base.py - 知识库构建命令

"""
知识库构建命令
支持批量导入文档、自动分块、向量化和索引构建

使用方法:
    python manage.py build_knowledge_base --source ./docs/ --batch-size 100

Author: AI Team
Date: 2025-09-15
"""

import os
import time
import hashlib
from pathlib import Path
from typing import List, Dict, Optional
import logging

from django.core.management.base import BaseCommand, CommandError
from django.conf import settings
import pymilvus
from pymilvus import Collection, CollectionSchema, FieldSchema, DataType, utility
import torch
from sentence_transformers import SentenceTransformer
import yaml

logger = logging.getLogger(__name__)


class KnowledgeBaseBuilder:
    """
    知识库构建器
    
    支持功能:
    - 多种文档格式解析(Markdown、TXT、HTML、PDF)
    - 智能分块(基于语义/段落/固定长度)
    - 批量向量化
    - Milvus索引构建
    """
    
    def __init__(self, config_path: str = None):
        """初始化"""
        # 加载配置
        if config_path and os.path.exists(config_path):
            with open(config_path, 'r', encoding='utf-8') as f:
                self.config = yaml.safe_load(f)
        else:
            self.config = self._default_config()
        
        # 初始化组件
        self._init_embedding_model()
        self._init_milvus()
    
    def _default_config(self) -> Dict:
        """默认配置"""
        return {
            "chunk_size": 512,
            "chunk_overlap": 50,
            "batch_size": 100,
            "enable_upsert": True,
            "collection_name": "knowledge_base"
        }
    
    def _init_embedding_model(self):
        """初始化Embedding模型"""
        model_name = self.config.get("embedding_model", 
                                      "paraphrase-multilingual-MiniLM-L12-v2")
        self.embedding_model = SentenceTransformer(model_name)
        
        # 获取向量维度
        test_embedding = self.embedding_model.encode("测试", convert_to_numpy=True)
        self.vector_dimension = len(test_embedding)
        
        logger.info(f"Embedding模型加载完成: {model_name}, 维度: {self.vector_dimension}")
    
    def _init_milvus(self):
        """初始化Milvus连接"""
        milvus_config = settings.VECTOR_DB
        
        self.milvus_host = milvus_config.get("host", "localhost")
        self.milvus_port = milvus_config.get("port", 19530)
        self.collection_name = self.config.get("collection_name", "knowledge_base")
        
        # 连接
        connections.connect(
            alias="default",
            host=self.milvus_host,
            port=self.milvus_port
        )
        
        logger.info(f"Milvus连接成功: {self.milvus_host}:{self.milvus_port}")
    
    def create_collection(self, drop_existing: bool = False):
        """
        创建Collection
        
        字段结构:
        - id: 文档唯一ID (VARCHAR)
        - content: 文档内容 (VARCHAR)
        - metadata: 元信息 (JSON/VARCHAR)
        - category: 分类 (VARCHAR)
        - source: 来源 (VARCHAR)
        - embedding: 向量 (FLOAT_VECTOR)
        - created_at: 创建时间 (INT64)
        - updated_at: 更新时间 (INT64)
        """
        # 检查是否已存在
        if utility.has_collection(self.collection_name):
            if drop_existing:
                utility.drop_collection(self.collection_name)
                logger.info(f"已删除现有Collection: {self.collection_name}")
            else:
                logger.info(f"Collection已存在: {self.collection_name}")
                self.collection = Collection(self.collection_name)
                self.collection.load()
                return
        
        # 定义Schema
        fields = [
            FieldSchema(name="id", dtype=DataType.VARCHAR, max_length=64, is_primary=True),
            FieldSchema(name="content", dtype=DataType.VARCHAR, max_length=4096),
            FieldSchema(name="metadata", dtype=DataType.VARCHAR, max_length=1024),
            FieldSchema(name="category", dtype=DataType.VARCHAR, max_length=64),
            FieldSchema(name="source", dtype=DataType.VARCHAR, max_length=256),
            FieldSchema(name="embedding", dtype=DataType.FLOAT_VECTOR, dim=self.vector_dimension),
            FieldSchema(name="created_at", dtype=DataType.INT64),
            FieldSchema(name="updated_at", dtype=DataType.INT64),
        ]
        
        schema = CollectionSchema(
            fields=fields,
            description="知识库集合"
        )
        
        # 创建Collection
        self.collection = Collection(
            name=self.collection_name,
            schema=schema
        )
        
        # 创建索引
        index_params = {
            "metric_type": "COSINE",
            "index_type": "HNSW",
            "params": {"M": 16, "efConstruction": 200}
        }
        
        self.collection.create_index(
            field_name="embedding",
            index_params=index_params
        )
        
        # 加载Collection
        self.collection.load()
        
        logger.info(f"Collection创建成功: {self.collection_name}")
    
    def parse_document(self, file_path: str) -> List[Dict]:
        """
        解析文档
        
        支持格式:
        - .txt: 纯文本
        - .md: Markdown
        - .html: HTML
        - .pdf: PDF (需要额外处理)
        """
        file_path = Path(file_path)
        
        if not file_path.exists():
            raise FileNotFoundError(f"文件不存在: {file_path}")
        
        suffix = file_path.suffix.lower()
        
        if suffix == '.txt':
            return self._parse_txt(file_path)
        elif suffix == '.md':
            return self._parse_markdown(file_path)
        elif suffix == '.html':
            return self._parse_html(file_path)
        else:
            logger.warning(f"不支持的文件格式: {suffix}")
            return []
    
    def _parse_txt(self, file_path: Path) -> List[Dict]:
        """解析纯文本"""
        with open(file_path, 'r', encoding='utf-8') as f:
            content = f.read()
        
        return [{
            "content": content,
            "metadata": {
                "filename": file_path.name,
                "filepath": str(file_path)
            }
        }]
    
    def _parse_markdown(self, file_path: Path) -> List[Dict]:
        """解析Markdown"""
        import re
        
        with open(file_path, 'r', encoding='utf-8') as f:
            content = f.read()
        
        # 提取标题和内容
        sections = []
        current_title = "概述"
        current_content = []
        
        for line in content.split('\n'):
            # 标题行
            title_match = re.match(r'^(#{1,6})\s+(.+)', line)
            if title_match:
                # 保存上一个section
                if current_content:
                    sections.append({
                        "title": current_title,
                        "content": '\n'.join(current_content)
                    })
                
                current_title = title_match.group(2)
                current_content = []
            else:
                current_content.append(line)
        
        # 最后一个section
        if current_content:
            sections.append({
                "title": current_title,
                "content": '\n'.join(current_content)
            })
        
        return [{
            "content": f"{s['title']}: {s['content']}",
            "metadata": {
                "filename": file_path.name,
                "title": s['title']
            }
        } for s in sections]
    
    def _parse_html(self, file_path: Path) -> List[Dict]:
        """解析HTML"""
        from bs4 import BeautifulSoup
        
        with open(file_path, 'r', encoding='utf-8') as f:
            soup = BeautifulSoup(f.read(), 'html.parser')
        
        # 提取文本
        texts = []
        for p in soup.find_all(['p', 'h1', 'h2', 'h3', 'h4', 'li']):
            text = p.get_text(strip=True)
            if text and len(text) > 10:
                texts.append(text)
        
        return [{
            "content": '\n'.join(texts),
            "metadata": {
                "filename": file_path.name,
                "title": soup.title.string if soup.title else file_path.stem
            }
        }]
    
    def chunk_text(self, text: str, chunk_size: int = None, 
                   overlap: int = None) -> List[str]:
        """
        文本分块
        
        策略:
        1. 按段落分割(优先)
        2. 按句子分割
        3. 固定长度分割(兜底)
        """
        chunk_size = chunk_size or self.config.get("chunk_size", 512)
        overlap = overlap or self.config.get("chunk_overlap", 50)
        
        # 按段落分割
        paragraphs = [p.strip() for p in text.split('\n\n') if p.strip()]
        
        chunks = []
        current_chunk = []
        current_length = 0
        
        for para in paragraphs:
            para_length = len(para)
            
            # 如果单个段落超过chunk_size,需要进一步分割
            if para_length > chunk_size:
                # 保存当前chunk
                if current_chunk:
                    chunks.append('\n'.join(current_chunk))
                    # 保留overlap
                    current_chunk = current_chunk[-2:] if len(current_chunk) >= 2 else []
                    current_length = sum(len(c) for c in current_chunk)
                
                # 分割长段落
                sentences = para.split('。')
                for sentence in sentences:
                    if current_length + len(sentence) > chunk_size:
                        if current_chunk:
                            chunks.append('\n'.join(current_chunk))
                        current_chunk = [sentence]
                        current_length = len(sentence)
                    else:
                        current_chunk.append(sentence)
                        current_length += len(sentence)
            
            # 正常段落
            elif current_length + para_length > chunk_size:
                chunks.append('\n'.join(current_chunk))
                # 保留overlap
                overlap_text = '\n'.join(current_chunk)[-overlap:]
                current_chunk = [overlap_text, para] if overlap_text else [para]
                current_length = len(overlap_text) + para_length
            else:
                current_chunk.append(para)
                current_length += para_length
        
        # 最后一个chunk
        if current_chunk:
            chunks.append('\n'.join(current_chunk))
        
        return [c for c in chunks if c]
    
    def encode_documents(self, documents: List[Dict]) -> List[Dict]:
        """
        向量化文档
        """
        contents = [doc["content"] for doc in documents]
        
        # 批量编码
        embeddings = self.embedding_model.encode(
            contents,
            batch_size=self.config.get("batch_size", 100),
            show_progress_bar=True,
            convert_to_numpy=True,
            normalize_embeddings=True
        )
        
        # 添加向量到文档
        for doc, embedding in zip(documents, embeddings):
            doc["embedding"] = embedding.tolist()
        
        return documents
    
    def insert_documents(self, documents: List[Dict], category: str = "general"):
        """
        插入文档到Milvus
        """
        if not documents:
            return 0
        
        import json
        
        # 准备数据
        ids = []
        contents = []
        metadatas = []
        categories = []
        sources = []
        embeddings = []
        created_ats = []
        updated_ats = []
        
        timestamp = int(time.time())
        
        for doc in documents:
            # 生成唯一ID
            content_hash = hashlib.md5(doc["content"].encode()).hexdigest()
            doc_id = f"{category}_{timestamp}_{content_hash[:8]}"
            
            ids.append(doc_id)
            contents.append(doc["content"][:4096])  # 限制长度
            metadatas.append(json.dumps(doc.get("metadata", {}), ensure_ascii=False))
            categories.append(category)
            sources.append(doc.get("metadata", {}).get("source", "manual"))
            embeddings.append(doc["embedding"])
            created_ats.append(timestamp)
            updated_ats.append(timestamp)
        
        # 批量插入
        data = [ids, contents, metadatas, categories, sources, embeddings, created_ats, updated_ats]
        
        self.collection.insert(data)
        self.collection.flush()
        
        logger.info(f"成功插入 {len(documents)} 篇文档")
        
        return len(documents)
    
    def build_from_directory(self, directory: str, category: str = "general",
                             recursive: bool = True, file_types: List[str] = None):
        """
        从目录构建知识库
        """
        directory = Path(directory)
        
        if not directory.exists():
            raise FileNotFoundError(f"目录不存在: {directory}")
        
        file_types = file_types or ['.txt', '.md', '.html']
        
        # 收集所有文件
        files = []
        if recursive:
            for ext in file_types:
                files.extend(directory.rglob(f"*{ext}"))
        else:
            for ext in file_types:
                files.extend(directory.glob(f"*{ext}"))
        
        logger.info(f"找到 {len(files)} 个文档文件")
        
        # 处理每个文件
        total_chunks = 0
        for file_path in files:
            try:
                # 解析文档
                documents = self.parse_document(str(file_path))
                
                # 分块
                all_chunks = []
                for doc in documents:
                    chunks = self.chunk_text(doc["content"])
                    for chunk in chunks:
                        all_chunks.append({
                            "content": chunk,
                            "metadata": {
                                **doc["metadata"],
                                "file": str(file_path)
                            }
                        })
                
                # 向量化
                if all_chunks:
                    all_chunks = self.encode_documents(all_chunks)
                    
                    # 插入
                    count = self.insert_documents(all_chunks, category=category)
                    total_chunks += count
                    
                    logger.info(f"处理完成: {file_path.name}, chunks: {count}")
                
            except Exception as e:
                logger.error(f"处理文件失败: {file_path}, error: {e}")
        
        logger.info(f"知识库构建完成,总计插入 {total_chunks} 个文档块")
        
        return total_chunks


class Command(BaseCommand):
    """Django管理命令"""
    
    help = '构建知识库索引'
    
    def add_arguments(self, parser):
        parser.add_argument(
            '--source',
            type=str,
            help='文档目录路径'
        )
        parser.add_argument(
            '--category',
            type=str,
            default='general',
            help='文档分类'
        )
        parser.add_argument(
            '--config',
            type=str,
            help='配置文件路径'
        )
        parser.add_argument(
            '--drop',
            action='store_true',
            help='删除现有Collection后重建'
        )
        parser.add_argument(
            '--batch-size',
            type=int,
            default=100,
            help='批处理大小'
        )
    
    def handle(self, *args, **options):
        source = options.get('source')
        category = options.get('category', 'general')
        config_path = options.get('config')
        drop = options.get('drop', False)
        batch_size = options.get('batch_size', 100)
        
        self.stdout.write(f"开始构建知识库...")
        
        try:
            # 创建构建器
            builder = KnowledgeBaseBuilder(config_path)
            builder.config["batch_size"] = batch_size
            
            # 创建Collection
            builder.create_collection(drop_existing=drop)
            
            # 如果指定了源目录,导入文档
            if source:
                count = builder.build_from_directory(source, category=category)
                self.stdout.write(
                    self.style.SUCCESS(f'知识库构建完成,共导入 {count} 个文档块')
                )
            else:
                self.stdout.write(
                    self.style.SUCCESS('Collection创建成功,未导入文档')
                )
                
        except Exception as e:
            raise CommandError(f'知识库构建失败: {e}')

五、Kubernetes 部署配置

5.1 部署配置

k8s/deployment.yaml - Kubernetes部署配置

# 企业智能客服系统 - Kubernetes部署配置
# 版本: v2.0.0
# 更新时间: 2025-09-15

---
# API服务 Deployment
apiVersion: apps/v1
kind: Deployment
metadata:
  name: customer-service-api
  namespace: customer-service
  labels:
    app: customer-service
    component: api
    version: v2.0.0
spec:
  replicas: 3
  selector:
    matchLabels:
      app: customer-service
      component: api
  strategy:
    type: RollingUpdate
    rollingUpdate:
      maxSurge: 1
      maxUnavailable: 0
  template:
    metadata:
      labels:
        app: customer-service
        component: api
        version: v2.0.0
      annotations:
        prometheus.io/scrape: "true"
        prometheus.io/port: "8000"
        prometheus.io/path: "/metrics"
    spec:
      # 亲和性配置:优先调度到GPU节点
      affinity:
        nodeAffinity:
          preferredDuringSchedulingIgnoredDuringExecution:
          - weight: 100
            preference:
              matchExpressions:
              - key: node-type
                operator: In
                values:
                - gpu-node
      #  Tolerations:允许调度到污点节点
      tolerations:
      - key: "gpu"
        operator: "Exists"
        effect: "NoSchedule"
      # 容器配置
      containers:
      - name: api
        image: registry.example.com/customer-service/api:v2.0.0
        imagePullPolicy: Always
        ports:
        - name: http
          containerPort: 8000
          protocol: TCP
        - name: grpc
          containerPort: 50051
          protocol: TCP
        # 环境变量
        env:
        - name: DJANGO_SETTINGS_MODULE
          value: config.settings.production
        - name: LOG_LEVEL
          value: "INFO"
        - name: WORKERS
          value: "4"
        - name: TIMEOUT
          value: "120"
        # 资源配置
        resources:
          requests:
            cpu: "500m"
            memory: "1Gi"
          limits:
            cpu: "2000m"
            memory: "4Gi"
        # 健康检查
        livenessProbe:
          httpGet:
            path: /health/
            port: http
          initialDelaySeconds: 30
          periodSeconds: 10
          timeoutSeconds: 5
          failureThreshold: 3
        readinessProbe:
          httpGet:
            path: /ready/
            port: http
          initialDelaySeconds: 10
          periodSeconds: 5
          timeoutSeconds: 3
          failureThreshold: 3
        # 挂载配置
        volumeMounts:
        - name: config
          mountPath: /app/config/settings.yaml
          subPath: settings.yaml
        - name: logs
          mountPath: /var/log
        # 环境特定配置
        envFrom:
        - configMapRef:
            name: customer-service-config
        - secretRef:
            name: customer-service-secrets
      # 初始化容器
      initContainers:
      - name: wait-for-dependencies
        image: busybox:1.36
        command:
        - sh
        - -c
        - |
          echo "Waiting for Redis..."
          until nc -z $REDIS_HOST $REDIS_PORT; do
            echo "Redis not ready, waiting..."
            sleep 2
          done
          echo "Redis is ready!"
          echo "Waiting for Milvus..."
          until nc -z $MILVUS_HOST $MILVUS_PORT; do
            echo "Milvus not ready, waiting..."
            sleep 2
          done
          echo "Milvus is ready!"
        envFrom:
        - configMapRef:
            name: customer-service-config
      volumes:
      - name: config
        configMap:
          name: customer-service-config
      - name: logs
        emptyDir: {}
      # 服务账户
      serviceAccountName: customer-service-sa
      securityContext:
        runAsNonRoot: true
        runAsUser: 1000
        fsGroup: 1000

---
# Worker服务 Deployment (处理异步任务)
apiVersion: apps/v1
kind: Deployment
metadata:
  name: customer-service-worker
  namespace: customer-service
  labels:
    app: customer-service
    component: worker
spec:
  replicas: 2
  selector:
    matchLabels:
      app: customer-service
      component: worker
  template:
    metadata:
      labels:
        app: customer-service
        component: worker
    spec:
      containers:
      - name: worker
        image: registry.example.com/customer-service/worker:v2.0.0
        imagePullPolicy: Always
        # 启动命令
        command: ["python", "manage.py", "celery", "worker"]
        args:
        - "--concurrency=4"
        - "--loglevel=info"
        - "--pool=prefork"
        env:
        - name: DJANGO_SETTINGS_MODULE
          value: config.settings.production
        resources:
          requests:
            cpu: "1000m"
            memory: "2Gi"
          limits:
            cpu: "2000m"
            memory: "4Gi"
      serviceAccountName: customer-service-sa

---
# Redis Deployment
apiVersion: apps/v1
kind: Deployment
metadata:
  name: redis
  namespace: customer-service
spec:
  replicas: 1
  selector:
    matchLabels:
      app: redis
  template:
    metadata:
      labels:
        app: redis
    spec:
      containers:
      - name: redis
        image: redis:7.2-alpine
        ports:
        - containerPort: 6379
        command: ["redis-server", "--appendonly", "yes"]
        resources:
          requests:
            cpu: "100m"
            memory: "512Mi"
          limits:
            cpu: "500m"
            memory: "2Gi"
        volumeMounts:
        - name: redis-data
          mountPath: /data
      volumes:
      - name: redis-data
        persistentVolumeClaim:
          claimName: redis-pvc

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

---
# PostgreSQL Deployment (可选,使用云数据库时可省略)
apiVersion: apps/v1
kind: Deployment
metadata:
  name: postgresql
  namespace: customer-service
spec:
  replicas: 1
  selector:
    matchLabels:
      app: postgresql
  template:
    metadata:
      labels:
        app: postgresql
    spec:
      containers:
      - name: postgresql
        image: postgres:15-alpine
        ports:
        - containerPort: 5432
        env:
        - name: POSTGRES_DB
          value: customer_service
        - name: POSTGRES_USER
          valueFrom:
            secretKeyRef:
              name: customer-service-secrets
              key: db_user
        - name: POSTGRES_PASSWORD
          valueFrom:
            secretKeyRef:
              name: customer-service-secrets
              key: db_password
        resources:
          requests:
            cpu: "250m"
            memory: "512Mi"
          limits:
            cpu: "1000m"
            memory: "2Gi"
        volumeMounts:
        - name: postgres-data
          mountPath: /var/lib/postgresql/data
      volumes:
      - name: postgres-data
        persistentVolumeClaim:
          claimName: postgresql-pvc

---
# API Service
apiVersion: v1
kind: Service
metadata:
  name: customer-service-api
  namespace: customer-service
  labels:
    app: customer-service
spec:
  type: ClusterIP
  ports:
  - name: http
    port: 8000
    targetPort: 8000
    protocol: TCP
  - name: grpc
    port: 50051
    targetPort: 50051
    protocol: TCP
  selector:
    app: customer-service
    component: api

---
# HPA 自动扩缩容配置
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
  name: customer-service-api-hpa
  namespace: customer-service
spec:
  scaleTargetRef:
    apiVersion: apps/v1
    kind: Deployment
    name: customer-service-api
  minReplicas: 3
  maxReplicas: 20
  metrics:
  - type: Resource
    resource:
      name: cpu
      target:
        type: Utilization
        averageUtilization: 70
  - type: Resource
    resource:
      name: memory
      target:
        type: Utilization
        averageUtilization: 80
  behavior:
    scaleUp:
      stabilizationWindowSeconds: 60
      policies:
      - type: Percent
        value: 100
        periodSeconds: 60
    scaleDown:
      stabilizationWindowSeconds: 300
      policies:
      - type: Percent
        value: 10
        periodSeconds: 60

5.2 ConfigMap和Secret配置

k8s/configmap.yaml - 配置和密钥

---
# ConfigMap 配置
apiVersion: v1
kind: ConfigMap
metadata:
  name: customer-service-config
  namespace: customer-service
data:
  # Django设置
  DJANGO_SETTINGS_MODULE: "config.settings.production"
  DJANGO_SECRET_KEY: "$(DJANGO_SECRET_KEY)"
  
  # Redis配置
  REDIS_HOST: "redis"
  REDIS_PORT: "6379"
  REDIS_DB: "0"
  
  # Milvus配置
  MILVUS_HOST: "milvus.milvus.svc.cluster.local"
  MILVUS_PORT: "19530"
  
  # LLM配置
  LLM_PROVIDER: "dashscope"
  LLM_MODEL: "qwen-max"
  
  # 日志配置
  LOG_LEVEL: "INFO"
  LOG_FORMAT: "%(asctime)s - %(name)s - %(levelname)s - %(message)s"
  
---
# Secret 配置
apiVersion: v1
kind: Secret
metadata:
  name: customer-service-secrets
  namespace: customer-service
type: Opaque
stringData:
  # 数据库
  db_user: "cs_user"
  db_password: "$(DB_PASSWORD)"
  db_host: "$(DB_HOST)"
  db_name: "customer_service"
  
  # Redis密码
  redis_password: "$(REDIS_PASSWORD)"
  
  # LLM API Key
  dashscope_api_key: "$(DASHSCOPE_API_KEY)"
  
  # Django Secret
  django_secret_key: "$(DJANGO_SECRET_KEY)"

---
# ServiceAccount
apiVersion: v1
kind: ServiceAccount
metadata:
  name: customer-service-sa
  namespace: customer-service

---
# RBAC Role
apiVersion: rbac.authorization.k8s.io/v1
kind: Role
metadata:
  name: customer-service-role
  namespace: customer-service
rules:
- apiGroups: [""]
  resources: ["configmaps", "secrets"]
  verbs: ["get", "list", "watch"]
- apiGroups: [""]
  resources: ["services"]
  verbs: ["get", "list", "watch", "create", "update"]

---
# RBAC RoleBinding
apiVersion: rbac.authorization.k8s.io/v1
kind: RoleBinding
metadata:
  name: customer-service-rolebinding
  namespace: customer-service
subjects:
- kind: ServiceAccount
  name: customer-service-sa
  namespace: customer-service
roleRef:
  kind: Role
  name: customer-service-role
  apiGroup: rbac.authorization.k8s.io

---
# PodDisruptionBudget
apiVersion: policy/v1
kind: PodDisruptionBudget
metadata:
  name: customer-service-api-pdb
  namespace: customer-service
spec:
  minAvailable: 2
  selector:
    matchLabels:
      app: customer-service
      component: api

六、系统优缺点分析

6.1 优点

优点

说明

架构解耦

各模块独立,便于维护和扩展

RAG增强

回答基于真实知识库,可溯源,减少幻觉

多策略意图识别

规则+语义+模型融合,准确率高

多轮对话支持

完整的上下文管理,支持复杂场景

人机协作

AI无法处理时自动转人工,体验好

可扩展性

支持水平扩展,应对流量高峰

可观测性

完善的监控和日志,便于问题排查

6.2 缺点与挑战

缺点/挑战

说明

缓解方案

响应延迟

RAG检索+LLM生成需要数秒

流式输出、预热缓存

向量数据库成本

Milvus需要独立部署

使用云服务或简化方案

模型依赖

受限于LLM能力

Prompt优化、模型选型

知识库维护

需要持续更新和优化

自动化ETL流程

冷启动问题

新品类知识可能不足

人工标注+迁移学习

复杂逻辑处理

多轮复杂对话仍有局限

人工接管机制


七、后续优化方向

7.1 短期优化(1-3个月)

  1. 性能优化

  • 引入Response Caching减少重复查询

  • 优化Embedding模型,缩短向量化时间

  • 实现Query改写,提升检索召回率

  1. 效果优化

  • 收集用户反馈数据,优化意图分类模型

  • 扩充知识库覆盖范围

  • 优化Prompt模板

  1. 稳定性提升

  • 完善熔断降级机制

  • 增加多级缓存策略

  • 优化错误处理

7.2 中期规划(3-6个月)

  1. 智能化提升

  • 引入Agent架构,支持复杂任务拆解

  • 实现个性化推荐

  • 开发情感识别能力

  1. 多模态支持

  • 支持图片理解

  • 支持语音输入/输出

  • 支持文档上传分析

  1. 数据分析

  • 构建数据标注平台

  • 实现对话质量自动评估

  • 开发运营Dashboard

7.3 长期愿景(6-12个月)

  1. 全渠道覆盖

  • 打通电话/短信渠道

  • 支持企业微信/钉钉深度集成

  • 实现跨渠道上下文共享

  1. 主动服务

  • 基于用户行为预测需求

  • 主动推送服务提醒

  • 智能营销推荐

  1. 行业深耕

  • 针对电商、金融、教育等行业的专属优化

  • 沉淀行业知识图谱

  • 提供行业解决方案


八、总结

本文详细介绍了企业级智能客服系统的完整实现方案,涵盖:

  1. 需求分析:深入剖析传统客服痛点,明确技术选型

  2. 架构设计:构建完整的系统架构,支持高并发、高可用

  3. 核心实现:RAG检索、多轮对话、意图识别、人机协作

  4. 部署运维:K8s部署、弹性扩缩容、监控告警

该方案已在实际生产环境中验证,单实例可支持 5000+ 并发会话,平均响应时间 <3秒**,知识问答准确率 **>85%

后续将持续优化系统效果,探索更多智能化能力,为企业创造更大价值。

0
  1. 支付宝打赏

    qrcode alipay
  2. 微信打赏

    qrcode weixin

评论区