餐厅问答智能体构建全流程指南:从需求分析到部署上线,手把手教你打造餐饮行业AI客服

一、故事的开端:一家连锁餐厅的数字化转型之痛

2024年深秋的一个傍晚,我接到了老朋友张磊的电话。张磊是一家名为"味之境"的中型连锁餐饮品牌的CTO,旗下拥有23家门店,覆盖北京、上海、成都三个城市。电话那头的声音带着疲惫:"兄弟,我快被客服问题逼疯了。"

餐厅AI智能体架构

图1:餐厅AI智能体分层架构——从用户交互到后端服务的完整链路

核心洞察:AI智能体在餐饮行业的价值不在于炫技,而在于解决真实的业务痛点。23家门店每天接到的重复性问题超过800个,这些问题的答案90%以上可以在知识库中找到。用AI自动处理这些问题,不仅能节省人力成本,更能让门店经理把时间花在真正需要人工判断的复杂场景上。

事情的起因是这样的:随着品牌扩张,每天涌入的顾客咨询量从三年前的日均200条飙升到日均3500条。这些咨询五花八门——"你们朝阳大悦城店今天营业吗?""菜单上有无 gluten-free 的选项吗?""我上周六在你们成都太古里店消费,小票上多收了一份茶位费,怎么退?""你们支持企业团餐预订吗?最少多少人起订?"问题种类繁多,而且很多涉及具体门店的实时信息。

张磊团队最初尝试了传统的关键词匹配客服系统。结果呢?准确率不到40%,顾客投诉不减反增。后来又上了基于规则的对话树系统,维护成本高得离谱——每次菜单调整、门店信息变更、促销活动上线,都要手动更新几十条规则分支。最让张磊崩溃的是,三名客服专员的离职导致整个规则体系无人维护,系统几乎瘫痪。

"我听说大模型可以做智能客服,但完全不知道从哪下手。你能帮我从零搭一个吗?"张磊问我。

这个问题,正是本文要解答的核心。接下来,我将完整复盘我们为"味之境"构建餐厅问答智能体的全过程——从需求分析、架构设计、技术选型,到核心代码实现和部署上线。这不是一篇纸上谈兵的理论文章,而是一个真实项目的实战记录。每一步踩过的坑、做过的权衡、最终的选择,我都会毫无保留地分享出来。

二、需求分析:别急着写代码,先想清楚要解决什么问题

2.1 业务场景梳理

在动手之前,我和张磊团队花了整整两天时间做需求梳理。这是整个项目中最重要的一步——很多人一上来就想着选什么大模型、用什么框架,却忽略了最根本的问题:你的智能体到底要回答什么?

通过分析"味之境"过去三个月的客服工单记录(共计约32万条),我们将顾客咨询归类为以下七大场景:

场景类别 占比 典型问题 难度评估 时效性要求
门店信息查询 28% 营业时间、地址、电话、是否停车 高(实时变更)
菜单与菜品咨询 22% 菜品推荐、过敏原、辣度、价格 中(随季调整)
预订与排队 18% 预订规则、排队叫号、包间情况 高(实时状态)
订单与支付问题 12% 退款、发票、优惠券使用 高(需系统联动)
促销活动咨询 10% 满减规则、会员权益、节日活动 中(活动周期内)
投诉与反馈 6% 菜品质量、服务态度、环境卫生 高(需人工介入)
其他闲聊 4% 你是机器人吗、推荐附近景点

这张表看起来简单,但它直接决定了后续的架构设计和技术选型。比如,"门店信息查询"占比最大且时效性要求高,这意味着我们需要一个实时同步的门店信息知识库;"订单与支付问题"难度高且需要系统联动,这意味着智能体不能只靠检索回答,还得有调用业务API的能力;"投诉与反馈"需要人工介入,这意味着必须设计优雅的转人工机制。

2.2 从业务需求到技术需求的转化

梳理完业务场景后,下一步是将业务需求翻译成技术需求。这个转化过程至关重要——它决定了我们后续的技术方案是否"够用但不过度"。

业务需求 技术需求 关键指标
准确回答顾客问题 知识检索 + 大模型生成 准确率 ≥ 85%
响应速度快 推理优化 + 缓存策略 P95 响应时间 ≤ 3秒
能查实时信息 API 工具调用能力 API 调用成功率 ≥ 99%
能处理多轮对话 对话状态管理 上下文窗口 ≥ 10轮
高并发可用 分布式部署 + 负载均衡 日活支撑 ≥ 10万次对话
可持续运营 知识库可视化管理 运营人员5分钟内完成知识更新
安全合规 敏感信息脱敏 + 审计日志 100% 审计可追溯

2.3 一个容易被忽视的需求:降级策略

在需求分析阶段,有一个需求差点被遗漏,后来却在实际运行中救了我们的命——降级策略。

张磊最初的需求里没有这条。是在我追问"如果大模型API挂了怎么办?"时,他才意识到这个问题的严重性。餐饮行业的客服高峰集中在午市(11:00-13:00)和晚市(17:00-21:00),如果在这个时段智能体宕机,大量顾客咨询会直接涌入人工客服,而人工客服的承载力只有智能体的十分之一。

所以我们增加了一条技术需求:当大模型API不可用时,系统自动降级为基于关键词匹配的FAQ检索模式,虽然体验差一些,但至少能保证基本可用。这个决策后来在上线第二个月遇到了一次大模型服务商区域性故障时发挥了关键作用——降级模式扛了40分钟,虽然回答质量下降,但系统没有完全瘫痪。

三、架构设计:分层解耦,让每一层都可独立演进

3.1 整体架构概览

经过需求分析,我们确定了整体架构采用五层分层设计。每一层职责清晰、接口明确,可以独立演进和替换。这个设计原则在后续的迭代中证明了自己的价值——我们在上线后替换了底层大模型服务商,整个过程只改动了一个适配层,其他模块完全不受影响。

架构层级 核心职责 关键技术组件 可独立替换
接入层 多渠道消息接入与统一调度 API Gateway、WebSocket
对话编排层 意图识别、对话管理、流程编排 LangChain、状态机引擎
知识检索层 向量检索、知识库管理、RAG Milvus、Embedding模型
能力执行层 工具调用、API编排、业务系统对接 Function Calling、微服务
基础设施层 模型推理、存储、监控、日志 GPU集群、Redis、Prometheus

3.2 架构详细设计

让我们逐层拆解这个架构的设计思路和关键决策。

接入层:统一入口,多渠道适配

接入层是整个系统的"前台",负责接收来自微信小程序、App、官网、第三方外卖平台的用户消息。这里的核心设计是"协议统一化"——无论消息来自哪个渠道,接入层都将其转化为统一的内部消息格式,后续层无需关心消息来源。

我们选择用Go语言编写接入层,原因是Go在处理高并发网络请求方面表现优异,goroutine的轻量级特性让我们可以轻松应对高峰时段的并发连接。同时,Go的编译型特性保证了较低的内存占用和快速的启动时间。

对话编排层:智能体的"大脑"

对话编排层是整个系统的核心,它负责理解用户意图、管理对话状态、编排执行流程。这一层我们采用了"意图路由 + 状态机"的混合模式:

  • 首先通过轻量级意图分类器快速判断用户意图类型(FAQ查询、预订、投诉、闲聊等)
  • 根据意图类型进入不同的对话流程
  • 在流程执行过程中,通过状态机管理对话进展,支持多轮交互
  • 对于复杂场景,动态编排知识检索和工具调用的执行顺序

这个设计的好处是:大部分简单问题可以走快速路径(意图分类 → 知识检索 → 直接回答),只有复杂问题才需要走完整的编排流程,有效降低了平均响应延迟。

知识检索层:RAG的核心引擎

知识检索层负责管理和检索餐厅相关的所有知识——门店信息、菜单详情、促销规则、FAQ等。我们采用RAG(Retrieval-Augmented Generation)架构,将知识检索与大模型生成相结合,既能保证回答的准确性,又能利用大模型的语言理解能力生成自然流畅的回复。

这一层的关键技术选型是向量数据库。我们需要一个能支撑百万级向量、支持毫秒级检索、具备高可用能力的向量数据库。经过对比,我们最终选择了Milvus(后面会详细说明选型理由)。

能力执行层:让智能体"长出手脚"

仅有知识检索是不够的。当顾客问"你们朝阳大悦城店现在排队几桌?"时,这不是知识库能回答的——这需要实时查询排队系统。能力执行层就是为这类场景设计的,它通过Function Calling机制让大模型能够调用外部API和业务系统。

我们为智能体配置了以下工具能力:

工具名称 功能描述 调用频率 平均延迟
query_store_info 查询门店实时信息 高(每分钟50+次) 120ms
check_queue_status 查询排队状态 高(每分钟30+次) 200ms
create_reservation 创建预订订单 中(每小时20+次) 500ms
query_order_detail 查询订单详情 中(每小时15+次) 300ms
process_refund 发起退款流程 低(每天5+次) 800ms
search_menu 搜索菜品信息 高(每分钟40+次) 150ms

基础设施层:稳定运行的基石

基础设施层提供模型推理、数据存储、监控告警等基础能力。这一层的设计原则是"可观测性优先"——我们希望在系统运行的每一刻,都能清楚地知道每个环节的健康状态、性能指标和资源消耗。因为对于线上系统来说,"不知道出了什么问题"比"出了问题"更可怕。

3.3 数据安全与隐私保护设计

在餐饮行业,智能体处理的数据涉及大量用户隐私——手机号、消费记录、预订信息等。如果在架构设计阶段不考虑数据安全,上线后才发现问题,修复成本将是巨大的。我们在架构层面做了以下安全设计:

安全层级 防护措施 实现方式
传输安全 全链路HTTPS加密 API Gateway强制TLS 1.3,内部服务间mTLS
存储安全 敏感字段加密存储 用户手机号、地址等AES-256加密后入库
处理安全 大模型调用前脱敏 PII信息在进入LLM前替换为占位符
审计安全 全量对话日志审计 每条对话记录意图、工具调用、回答内容,保留90天
访问安全 RBAC权限控制 知识库管理、日志查看等操作需对应角色权限

其中,"大模型调用前脱敏"这一层尤其关键。当我们把用户消息发送给闭源大模型API时,消息中可能包含用户的手机号、姓名等信息。我们的做法是:在消息进入编排层之前,先用正则和NER模型识别出PII信息,替换为占位符(如[PHONE]、[NAME]),然后在生成回答时再还原。这样即使大模型服务商记录了请求内容,也无法获取到用户的真实隐私信息。

张磊最初觉得这一层"多此一举"——"大模型厂商应该有自己的隐私保护吧?"我给他算了一笔账:一旦用户隐私泄露,根据《个人信息保护法》,企业面临的罚款最高可达5000万元或上一年度营业额的5%。对于一家年营收约2亿的餐饮品牌来说,这个风险完全不值得冒。张磊听完立刻拍板同意了这套方案。

3.4 数据流设计

让我们通过一个完整的请求流程来理解系统的数据流:

  1. 用户在微信小程序输入:"你们朝阳大悦城店现在排队几桌?"
  2. 接入层接收消息,转换为统一格式,写入消息队列
  3. 对话编排层消费消息,意图分类器判断为"排队查询"意图
  4. 编排引擎决定执行路径:先调用 check_queue_status 工具查询实时排队状态
  5. 能力执行层调用排队系统API,获取"当前排队8桌,预计等待35分钟"
  6. 编排引擎将工具返回结果连同用户问题一起送入大模型
  7. 大模型生成自然语言回复:"您好,朝阳大悦城店目前排队8桌,预计等待时间约35分钟。您可以在小程序上提前取号哦~"
  8. 接入层将回复推送给用户,同时记录对话日志

整个流程从用户发送消息到收到回复,P95延迟控制在2.8秒以内。这个数字看起来不算快,但对于需要调用外部API和等待大模型生成的场景来说,已经是比较理想的水平了。

四、技术选型:每个决策背后的权衡

技术选型是整个项目中最耗费精力的环节之一。每一个选择都不是拍脑袋决定的,而是经过多维度对比、实际测试后的结果。我将几个关键选型决策的对比过程完整记录下来,希望能为你的项目提供参考。

RAG知识库构建流程

图2:RAG知识库构建流程——准确率从60%到91%的优化路径

范式变革:RAG架构的核心优势是"知识的可维护性"。传统FAQ机器人每次更新都需要重新训练模型,而RAG只需要往向量数据库里添加新文档。这意味着餐厅新增菜品、修改价格、调整营业时间,客服系统可以在分钟级别内更新知识,而不是等待数天的模型重训。

4.1 大语言模型选型

大模型是智能体的核心引擎,选型时我们重点考察了中文理解能力、推理质量、API稳定性、响应速度和成本五个维度。我们对四个主流方案进行了为期两周的平行测试,使用"味之境"的真实客服数据进行评估。

评估维度 方案A:自研微调 方案B:闭源大模型API 方案C:开源模型本地部署 方案D:闭源+开源混合
中文理解能力 ★★★★☆(领域内优秀) ★★★★★(通用优秀) ★★★★☆(通用良好) ★★★★★(互补优势)
推理质量 ★★★☆☆(泛化弱) ★★★★★(最强) ★★★★☆(良好) ★★★★★(最强)
API稳定性 ★★★★★(完全自控) ★★★★☆(依赖供应商) ★★★★★(完全自控) ★★★★☆(混合依赖)
响应速度 ★★★★☆(可优化) ★★★★☆(网络依赖) ★★★☆☆(受GPU限制) ★★★★☆(可路由优化)
月度成本(估算) ¥80,000(GPU+人力) ¥45,000(按量计费) ¥60,000(GPU租赁) ¥50,000(混合计费)
数据隐私 ★★★★★(完全私有) ★★★☆☆(需审查条款) ★★★★★(完全私有) ★★★★☆(部分私有)
运维复杂度 高(需ML团队) 低(API调用) 高(需GPU运维) 中(需路由策略)

经过深入对比,我们最终选择了方案D——闭源+开源混合方案。具体策略是:

  • 日常对话、简单FAQ使用开源模型本地部署,降低成本
  • 复杂推理、投诉处理等高质量要求的场景路由到闭源大模型API
  • 通过一个智能路由层根据问题复杂度自动选择模型
  • 闭源API的数据传输做了脱敏处理,不传输用户PII信息

这个方案的核心优势是"成本可控、质量有保、降级有路"。当闭源API不可用时,所有流量自动降级到本地开源模型,虽然质量有所下降,但系统保持可用。

4.2 向量数据库选型

向量数据库是RAG架构的基石。我们考察了五个主流方案,重点关注检索性能、扩展性、运维成本和生态成熟度。

向量数据库 检索算法 最大向量数 QPS(百万级) 部署复杂度 是否支持标量过滤 社区活跃度
Milvus HNSW/IVF/ANNOY 十亿级 ~8000 中(需etcd+MinIO) 是(原生支持) 非常高
Weaviate HNSW 亿级 ~5000 低(单二进制)
Qdrant HNSW 亿级 ~6000 低(Rust实现)
pgvector IVFFlat/HNSW 千万级 ~2000 低(PostgreSQL扩展) 是(SQL原生)
FAISS IVF/HNSW/PQ 十亿级 ~10000 高(需自建服务层) 否(需额外实现) 中(库而非服务)

最终我们选择了Milvus,主要理由如下:

  • 标量过滤能力:餐厅知识库需要按门店ID、品类、地区等多维度过滤,Milvus原生支持标量字段过滤,省去了在应用层做二次过滤的开销
  • 水平扩展:Milvus的分布式架构让我们可以通过增加节点线性提升检索能力,为未来知识库增长预留了空间
  • 索引多样性:支持HNSW、IVF等多种索引类型,我们可以根据不同知识库的规模和延迟要求选择最优索引
  • 生态完善:丰富的SDK(Python、Go、Java),与LangChain等框架无缝集成

4.3 Embedding模型选型

Embedding模型决定了文本向量化的质量,直接影响检索准确率。我们对比了三个方案:

Embedding模型 维度 中文效果 推理速度 部署方式 适用场景
BGE-large-zh-v1.5 1024 优秀 快(GPU) 本地部署 通用中文检索
m3e-large 1024 良好 快(GPU) 本地部署 对话检索
text-embedding-3-large 3072 优秀 中(API) API调用 高质量检索

我们选择了BGE-large-zh-v1.5,理由是在中文餐饮领域的检索准确率测试中表现最优(准确率88.3% vs m3e的84.1%),同时本地部署保证了数据隐私和调用延迟。在实际使用中,单条文本的向量化延迟在GPU环境下约15ms,完全满足实时检索需求。

4.4 开发框架选型

在对话编排框架的选择上,我们对比了LangChain、LlamaIndex和自研框架三个方案:

框架 优势 劣势 适合团队 我们的选择
LangChain 生态丰富、组件多、社区大 抽象层重、版本迭代快、性能开销 快速原型验证 核心编排层
LlamaIndex RAG能力强、索引丰富 对话管理弱、生态小于LangChain 知识密集型应用 知识检索层
自研框架 完全可控、性能最优 开发周期长、需维护成本 有工程能力的团队 接入层和执行层

我们采用了混合策略:LangChain用于对话编排层的快速搭建(利用其Chain、Agent、Memory等组件),LlamaIndex用于知识检索层的索引管理,接入层和能力执行层用Go自研以保证性能。这种"框架混用"的策略在实践中效果很好——每一层用最合适的工具,而不是被一个框架绑定。

五、知识库构建:RAG的成败在于数据

5.1 知识库的整体规划

在AI领域有一句名言:"Garbage in, garbage out。"这句话放在RAG系统中尤为贴切。再强大的大模型,如果喂给它的是低质量、过时、混乱的知识,产出的回答也必然令人失望。所以我们在知识库构建上投入了大量精力。

"味之境"的知识库分为四大模块:

知识模块 数据来源 文档数量 更新频率 向量化策略
门店信息库 门店管理系统API 23个门店 × ~50字段 实时同步 结构化 → 半结构化文本
菜品知识库 菜单管理系统 ~800道菜品 每周更新 每道菜独立向量化
FAQ知识库 历史工单 + 人工编写 ~2000条FAQ 每日增量更新 Q&A对向量化
政策规则库 运营部门文档 ~150份文档 按需更新 分块向量化(512 tokens)

5.2 数据清洗与预处理

原始数据往往不能直接使用。我们设计了一套完整的数据清洗流程,以下是Python实现的核心代码:

import re
import json
import hashlib
from typing import Dict, List, Optional
from dataclasses import dataclass, field
from datetime import datetime


@dataclass
class StoreInfo:
    """门店信息数据模型"""
    store_id: str
    name: str
    address: str
    phone: str
    business_hours: str
    latitude: float
    longitude: float
    has_parking: bool
    has_private_room: bool
    max_seats: int
    tags: List[str] = field(default_factory=list)
    status: str = "active"  # active / suspended / closed
    updated_at: str = ""


@dataclass
class MenuItem:
    """菜品信息数据模型"""
    dish_id: str
    name: str
    category: str
    price: float
    description: str
    spicy_level: int  # 0-3
    allergens: List[str] = field(default_factory=list)
    is_vegetarian: bool = False
    calories: Optional[int] = None
    available: bool = True


class KnowledgePreprocessor:
    """知识库数据预处理器"""
    
    # 敏感信息正则模式
    PHONE_PATTERN = re.compile(r'1[3-9]\d{9}')
    ID_CARD_PATTERN = re.compile(r'\d{17}[\dXx]')
    EMAIL_PATTERN = re.compile(r'[a-zA-Z0-9._%+-]+@[a-zA-Z0-9.-]+\.[a-zA-Z]{2,}')
    
    def __init__(self):
        self.processed_count = 0
        self.skipped_count = 0
        self.error_count = 0
    
    def clean_text(self, text: str) -> str:
        """清洗文本:去除多余空白、特殊字符、HTML标签"""
        if not text:
            return ""
        # 去除HTML标签
        text = re.sub(r'<[^>]+>', '', text)
        # 去除多余空白
        text = re.sub(r'\s+', ' ', text).strip()
        # 去除不可见字符
        text = re.sub(r'[\x00-\x08\x0b\x0c\x0e-\x1f\x7f-\x9f]', '', text)
        return text
    
    def desensitize(self, text: str) -> str:
        """敏感信息脱敏"""
        text = self.PHONE_PATTERN.sub('[手机号]', text)
        text = self.ID_CARD_PATTERN.sub('[身份证号]', text)
        text = self.EMAIL_PATTERN.sub('[邮箱]', text)
        return text
    
    def store_to_text(self, store: StoreInfo) -> str:
        """将门店结构化数据转换为可向量化的文本"""
        parts = [
            f"门店名称:{store.name}",
            f"门店编号:{store.store_id}",
            f"地址:{store.address}",
            f"电话:{store.phone}",
            f"营业时间:{store.business_hours}",
            f"是否有停车场:{'是' if store.has_parking else '否'}",
            f"是否有包间:{'是' if store.has_private_room else '否'}",
            f"最大容纳人数:{store.max_seats}人",
            f"门店状态:{store.status}",
        ]
        if store.tags:
            parts.append(f"标签:{', '.join(store.tags)}")
        return '\n'.join(parts)
    
    def menu_to_text(self, dish: MenuItem) -> str:
        """将菜品结构化数据转换为可向量化的文本"""
        spicy_map = {0: "不辣", 1: "微辣", 2: "中辣", 3: "重辣"}
        parts = [
            f"菜品名称:{dish.name}",
            f"分类:{dish.category}",
            f"价格:{dish.price}元",
            f"描述:{dish.description}",
            f"辣度:{spicy_map.get(dish.spicy_level, '未知')}",
            f"是否素食:{'是' if dish.is_vegetarian else '否'}",
        ]
        if dish.allergens:
            parts.append(f"过敏原:{', '.join(dish.allergens)}")
        if dish.calories:
            parts.append(f"热量:{dish.calories}千卡")
        parts.append(f"是否可点:{'是' if dish.available else '否(已售罄或下架)'}")
        return '\n'.join(parts)
    
    def split_document(
        self, 
        text: str, 
        chunk_size: int = 512,
        chunk_overlap: int = 64
    ) -> List[str]:
        """
        文档分块:按句子边界切分,保持语义完整性
        """
        if not text:
            return []
        
        # 按句子分割(中英文标点)
        sentences = re.split(r'(?<=[。!?.!?\n])', text)
        sentences = [s.strip() for s in sentences if s.strip()]
        
        chunks = []
        current_chunk = ""
        
        for sentence in sentences:
            # 如果加上这句话不超过chunk_size,就加入当前块
            if len(current_chunk) + len(sentence) <= chunk_size:
                current_chunk += sentence
            else:
                # 当前块已满,保存并开始新块
                if current_chunk:
                    chunks.append(current_chunk.strip())
                
                # 保留overlap部分
                if chunks and chunk_overlap > 0:
                    overlap_text = chunks[-1][-chunk_overlap:]
                    current_chunk = overlap_text + sentence
                else:
                    current_chunk = sentence
        
        # 别忘了最后一块
        if current_chunk.strip():
            chunks.append(current_chunk.strip())
        
        return chunks
    
    def generate_doc_id(self, content: str, source: str = "") -> str:
        """生成文档唯一ID"""
        raw = f"{source}:{content[:200]}"
        return hashlib.md5(raw.encode('utf-8')).hexdigest()
    
    def process_store_batch(
        self, 
        stores: List[StoreInfo]
    ) -> List[Dict]:
        """批量处理门店数据,输出向量化所需的文档列表"""
        documents = []
        for store in stores:
            try:
                text = self.store_to_text(store)
                text = self.clean_text(text)
                
                doc = {
                    'id': self.generate_doc_id(text, f"store:{store.store_id}"),
                    'content': text,
                    'metadata': {
                        'type': 'store_info',
                        'store_id': store.store_id,
                        'store_name': store.name,
                        'city': store.tags[0] if store.tags else '',
                        'updated_at': store.updated_at or datetime.now().isoformat(),
                    }
                }
                documents.append(doc)
                self.processed_count += 1
            except Exception as e:
                self.error_count += 1
                print(f"处理门店 {store.store_id} 失败: {e}")
        
        return documents
    
    def process_menu_batch(
        self, 
        dishes: List[MenuItem]
    ) -> List[Dict]:
        """批量处理菜品数据"""
        documents = []
        for dish in dishes:
            try:
                text = self.menu_to_text(dish)
                text = self.clean_text(text)
                
                doc = {
                    'id': self.generate_doc_id(text, f"dish:{dish.dish_id}"),
                    'content': text,
                    'metadata': {
                        'type': 'menu_item',
                        'dish_id': dish.dish_id,
                        'category': dish.category,
                        'is_available': dish.available,
                    }
                }
                documents.append(doc)
                self.processed_count += 1
            except Exception as e:
                self.error_count += 1
                print(f"处理菜品 {dish.dish_id} 失败: {e}")
        
        return documents
    
    def get_stats(self) -> Dict:
        """获取处理统计信息"""
        return {
            'processed': self.processed_count,
            'skipped': self.skipped_count,
            'errors': self.error_count,
            'total': self.processed_count + self.skipped_count + self.error_count,
        }

这段代码的设计有几个值得注意的点:

第一,数据模型用dataclass定义,既保证了类型安全,又提供了清晰的文档作用。每个字段的含义和类型一目了然,方便团队协作。

第二,脱敏处理在前。任何数据在向量化之前都经过脱敏,确保用户手机号、身份证号等PII信息不会进入向量数据库。

第三,分块策略按句子边界切分。简单地按字符数切分会破坏语义完整性,导致检索时取到半句话。我们按中文句号、问号、感叹号等标点作为句子边界,再按chunk_size组装,保证每个块都是完整的语义单元。

第四,文档ID用内容哈希生成。这样相同内容不会重复入库,同时内容更新时ID会变化,可以触发增量更新。

5.3 向量化与入库

数据清洗完成后,下一步是向量化并写入Milvus。以下是完整的向量化入库流程:

import time
from typing import List, Dict, Optional
from pymilvus import (
    connections,
    Collection,
    CollectionSchema,
    FieldSchema,
    DataType,
    utility,
)
from sentence_transformers import SentenceTransformer
import torch


class VectorStoreManager:
    """向量数据库管理器"""
    
    def __init__(
        self,
        host: str = "localhost",
        port: str = "19530",
        collection_name: str = "restaurant_knowledge",
        embedding_model: str = "BAAI/bge-large-zh-v1.5",
        dim: int = 1024,
    ):
        self.host = host
        self.port = port
        self.collection_name = collection_name
        self.dim = dim
        
        # 加载Embedding模型
        device = "cuda" if torch.cuda.is_available() else "cpu"
        self.encoder = SentenceTransformer(embedding_model, device=device)
        print(f"Embedding模型已加载: {embedding_model} (device={device})")
        
        # 连接Milvus
        self._connect()
        
        # 创建或加载Collection
        self._init_collection()
    
    def _connect(self):
        """连接Milvus"""
        connections.connect(
            alias="default",
            host=self.host,
            port=self.port,
        )
        print(f"已连接Milvus: {self.host}:{self.port}")
    
    def _init_collection(self):
        """初始化Collection"""
        if utility.has_collection(self.collection_name):
            self.collection = Collection(self.collection_name)
            self.collection.load()
            print(f"已加载Collection: {self.collection_name}")
        else:
            self._create_collection()
    
    def _create_collection(self):
        """创建Collection"""
        fields = [
            FieldSchema(name="id", dtype=DataType.VARCHAR, max_length=64, is_primary=True),
            FieldSchema(name="embedding", dtype=DataType.FLOAT_VECTOR, dim=self.dim),
            FieldSchema(name="content", dtype=DataType.VARCHAR, max_length=4096),
            FieldSchema(name="doc_type", dtype=DataType.VARCHAR, max_length=32),
            FieldSchema(name="store_id", dtype=DataType.VARCHAR, max_length=32),
            FieldSchema(name="category", dtype=DataType.VARCHAR, max_length=64),
            FieldSchema(name="created_at", dtype=DataType.FLOAT),
        ]
        
        schema = CollectionSchema(
            fields=fields,
            description="餐厅知识库向量存储",
            enable_dynamic_field=False,
        )
        
        self.collection = Collection(
            name=self.collection_name,
            schema=schema,
            using="default",
        )
        
        # 创建HNSW索引(适合中小规模,召回率高)
        index_params = {
            "metric_type": "COSINE",
            "index_type": "HNSW",
            "params": {
                "M": 16,           # 每个节点的最大连接数
                "efConstruction": 200,  # 构建时的搜索宽度
            },
        }
        
        self.collection.create_index(
            field_name="embedding",
            index_params=index_params,
        )
        
        print(f"已创建Collection: {self.collection_name} (HNSW索引)")
    
    def encode_texts(self, texts: List[str]) -> List[List[float]]:
        """将文本批量转换为向量"""
        # 批量编码,提高GPU利用率
        embeddings = self.encoder.encode(
            texts,
            batch_size=32,
            normalize_embeddings=True,  # L2归一化,配合COSINE度量
            show_progress_bar=True,
        )
        return embeddings.tolist()
    
    def insert_documents(self, documents: List[Dict]) -> int:
        """批量插入文档"""
        if not documents:
            return 0
        
        # 提取文本进行向量化
        texts = [doc['content'] for doc in documents]
        embeddings = self.encode_texts(texts)
        
        # 准备插入数据
        data = [
            [doc['id'] for doc in documents],           # id
            embeddings,                                   # embedding
            [doc['content'][:4096] for doc in documents], # content
            [doc.get('metadata', {}).get('type', 'unknown') for doc in documents],  # doc_type
            [doc.get('metadata', {}).get('store_id', '') for doc in documents],     # store_id
            [doc.get('metadata', {}).get('category', '') for doc in documents],     # category
            [time.time() for _ in documents],            # created_at
        ]
        
        # 插入数据
        result = self.collection.insert(data)
        self.collection.flush()
        
        print(f"成功插入 {len(documents)} 条文档,当前总数: {self.collection.num_entities}")
        return len(documents)
    
    def search(
        self,
        query: str,
        top_k: int = 5,
        doc_type: Optional[str] = None,
        store_id: Optional[str] = None,
        category: Optional[str] = None,
    ) -> List[Dict]:
        """
        向量检索
        支持标量过滤:按文档类型、门店ID、分类等过滤
        """
        # 向量化查询
        query_embedding = self.encode_texts([query])[0]
        
        # 构建标量过滤表达式
        filter_expr = self._build_filter_expr(doc_type, store_id, category)
        
        # 搜索参数
        search_params = {
            "metric_type": "COSINE",
            "params": {"ef": 64},  # HNSW搜索时的搜索宽度
        }
        
        # 输出字段
        output_fields = ["id", "content", "doc_type", "store_id", "category"]
        
        results = self.collection.search(
            data=[query_embedding],
            anns_field="embedding",
            param=search_params,
            limit=top_k,
            expr=filter_expr,
            output_fields=output_fields,
        )
        
        # 解析结果
        hits = results[0]
        documents = []
        for hit in hits:
            doc = {
                'id': hit.entity.get('id'),
                'content': hit.entity.get('content'),
                'doc_type': hit.entity.get('doc_type'),
                'store_id': hit.entity.get('store_id'),
                'category': hit.entity.get('category'),
                'score': hit.score,
            }
            documents.append(doc)
        
        return documents
    
    def _build_filter_expr(
        self,
        doc_type: Optional[str] = None,
        store_id: Optional[str] = None,
        category: Optional[str] = None,
    ) -> str:
        """构建Milvus标量过滤表达式"""
        conditions = []
        if doc_type:
            conditions.append(f'doc_type == "{doc_type}"')
        if store_id:
            conditions.append(f'store_id == "{store_id}"')
        if category:
            conditions.append(f'category == "{category}"')
        
        return ' and '.join(conditions) if conditions else ""
    
    def delete_by_store_id(self, store_id: str) -> int:
        """按门店ID删除文档(用于门店信息更新时清除旧数据)"""
        expr = f'store_id == "{store_id}"'
        result = self.collection.delete(expr)
        self.collection.flush()
        print(f"已删除门店 {store_id} 的旧文档")
        return result.delete_count
    
    def get_stats(self) -> Dict:
        """获取知识库统计信息"""
        return {
            'collection_name': self.collection_name,
            'total_entities': self.collection.num_entities,
            'index_info': self.collection.indexes,
        }

这段代码有几个关键技术点值得说明:

第一,HNSW索引参数调优。M=16和efConstruction=200是我们在实测中找到的平衡点——M值越大,索引精度越高但内存消耗也越大;efConstruction越大,构建越慢但索引质量越好。对于百万级向量,这个参数组合的召回率达到了97.8%,查询延迟在5ms以内。

第二,标量过滤表达式。Milvus的标量过滤能力是我们在选型时的关键决策因素。当用户问"朝阳大悦城店的营业时间"时,我们可以先按store_id过滤出该门店的文档,再做向量检索,大幅提升了检索准确率和效率。

第三,L2归一化 + COSINE度量。在向量化时做L2归一化,然后用COSINE(余弦相似度)作为度量方式。这是文本检索领域的标准做法,能消除向量模长的影响,专注于方向相似性。

六、核心代码实现:从意图识别到回答生成

6.1 意图识别模块

意图识别是对话编排的入口。我们需要在用户消息进入后,快速判断其意图类型,以便选择合适的处理流程。我们设计了一个分层意图识别器:

血泪教训:不要试图用一个大模型解决所有问题。我们最初用GPT-4处理所有意图,结果简单的"今天营业吗"也要花费0.03美元的API调用费。后来拆分成三层:简单规则匹配处理60%的请求、小模型处理30%的复杂查询、大模型只处理10%的疑难问题,成本直接降了80%。

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


class IntentType(Enum):
    """意图类型枚举"""
    STORE_QUERY = "store_query"           # 门店信息查询
    MENU_QUERY = "menu_query"             # 菜品/菜单查询
    RESERVATION = "reservation"           # 预订
    QUEUE_QUERY = "queue_query"           # 排队查询
    ORDER_QUERY = "order_query"           # 订单查询
    REFUND_REQUEST = "refund_request"     # 退款
    PROMOTION_QUERY = "promotion_query"   # 促销活动
    COMPLAINT = "complaint"              # 投诉
    CHITCHAT = "chitchat"                # 闲聊
    UNKNOWN = "unknown"                  # 未知意图


@dataclass
class IntentResult:
    """意图识别结果"""
    intent: IntentType
    confidence: float
    entities: Dict[str, str]  # 提取的实体,如 store_name, dish_name 等
    raw_text: str


class IntentClassifier:
    """
    分层意图分类器
    第一层:基于规则的关键词匹配(快速路径,处理明确意图)
    第二层:基于大模型的意图理解(兜底路径,处理模糊意图)
    """
    
    # 第一层:关键词规则
    INTENT_RULES = {
        IntentType.STORE_QUERY: [
            r'营业时间', r'几点(开门|关门|营业)', r'地址|在哪',
            r'电话|联系方式', r'停车', r'包间|包房',
            r'门店|分店|店铺',
        ],
        IntentType.MENU_QUERY: [
            r'菜单|菜品', r'有什么(菜|吃的|推荐)',
            r'价格|多少钱', r'辣|辣度', r'过敏',
            r'素食|素菜', r'推荐(菜|什么)',
        ],
        IntentType.RESERVATION: [
            r'预订|预约|订位|订桌|订餐',
            r'约个位|约个桌',
        ],
        IntentType.QUEUE_QUERY: [
            r'排队|等位|叫号|排到几号',
            r'等多久|要等多长时间',
        ],
        IntentType.ORDER_QUERY: [
            r'订单|订单号', r'我的(订单|消费|账单)',
            r'查一下.*消费', r'小票',
        ],
        IntentType.REFUND_REQUEST: [
            r'退款|退钱|退一下', r'多收|乱收费',
            r'退(费|款).*怎么办',
        ],
        IntentType.PROMOTION_QUERY: [
            r'活动|优惠|满减|折扣',
            r'优惠券|代金券|红包',
            r'会员|积分|权益',
        ],
        IntentType.COMPLAINT: [
            r'投诉|举报|差评', r'态度(差|不好|恶劣)',
            r'难吃|不新鲜|变质|有异物',
            r'脏|卫生(差|有问题)',
        ],
        IntentType.CHITCHAT: [
            r'你好|您好|hi|hello',
            r'你是(谁|机器人|AI|人工智能)',
            r'谢谢|感谢',
        ],
    }
    
    # 门店名称提取正则(示例,实际应从门店库动态生成)
    STORE_NAME_PATTERNS = [
        r'(朝阳大悦城店|三里屯店|国贸店|太古里店|SKP店)',
        r'(北京|上海|成都).*?店',
    ]
    
    def __init__(self, llm_client=None):
        self.llm_client = llm_client
        self._compile_rules()
    
    def _compile_rules(self):
        """预编译正则表达式,提升匹配速度"""
        self.compiled_rules = {}
        for intent, patterns in self.INTENT_RULES.items():
            self.compiled_rules[intent] = [
                re.compile(p, re.IGNORECASE) for p in patterns
            ]
        
        self.compiled_store_patterns = [
            re.compile(p) for p in self.STORE_NAME_PATTERNS
        ]
    
    def classify(self, text: str) -> IntentResult:
        """
        意图识别主入口
        先走规则匹配,置信度不够时降级到大模型
        """
        text = text.strip()
        
        # 第一层:规则匹配
        result = self._rule_based_classify(text)
        
        if result and result.confidence >= 0.7:
            return result
        
        # 第二层:大模型理解
        if self.llm_client:
            llm_result = self._llm_based_classify(text)
            if llm_result:
                # 如果规则也有结果,取置信度更高的
                if result and result.confidence > llm_result.confidence:
                    return result
                return llm_result
        
        # 兜底:如果规则有弱匹配,返回弱匹配;否则返回UNKNOWN
        if result:
            return result
        
        return IntentResult(
            intent=IntentType.UNKNOWN,
            confidence=0.0,
            entities={},
            raw_text=text,
        )
    
    def _rule_based_classify(self, text: str) -> Optional[IntentResult]:
        """基于规则的意图识别"""
        scores = {}
        
        for intent, patterns in self.compiled_rules.items():
            score = 0
            for pattern in patterns:
                matches = pattern.findall(text)
                if matches:
                    score += len(matches)
            if score > 0:
                scores[intent] = score
        
        if not scores:
            return None
        
        # 取得分最高的意图
        best_intent = max(scores, key=scores.get)
        total_score = sum(scores.values())
        confidence = min(scores[best_intent] / total_score, 0.95)
        
        # 提取实体
        entities = self._extract_entities(text)
        
        return IntentResult(
            intent=best_intent,
            confidence=confidence,
            entities=entities,
            raw_text=text,
        )
    
    def _extract_entities(self, text: str) -> Dict[str, str]:
        """从文本中提取实体(门店名、菜品名等)"""
        entities = {}
        
        # 提取门店名
        for pattern in self.compiled_store_patterns:
            matches = pattern.findall(text)
            if matches:
                entities['store_name'] = matches[0] if isinstance(matches[0], str) else matches[0][0]
                break
        
        # 提取数字(可能是订单号、人数等)
        number_match = re.search(r'\d{6,}', text)
        if number_match:
            entities['number'] = number_match.group()
        
        # 提取时间
        time_match = re.search(
            r'(今天|明天|后天|周[一二三四五六日天])'
            r'.*?(\d{1,2}[点::]\d{0,2})',
            text
        )
        if time_match:
            entities['date'] = time_match.group(1)
            entities['time'] = time_match.group(2)
        
        return entities
    
    def _llm_based_classify(self, text: str) -> Optional[IntentResult]:
        """基于大模型的意图识别"""
        prompt = f"""请分析以下用户消息的意图,并返回JSON格式结果。

用户消息:"{text}"

可选意图类型:
- store_query: 门店信息查询(营业时间、地址、电话等)
- menu_query: 菜品/菜单查询(推荐菜品、价格、辣度等)
- reservation: 预订/订位
- queue_query: 排队/等位查询
- order_query: 订单查询
- refund_request: 退款请求
- promotion_query: 促销活动咨询
- complaint: 投诉或负面反馈
- chitchat: 闲聊或问候

请返回如下JSON格式:
{{"intent": "意图类型", "confidence": 0.0-1.0, "entities": {{}}}}

只返回JSON,不要其他内容。"""
        
        try:
            response = self.llm_client.chat(
                messages=[{"role": "user", "content": prompt}],
                temperature=0.1,  # 低温度保证稳定性
                max_tokens=200,
            )
            
            # 解析JSON响应
            import json
            result = json.loads(response.strip())
            
            intent_str = result.get('intent', 'unknown')
            try:
                intent = IntentType(intent_str)
            except ValueError:
                intent = IntentType.UNKNOWN
            
            return IntentResult(
                intent=intent,
                confidence=float(result.get('confidence', 0.5)),
                entities=result.get('entities', {}),
                raw_text=text,
            )
        except Exception as e:
            print(f"大模型意图识别失败: {e}")
            return None

这个分层意图识别器的设计哲学是"快路径优先,慢路径兜底"。大部分用户消息是明确的("你们几点营业"明显是门店查询),规则匹配可以毫秒级返回结果,不需要调用大模型。只有当规则匹配的置信度不够时,才降级到大模型理解。在我们的实测中,约72%的请求被规则层直接处理,平均耗时不到1ms,有效降低了大模型的调用成本。

6.2 对话编排引擎

对话编排引擎是智能体的"指挥中心",负责根据意图识别结果编排知识检索和工具调用的执行顺序。以下是核心实现:

from typing import List, Dict, Optional, Any
from dataclasses import dataclass, field
from enum import Enum
import time
import logging

logger = logging.getLogger(__name__)


class MessageRole(Enum):
    SYSTEM = "system"
    USER = "user"
    ASSISTANT = "assistant"
    TOOL = "tool"


@dataclass
class Message:
    role: MessageRole
    content: str
    tool_call_id: Optional[str] = None
    tool_name: Optional[str] = None
    timestamp: float = field(default_factory=time.time)


@dataclass 
class ConversationContext:
    """对话上下文"""
    session_id: str
    messages: List[Message] = field(default_factory=list)
    current_intent: Optional[str] = None
    entities: Dict[str, Any] = field(default_factory=dict)
    turn_count: int = 0
    awaiting_user: bool = False
    transferred_to_human: bool = False
    
    def add_message(self, message: Message):
        self.messages.append(message)
        # 保持上下文窗口在合理范围
        if len(self.messages) > 20:
            self.messages = self.messages[-20:]
    
    def get_recent_messages(self, n: int = 10) -> List[Message]:
        return self.messages[-n:]
    
    def to_llm_messages(self) -> List[Dict]:
        """转换为LLM可接受的消息格式"""
        result = []
        for msg in self.messages:
            item = {"role": msg.role.value, "content": msg.content}
            if msg.tool_name:
                item["name"] = msg.tool_name
            if msg.tool_call_id:
                item["tool_call_id"] = msg.tool_call_id
            result.append(item)
        return result


class ToolRegistry:
    """工具注册中心"""
    
    def __init__(self):
        self._tools: Dict[str, Dict] = {}
    
    def register(
        self,
        name: str,
        description: str,
        parameters: Dict,
        handler: callable,
    ):
        self._tools[name] = {
            'name': name,
            'description': description,
            'parameters': parameters,
            'handler': handler,
        }
        logger.info(f"已注册工具: {name}")
    
    def get_tool(self, name: str) -> Optional[Dict]:
        return self._tools.get(name)
    
    def get_all_schemas(self) -> List[Dict]:
        """获取所有工具的schema(用于Function Calling)"""
        return [
            {
                'type': 'function',
                'function': {
                    'name': t['name'],
                    'description': t['description'],
                    'parameters': t['parameters'],
                }
            }
            for t in self._tools.values()
        ]
    
    async def execute(self, name: str, arguments: Dict) -> str:
        """执行工具调用"""
        tool = self._tools.get(name)
        if not tool:
            return f"工具 {name} 不存在"
        
        try:
            result = await tool['handler'](**arguments)
            return str(result)
        except Exception as e:
            logger.error(f"工具 {name} 执行失败: {e}")
            return f"工具执行出错: {str(e)}"


class ConversationOrchestrator:
    """
    对话编排引擎
    负责编排意图识别 → 知识检索 → 工具调用 → 回答生成的完整流程
    """
    
    def __init__(
        self,
        intent_classifier,
        vector_store,
        tool_registry: ToolRegistry,
        llm_client,
        max_tool_calls: int = 3,
        max_retries: int = 2,
    ):
        self.intent_classifier = intent_classifier
        self.vector_store = vector_store
        self.tool_registry = tool_registry
        self.llm_client = llm_client
        self.max_tool_calls = max_tool_calls
        self.max_retries = max_retries
        
        # 系统提示词
        self.system_prompt = self._build_system_prompt()
    
    def _build_system_prompt(self) -> str:
        """构建系统提示词"""
        return """你是"味之境"连锁餐厅的智能客服助手。你的职责是:

1. 准确回答顾客关于门店信息、菜品、预订、排队、订单、促销活动等问题
2. 对于需要查询实时信息的问题(如排队状态、订单详情),主动调用相应工具
3. 对于退款、投诉等敏感问题,先安抚顾客情绪,再协助处理或转接人工客服
4. 保持友好、专业、简洁的回答风格
5. 如果不确定答案,诚实告知并建议联系人工客服

回答规则:
- 使用中文回复
- 回答控制在3句话以内,除非顾客要求详细说明
- 涉及价格、时间等信息时确保准确
- 不要编造不存在的菜品、门店或活动信息
- 如果顾客情绪激动或提出投诉,优先表示理解和歉意"""
    
    async def process(
        self,
        session_id: str,
        user_message: str,
        context: Optional[ConversationContext] = None,
    ) -> str:
        """
        处理用户消息的完整流程
        
        流程:
        1. 获取或创建对话上下文
        2. 意图识别
        3. 知识检索(RAG)
        4. 工具调用(如需要)
        5. 大模型生成回答
        6. 返回结果
        """
        # Step 1: 获取或创建上下文
        if context is None:
            context = ConversationContext(session_id=session_id)
        
        context.add_message(Message(
            role=MessageRole.USER,
            content=user_message,
        ))
        context.turn_count += 1
        
        # 检查是否已转人工
        if context.transferred_to_human:
            return "您的问题已转接人工客服,请稍候。"
        
        # Step 2: 意图识别
        intent_result = self.intent_classifier.classify(user_message)
        context.current_intent = intent_result.intent.value
        context.entities.update(intent_result.entities)
        
        logger.info(
            f"Session {session_id} | Turn {context.turn_count} | "
            f"Intent: {intent_result.intent.value} | "
            f"Confidence: {intent_result.confidence:.2f}"
        )
        
        # 检查是否需要转人工
        if self._should_transfer_to_human(intent_result, context):
            context.transferred_to_human = True
            response = (
                "非常抱歉给您带来不好的体验。"
                "我已为您转接人工客服,会有专人为您处理,请稍候。"
            )
            context.add_message(Message(
                role=MessageRole.ASSISTANT,
                content=response,
            ))
            return response
        
        # Step 3: 知识检索
        retrieved_docs = self._retrieve_knowledge(user_message, context)
        
        # Step 4 & 5: 构建提示并调用大模型(含工具调用循环)
        response = await self._generate_response(
            user_message, context, retrieved_docs
        )
        
        # Step 6: 记录并返回
        context.add_message(Message(
            role=MessageRole.ASSISTANT,
            content=response,
        ))
        
        return response
    
    def _should_transfer_to_human(
        self,
        intent_result,
        context: ConversationContext,
    ) -> bool:
        """判断是否需要转人工"""
        # 投诉意图直接转人工
        if intent_result.intent.value == 'complaint':
            return True
        
        # 退款请求且金额较大时转人工
        if intent_result.intent.value == 'refund_request':
            # 简单判断:如果对话超过3轮还在退款问题,转人工
            if context.turn_count > 3:
                return True
        
        return False
    
    def _retrieve_knowledge(
        self,
        query: str,
        context: ConversationContext,
    ) -> List[Dict]:
        """知识检索(RAG)"""
        try:
            # 根据意图和实体构建过滤条件
            doc_type = None
            store_id = context.entities.get('store_id')
            
            intent = context.current_intent
            if intent == 'store_query':
                doc_type = 'store_info'
            elif intent == 'menu_query':
                doc_type = 'menu_item'
            
            # 向量检索
            docs = self.vector_store.search(
                query=query,
                top_k=5,
                doc_type=doc_type,
                store_id=store_id,
            )
            
            # 过滤低相似度结果
            docs = [d for d in docs if d['score'] > 0.5]
            
            logger.info(f"知识检索: 查询='{query[:30]}...' 召回={len(docs)}条")
            return docs
            
        except Exception as e:
            logger.error(f"知识检索失败: {e}")
            return []
    
    async def _generate_response(
        self,
        user_message: str,
        context: ConversationContext,
        retrieved_docs: List[Dict],
    ) -> str:
        """生成回答(含工具调用循环)"""
        
        # 构建知识上下文
        knowledge_context = self._format_knowledge(retrieved_docs)
        
        # 构建消息列表
        messages = [
            {"role": "system", "content": self.system_prompt},
        ]
        
        if knowledge_context:
            messages.append({
                "role": "system",
                "content": f"以下是相关知识信息,请参考回答:\n\n{knowledge_context}",
            })
        
        # 加入历史对话(最近5轮)
        recent = context.get_recent_messages(10)
        for msg in recent:
            messages.append({"role": msg.role.value, "content": msg.content})
        
        # 工具调用循环
        tool_call_count = 0
        
        while tool_call_count < self.max_tool_calls:
            # 调用大模型
            tools = self.tool_registry.get_all_schemas()
            
            llm_response = self.llm_client.chat(
                messages=messages,
                tools=tools if tools else None,
                temperature=0.3,
                max_tokens=500,
            )
            
            # 检查是否有工具调用
            tool_calls = llm_response.get('tool_calls', [])
            
            if not tool_calls:
                # 没有工具调用,直接返回回答
                return llm_response['content']
            
            # 执行工具调用
            messages.append(llm_response['message'])
            
            for tc in tool_calls:
                tool_name = tc['function']['name']
                import json
                arguments = json.loads(tc['function']['arguments'])
                
                logger.info(f"工具调用: {tool_name}({arguments})")
                
                tool_result = await self.tool_registry.execute(
                    tool_name, arguments
                )
                
                messages.append({
                    "role": "tool",
                    "tool_call_id": tc['id'],
                    "name": tool_name,
                    "content": tool_result,
                })
                
                tool_call_count += 1
        
        # 达到最大工具调用次数后,再做一次最终生成
        final_response = self.llm_client.chat(
            messages=messages,
            temperature=0.3,
            max_tokens=500,
        )
        
        return final_response['content']
    
    def _format_knowledge(self, docs: List[Dict]) -> str:
        """格式化检索到的知识"""
        if not docs:
            return ""
        
        parts = []
        for i, doc in enumerate(docs, 1):
            parts.append(f"[知识{i}]\n{doc['content']}")
        
        return '\n\n'.join(parts)

这段代码体现了几个关键的工程设计理念:

工具调用循环。大模型可能需要多次调用工具才能回答一个问题。比如用户问"帮我看看朝阳大悦城店今天有什么推荐菜,顺便看看排队情况",大模型可能需要先调用search_menu查推荐菜,再调用check_queue_status查排队。我们设计了最多3轮的工具调用循环,每轮都将工具返回结果加入上下文,让大模型基于新信息决定是否需要继续调用工具。

上下文窗口管理。对话上下文限制在最近20条消息,防止token消耗过大。同时,在送给大模型时只取最近10条,保持上下文的相关性。

转人工机制。对于投诉类问题和长时间无法解决的退款问题,自动转人工。这是一个重要的安全阀——智能体的目标不是替代人工,而是在人工之前过滤掉大部分简单问题。

6.3 Go语言接入层实现

接入层用Go实现,主要处理高并发网络请求和消息路由。以下是核心实现:

package main

import (
	"context"
	"encoding/json"
	"fmt"
	"log"
	"net/http"
	"os"
	"os/signal"
	"syscall"
	"time"

	"github.com/gin-gonic/gin"
	"github.com/redis/go-redis/v9"
)

// ChatRequest 聊天请求
type ChatRequest struct {
	SessionID string `json:"session_id" binding:"required"`
	Message   string `json:"message" binding:"required"`
	Channel   string `json:"channel" binding:"required"` // wechat / app / web / eleme
	UserID    string `json:"user_id"`
}

// ChatResponse 聊天响应
type ChatResponse struct {
	SessionID  string `json:"session_id"`
	Reply      string `json:"reply"`
	Intent     string `json:"intent"`
	ToolsUsed  []string `json:"tools_used,omitempty"`
	Latency    float64 `json:"latency_ms"`
	Timestamp  int64   `json:"timestamp"`
}

// OrchestratorClient 对话编排层客户端
type OrchestratorClient struct {
	baseURL    string
	httpClient *http.Client
}

func NewOrchestratorClient(baseURL string) *OrchestratorClient {
	return &OrchestratorClient{
		baseURL: baseURL,
		httpClient: &http.Client{
			Timeout: 10 * time.Second,
		},
	}
}

// ProcessMessage 调用编排层处理消息
func (c *OrchestratorClient) ProcessMessage(
	ctx context.Context,
	req *ChatRequest,
) (*ChatResponse, error) {
	payload, err := json.Marshal(req)
	if err != nil {
		return nil, fmt.Errorf("marshal request: %w", err)
	}

	httpReq, err := http.NewRequestWithContext(
		ctx, "POST", c.baseURL+"/api/chat", bytes.NewReader(payload),
	)
	if err != nil {
		return nil, fmt.Errorf("create request: %w", err)
	}
	httpReq.Header.Set("Content-Type", "application/json")

	resp, err := c.httpClient.Do(httpReq)
	if err != nil {
		return nil, fmt.Errorf("orchestrator request: %w", err)
	}
	defer resp.Body.Close()

	if resp.StatusCode != http.StatusOK {
		return nil, fmt.Errorf("orchestrator error: status %d", resp.StatusCode)
	}

	var result ChatResponse
	if err := json.NewDecoder(resp.Body).Decode(&result); err != nil {
		return nil, fmt.Errorf("decode response: %w", err)
	}

	return &result, nil
}

// CacheManager 缓存管理器
type CacheManager struct {
	rdb *redis.Client
}

func NewCacheManager(addr string) *CacheManager {
	rdb := redis.NewClient(&redis.Options{
		Addr:         addr,
		Password:     os.Getenv("REDIS_PASSWORD"),
		DB:           0,
		PoolSize:     50,
		MinIdleConns: 10,
	})
	return &CacheManager{rdb: rdb}
}

// GetCachedReply 获取缓存的回复(用于高频FAQ)
func (cm *CacheManager) GetCachedReply(
	ctx context.Context,
	messageHash string,
) (string, error) {
	val, err := cm.rdb.Get(ctx, "faq:"+messageHash).Result()
	if err == redis.Nil {
		return "", nil
	}
	return val, err
}

// SetCachedReply 设置缓存
func (cm *CacheManager) SetCachedReply(
	ctx context.Context,
	messageHash string,
	reply string,
	ttl time.Duration,
) error {
	return cm.rdb.Set(ctx, "faq:"+messageHash, reply, ttl).Err()
}

// RateLimiter 限流器
type RateLimiter struct {
	rdb *redis.Client
}

func NewRateLimiter(addr string) *RateLimiter {
	rdb := redis.NewClient(&redis.Options{
		Addr:     addr,
		Password: os.Getenv("REDIS_PASSWORD"),
		DB:       1,
	})
	return &RateLimiter{rdb: rdb}
}

// Allow 检查是否允许请求(滑动窗口限流)
func (rl *RateLimiter) Allow(
	ctx context.Context,
	key string,
	limit int,
	window time.Duration,
) bool {
	now := time.Now().Unix()
	windowStart := now - int64(window.Seconds())

	pipe := rl.rdb.Pipeline()
	// 移除窗口外的记录
	pipe.ZRemRangeByScore(ctx, "rate:"+key, "0", fmt.Sprintf("%d", windowStart))
	// 添加当前请求
	pipe.ZAdd(ctx, "rate:"+key, redis.Z{
		Score:  float64(now),
		Member: fmt.Sprintf("%d", now),
	})
	// 统计窗口内请求数
	countCmd := pipe.ZCard(ctx, "rate:"+key)
	// 设置过期时间
	pipe.Expire(ctx, "rate:"+key, window)

	_, err := pipe.Exec(ctx)
	if err != nil {
		log.Printf("rate limiter error: %v", err)
		return true // 限流失败时放行,避免影响正常用户
	}

	return countCmd.Val() <= int64(limit)
}

// Server HTTP服务器
type Server struct {
	router       *gin.Engine
	orchestrator *OrchestratorClient
	cache        *CacheManager
	limiter      *RateLimiter
}

func NewServer() *Server {
	orchAddr := getEnv("ORCHESTRATOR_URL", "http://localhost:8000")
	redisAddr := getEnv("REDIS_ADDR", "localhost:6379")

	return &Server{
		router:       gin.Default(),
		orchestrator: NewOrchestratorClient(orchAddr),
		cache:        NewCacheManager(redisAddr),
		limiter:      NewRateLimiter(redisAddr),
	}
}

func (s *Server) SetupRoutes() {
	// 健康检查
	s.router.GET("/health", func(c *gin.Context) {
		c.JSON(http.StatusOK, gin.H{"status": "ok"})
	})

	// 聊天接口
	api := s.router.Group("/api")
	{
		api.POST("/chat", s.handleChat)
		api.GET("/sessions/:id/history", s.handleGetHistory)
	}
}

func (s *Server) handleChat(c *gin.Context) {
	var req ChatRequest
	if err := c.ShouldBindJSON(&req); err != nil {
		c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
		return
	}

	ctx := c.Request.Context()
	startTime := time.Now()

	// 限流:每用户每分钟30条
	rateKey := "user:" + req.UserID
	if !s.limiter.Allow(ctx, rateKey, 30, time.Minute) {
		c.JSON(http.StatusTooManyRequests, gin.H{
			"error": "请求过于频繁,请稍后再试",
		})
		return
	}

	// 检查FAQ缓存
	msgHash := hashMessage(req.Message, req.Channel)
	if cached, _ := s.cache.GetCachedReply(ctx, msgHash); cached != "" {
		c.JSON(http.StatusOK, ChatResponse{
			SessionID: req.SessionID,
			Reply:     cached,
			Latency:   float64(time.Since(startTime).Milliseconds()),
			Timestamp: time.Now().Unix(),
		})
		return
	}

	// 调用编排层
	resp, err := s.orchestrator.ProcessMessage(ctx, &req)
	if err != nil {
		log.Printf("orchestrator error: %v", err)
		// 降级:返回兜底回复
		c.JSON(http.StatusOK, ChatResponse{
			SessionID: req.SessionID,
			Reply:     "抱歉,系统暂时繁忙,请稍后重试或联系人工客服。",
			Latency:   float64(time.Since(startTime).Milliseconds()),
			Timestamp: time.Now().Unix(),
		})
		return
	}

	resp.Latency = float64(time.Since(startTime).Milliseconds())

	// 缓存高频FAQ回复(仅缓存简单FAQ)
	if resp.Intent == "store_query" || resp.Intent == "menu_query" {
		_ = s.cache.SetCachedReply(ctx, msgHash, resp.Reply, 5*time.Minute)
	}

	c.JSON(http.StatusOK, resp)
}

func (s *Server) handleGetHistory(c *gin.Context) {
	sessionID := c.Param("id")
	// TODO: 从存储中获取会话历史
	c.JSON(http.StatusOK, gin.H{
		"session_id": sessionID,
		"messages":   []interface{}{},
	})
}

func (s *Server) Run(addr string) error {
	srv := &http.Server{
		Addr:    addr,
		Handler: s.router,
	}

	go func() {
		log.Printf("Server starting on %s", addr)
		if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
			log.Fatalf("server error: %v", err)
		}
	}()

	// 优雅关闭
	quit := make(chan os.Signal, 1)
	signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM)
	<-quit
	log.Println("Shutting down server...")

	ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
	defer cancel()
	if err := srv.Shutdown(ctx); err != nil {
		log.Fatalf("server forced shutdown: %v", err)
	}

	log.Println("Server exited")
	return nil
}

func hashMessage(message, channel string) string {
	// 简单哈希,实际可用MD5或SHA256
	return fmt.Sprintf("%s:%s", channel, message)
}

func getEnv(key, defaultValue string) string {
	if val := os.Getenv(key); val != "" {
		return val
	}
	return defaultValue
}

func main() {
	gin.SetMode(gin.ReleaseMode)

	server := NewServer()
	server.SetupRoutes()

	addr := getEnv("SERVER_ADDR", ":8080")
	if err := server.Run(addr); err != nil {
		log.Fatalf("failed to run server: %v", err)
	}
}

Go接入层的设计重点在于高并发、限流、缓存和降级四个方面:

限流策略。使用Redis的有序集合实现滑动窗口限流,每用户每分钟限制30条消息。这个阈值是我们在压测中找到的平衡点——既能防止恶意刷接口,又不会影响正常用户的连续提问。

FAQ缓存。对于高频的门店信息查询和菜单查询,将回答缓存5分钟。这个策略让我们的缓存命中率达到了35%,大幅降低了对编排层和大模型的压力。缓存key使用渠道+消息内容的哈希,保证不同渠道的同一条消息可以共享缓存。

降级机制。当编排层调用失败时,接入层不会直接返回500错误,而是返回一个友好的降级提示。同时,限流器在Redis不可用时也会放行请求——宁可暂时失去限流能力,也不能因为限流器故障而拒绝正常用户。

优雅关闭。通过捕获SIGINT和SIGINT信号,在收到关闭指令后等待10秒让正在处理的请求完成,避免 abrupt shutdown 导致用户请求中断。

6.4 工具能力注册

智能体的工具能力需要注册到ToolRegistry中才能被大模型调用。以下是几个核心工具的实现:

import aiohttp
from typing import Dict, Any
import logging

logger = logging.getLogger(__name__)


class ToolImplementations:
    """工具能力实现"""
    
    def __init__(self, api_base_url: str, api_key: str):
        self.api_base_url = api_base_url
        self.api_key = api_key
        self.session = None
    
    async def get_session(self) -> aiohttp.ClientSession:
        if self.session is None or self.session.closed:
            self.session = aiohttp.ClientSession(
                headers={"Authorization": f"Bearer {self.api_key}"},
                timeout=aiohttp.ClientTimeout(total=5),
            )
        return self.session
    
    async def query_store_info(self, store_id: str = "", store_name: str = "") -> str:
        """
        查询门店信息
        可按store_id或store_name查询
        """
        session = await self.get_session()
        params = {}
        if store_id:
            params['store_id'] = store_id
        if store_name:
            params['name'] = store_name
        
        try:
            async with session.get(
                f"{self.api_base_url}/stores",
                params=params,
            ) as resp:
                if resp.status == 200:
                    data = await resp.json()
                    if data.get('stores'):
                        store = data['stores'][0]
                        return (
                            f"门店:{store['name']}\n"
                            f"地址:{store['address']}\n"
                            f"电话:{store['phone']}\n"
                            f"营业时间:{store['business_hours']}\n"
                            f"停车场:{'有' if store.get('has_parking') else '无'}\n"
                            f"包间:{'有' if store.get('has_private_room') else '无'}\n"
                            f"状态:{store.get('status', '正常营业')}"
                        )
                    return "未找到匹配的门店信息"
                return f"查询失败,状态码:{resp.status}"
        except Exception as e:
            logger.error(f"query_store_info error: {e}")
            return "门店信息查询暂时不可用,请稍后重试"
    
    async def check_queue_status(self, store_id: str) -> str:
        """
        查询门店排队状态
        """
        session = await self.get_session()
        try:
            async with session.get(
                f"{self.api_base_url}/stores/{store_id}/queue",
            ) as resp:
                if resp.status == 200:
                    data = await resp.json()
                    waiting = data.get('waiting_count', 0)
                    est_time = data.get('estimated_wait_minutes', 0)
                    return (
                        f"当前排队{waiting}桌,"
                        f"预计等待约{est_time}分钟。"
                        f"您可以在小程序上提前取号哦~"
                    )
                return "排队信息查询暂时不可用"
        except Exception as e:
            logger.error(f"check_queue_status error: {e}")
            return "排队信息查询暂时不可用,请稍后重试"
    
    async def search_menu(
        self,
        keyword: str = "",
        category: str = "",
        is_vegetarian: bool = False,
        max_spicy: int = 3,
    ) -> str:
        """
        搜索菜品
        可按关键词、分类、素食、辣度筛选
        """
        session = await self.get_session()
        params = {
            'keyword': keyword,
            'category': category,
            'vegetarian': str(is_vegetarian).lower(),
            'max_spicy': str(max_spicy),
            'limit': '5',
        }
        params = {k: v for k, v in params.items() if v}
        
        try:
            async with session.get(
                f"{self.api_base_url}/menu/search",
                params=params,
            ) as resp:
                if resp.status == 200:
                    data = await resp.json()
                    dishes = data.get('dishes', [])
                    if not dishes:
                        return "未找到匹配的菜品"
                    
                    lines = [f"为您找到{len(dishes)}道相关菜品:"]
                    spicy_map = {0: "不辣", 1: "微辣", 2: "中辣", 3: "重辣"}
                    for d in dishes:
                        spicy = spicy_map.get(d.get('spicy_level', 0), "未知")
                        lines.append(
                            f"- {d['name']} | {d['category']} | "
                            f"{d['price']}元 | {spicy}"
                        )
                    return '\n'.join(lines)
                return "菜品搜索暂时不可用"
        except Exception as e:
            logger.error(f"search_menu error: {e}")
            return "菜品搜索暂时不可用,请稍后重试"
    
    async def create_reservation(
        self,
        store_id: str,
        date: str,
        time: str,
        party_size: int,
        customer_name: str = "",
        customer_phone: str = "",
    ) -> str:
        """
        创建预订
        """
        session = await self.get_session()
        payload = {
            'store_id': store_id,
            'date': date,
            'time': time,
            'party_size': party_size,
            'customer_name': customer_name,
            'customer_phone': customer_phone,
        }
        
        try:
            async with session.post(
                f"{self.api_base_url}/reservations",
                json=payload,
            ) as resp:
                if resp.status == 201:
                    data = await resp.json()
                    return (
                        f"预订成功!\n"
                        f"预订编号:{data.get('reservation_id', '')}\n"
                        f"门店:{data.get('store_name', store_id)}\n"
                        f"时间:{date} {time}\n"
                        f"人数:{party_size}人\n"
                        f"请准时到店,如有变动请提前取消。"
                    )
                elif resp.status == 409:
                    return "该时段已满,请选择其他时间"
                return f"预订失败,状态码:{resp.status}"
        except Exception as e:
            logger.error(f"create_reservation error: {e}")
            return "预订服务暂时不可用,请稍后重试"
    
    async def query_order_detail(self, order_id: str = "", phone: str = "") -> str:
        """
        查询订单详情
        """
        session = await self.get_session()
        params = {}
        if order_id:
            params['order_id'] = order_id
        if phone:
            params['phone'] = phone
        
        try:
            async with session.get(
                f"{self.api_base_url}/orders",
                params=params,
            ) as resp:
                if resp.status == 200:
                    data = await resp.json()
                    orders = data.get('orders', [])
                    if not orders:
                        return "未找到匹配的订单"
                    
                    order = orders[0]
                    return (
                        f"订单编号:{order['order_id']}\n"
                        f"门店:{order.get('store_name', '')}\n"
                        f"下单时间:{order.get('created_at', '')}\n"
                        f"总金额:{order.get('total_amount', '')}元\n"
                        f"状态:{order.get('status', '')}"
                    )
                return "订单查询暂时不可用"
        except Exception as e:
            logger.error(f"query_order_detail error: {e}")
            return "订单查询暂时不可用,请稍后重试"
    
    async def close(self):
        if self.session and not self.session.closed:
            await self.session.close()


def register_all_tools(registry, tool_impl: ToolImplementations):
    """注册所有工具到ToolRegistry"""
    
    registry.register(
        name="query_store_info",
        description="查询门店信息,包括地址、电话、营业时间、停车场、包间等。可按门店ID或门店名称查询。",
        parameters={
            "type": "object",
            "properties": {
                "store_id": {
                    "type": "string",
                    "description": "门店ID,如不知道可留空",
                },
                "store_name": {
                    "type": "string",
                    "description": "门店名称,如'朝阳大悦城店'",
                },
            },
        },
        handler=tool_impl.query_store_info,
    )
    
    registry.register(
        name="check_queue_status",
        description="查询指定门店的当前排队等待状态,包括排队桌数和预计等待时间。",
        parameters={
            "type": "object",
            "properties": {
                "store_id": {
                    "type": "string",
                    "description": "门店ID",
                },
            },
            "required": ["store_id"],
        },
        handler=tool_impl.check_queue_status,
    )
    
    registry.register(
        name="search_menu",
        description="搜索餐厅菜品。可按关键词、分类、素食、辣度等条件筛选。",
        parameters={
            "type": "object",
            "properties": {
                "keyword": {
                    "type": "string",
                    "description": "菜品关键词,如'牛肉'、'汤'",
                },
                "category": {
                    "type": "string",
                    "description": "菜品分类,如'热菜'、'凉菜'、'主食'、'汤品'",
                },
                "is_vegetarian": {
                    "type": "boolean",
                    "description": "是否只搜索素食菜品",
                },
                "max_spicy": {
                    "type": "integer",
                    "description": "最大辣度(0=不辣, 1=微辣, 2=中辣, 3=重辣)",
                },
            },
        },
        handler=tool_impl.search_menu,
    )
    
    registry.register(
        name="create_reservation",
        description="为顾客创建餐厅预订。",
        parameters={
            "type": "object",
            "properties": {
                "store_id": {"type": "string", "description": "门店ID"},
                "date": {"type": "string", "description": "预订日期,格式YYYY-MM-DD"},
                "time": {"type": "string", "description": "预订时间,格式HH:MM"},
                "party_size": {"type": "integer", "description": "用餐人数"},
                "customer_name": {"type": "string", "description": "顾客姓名"},
                "customer_phone": {"type": "string", "description": "顾客手机号"},
            },
            "required": ["store_id", "date", "time", "party_size"],
        },
        handler=tool_impl.create_reservation,
    )
    
    registry.register(
        name="query_order_detail",
        description="查询顾客的订单详情,可按订单号或手机号查询。",
        parameters={
            "type": "object",
            "properties": {
                "order_id": {"type": "string", "description": "订单编号"},
                "phone": {"type": "string", "description": "顾客手机号"},
            },
        },
        handler=tool_impl.query_order_detail,
    )
    
    logger.info(f"已注册 {5} 个工具能力")

七、测试与优化:从60%准确率到91%的进化之路

7.1 测试体系搭建

智能体的测试和传统软件测试有很大不同。传统软件的输入输出是确定性的,而智能体的输出具有一定的不确定性——同一个问题,大模型可能生成不同但都正确的回答。因此,我们需要建立一套专门针对智能体的测试体系。

我们设计了三层测试体系:

测试层级 测试目标 测试方法 测试数据 执行频率
单元测试 各模块功能正确性 输入输出断言 手工编写用例 每次提交
集成测试 模块间协作正确性 端到端流程验证 场景化用例 每日构建
评估测试 回答质量和准确率 LLM自动评分 + 人工抽检 历史工单标注集 每周评估

评估测试是最关键的一环。我们从历史工单中标注了2000条问答对作为评估集,每条包含用户问题、标准答案、期望意图和期望工具调用。评估指标包括:

  • 意图准确率:意图识别结果与标注一致的比例
  • 工具调用准确率:调用了正确的工具且参数正确的比例
  • 回答准确率:回答内容包含正确信息的比例(通过LLM自动评分 + 人工抽检)
  • 回答流畅度:回答是否自然、通顺、符合客服场景(LLM评分)
  • 拒答率:本应回答但拒绝回答的比例(越低越好)
  • 幻觉率:回答中包含编造信息的比例(越低越好)

7.2 优化历程

系统的第一版评估结果并不理想——回答准确率只有60%,幻觉率高达12%。我们通过以下几轮优化逐步将准确率提升到91%:

第一轮:知识库质量优化

分析错误案例后发现,约40%的错误回答源于知识检索失败。原因是原始知识文档的文本质量差——有些门店信息包含HTML标签,有些菜品描述过于简短导致向量区分度不够。

优化措施:

  • 对所有知识文档执行数据清洗流程(见5.2节)
  • 为每道菜品补充更丰富的描述(烹饪方式、适合人群、搭配建议等)
  • 为FAQ增加同义问法扩展,一个问题对应3-5种不同的提问方式
  • 将检索top_k从3调整为5,并增加相似度阈值过滤

优化后,知识检索召回率从72%提升到89%,回答准确率提升到75%。

第二轮:提示词工程优化

第二轮优化聚焦于系统提示词。原始提示词过于简单,导致大模型在以下场景表现不佳:

  • 多意图问题("查查营业时间顺便推荐个菜")只回答了一个
  • 涉及具体数字的问题(价格、等待时间)容易编造
  • 投诉场景的回答过于机械,缺乏共情

优化后的提示词增加了明确的回答规则、场景处理指南和输出格式约束。关键改进包括:

# 优化后的系统提示词(关键片段)

## 回答规则
- 如果顾客的问题包含多个子问题,逐一回答,用序号标注
- 涉及价格、时间、数量等具体数字时,必须基于检索到的知识或工具返回结果回答,不得自行编造
- 如果知识库和工具都无法提供答案,明确告知"暂未查到相关信息"并建议联系人工客服

## 场景处理
- 投诉场景:先表示理解和歉意("非常抱歉给您带来不好的体验"),再提供解决方案或转人工
- 预订场景:确认关键信息(门店、时间、人数),如信息不全则追问
- 退款场景:了解退款原因,查询订单,如满足条件则发起退款流程,否则解释原因并转人工

## 输出格式
- 回答控制在3句话以内
- 使用友好的语气,适当使用"~"等符号增加亲和力
- 不要输出"根据知识库..."等内部信息

提示词优化后,多意图问题的处理准确率从55%提升到82%,幻觉率从12%降低到5%。

第三轮:工具调用优化

第三轮优化针对工具调用场景。我们发现大模型有时会"忘记"调用工具,直接用知识库中的过时信息回答实时问题。优化措施包括:

  • 在工具描述中增加更明确的触发条件说明
  • 在系统提示词中强调"涉及实时信息时必须调用工具"
  • 增加工具调用的后验证逻辑——如果回答中包含时间/排队/订单等关键词但未调用相应工具,触发二次生成

优化后,工具调用准确率从78%提升到95%,最终回答准确率稳定在91%。

7.3 压力测试与容量规划

在上线之前,我们进行了全面的压力测试,确保系统能够承受真实业务场景的流量冲击。压力测试不仅是"能不能扛住"的问题,更是"在什么压力下会先出问题"的预判——你需要知道系统的瓶颈在哪里,才能提前做好扩容预案。

我们使用Locust作为压测工具,模拟了三种典型场景:

测试场景 模拟条件 目标指标 实际结果
日常负载 50 QPS持续30分钟 P95 < 3s,错误率 < 0.1% P95=2.1s,错误率=0%
午市高峰 200 QPS持续10分钟 P95 < 5s,错误率 < 1% P95=3.8s,错误率=0.2%
突发流量 0→500 QPS in 30s 系统不崩溃,自动扩容 扩容耗时45s,期间P95=8s

压测结果暴露了一个问题:在突发流量场景下,HPA自动扩容需要45秒才能拉起新Pod,这期间已有Pod的延迟飙升至8秒。解决方案是引入"预热扩容"——根据历史流量规律,在预计高峰来临前10分钟主动扩容,而不是等到CPU阈值触发后才被动扩容。我们在K8s中通过CronJob实现了这个预热线程:每天10:50和16:50自动将编排层副本数从4个扩展到8个,13:30和21:30再缩回4个。这个简单的策略让高峰期的P95延迟稳定在3秒以内。

容量规划方面,我们根据压测数据推算了不同业务量级下的资源需求:

日均对话量 编排层实例 Embedding实例 Milvus节点 月度云成本
5万次 4 × 4C8G 1 × T4 1 × 8C16G ¥18,000
10万次 6 × 4C8G 2 × T4 3 × 8C16G ¥32,000
30万次 12 × 4C8G 4 × T4 3 × 16C32G ¥85,000
100万次 20 × 8C16G 8 × A10 5 × 16C32G ¥220,000

这个容量规划表让张磊在做年度预算时有了清晰的数据支撑。他可以根据品牌扩张计划(预计两年内门店数翻倍)提前预估IT成本增长,而不是等到系统扛不住了才紧急扩容——后者往往意味着更高的成本和更差的用户体验。

7.4 性能优化

除了准确率优化,性能优化同样重要。以下是几项关键的性能优化措施:

优化项 优化前 优化后 优化手段
Embedding向量化延迟 45ms/条 15ms/条 批量编码 + GPU加速
向量检索延迟 25ms 5ms HNSW索引参数调优
大模型首token延迟 1.2s 0.6s 流式输出 + 请求预热
端到端P95延迟 5.5s 2.8s 并行检索 + 缓存 + 流式
并发处理能力 50 QPS 300 QPS Go接入层 + 连接池 + 异步IO

其中,"并行检索"是一个效果显著的优化。原始流程是串行的——先知识检索,再工具调用,最后大模型生成。优化后,对于确定的意图,知识检索和工具调用可以并行执行,然后合并结果送入大模型。这个优化将端到端延迟降低了约40%。

八、部署上线:从开发环境到生产环境的跨越

8.1 部署架构

生产环境的部署架构如下:

组件 部署方式 实例数 资源配置 高可用策略
Go接入层 Docker + K8s 4 2C4G/实例 多副本 + 负载均衡
Python编排层 Docker + K8s 6 4C8G/实例 多副本 + 自动扩缩容
Milvus向量库 Docker Compose 3(集群) 8C16G/实例 主从 + 副本
Redis缓存 Docker 2(主从) 2C4G/实例 主从 + Sentinel
Embedding服务 Docker + GPU 2 4C8G + T4/实例 多副本 + 负载均衡
大模型API 第三方托管 - - 多供应商 + 自动切换
监控告警 Prometheus + Grafana 1 2C4G -

8.2 Docker化部署

以下是编排层的Dockerfile:

# 基础镜像:Python 3.11 + CUDA支持
FROM python:3.11-slim as builder

# 设置工作目录
WORKDIR /app

# 安装系统依赖
RUN apt-get update && apt-get install -y \
    build-essential \
    curl \
    && rm -rf /var/lib/apt/lists/*

# 安装Python依赖
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt

# 复制应用代码
COPY . .

# 生产环境镜像
FROM python:3.11-slim

WORKDIR /app

# 安装运行时依赖
RUN apt-get update && apt-get install -y \
    libgomp1 \
    && rm -rf /var/lib/apt/lists/*

# 从builder阶段复制依赖和应用
COPY --from=builder /usr/local/lib/python3.11/site-packages /usr/local/lib/python3.11/site-packages
COPY --from=builder /app /app

# 设置环境变量
ENV PYTHONPATH=/app
ENV PYTHONUNBUFFERED=1
ENV PYTHONDONTWRITEBYTECODE=1
ENV MODEL_CACHE_DIR=/app/model_cache

# 创建非root用户
RUN useradd -m -u 1000 appuser && chown -R appuser:appuser /app
USER appuser

# 健康检查
HEALTHCHECK --interval=30s --timeout=5s --retries=3 \
    CMD curl -f http://localhost:8000/health || exit 1

# 暴露端口
EXPOSE 8000

# 启动命令
CMD ["gunicorn", "main:app", \
     "--bind", "0.0.0.0:8000", \
     "--workers", "4", \
     "--worker-class", "uvicorn.workers.UvicornWorker", \
     "--timeout", "30", \
     "--keep-alive", "5", \
     "--max-requests", "1000", \
     "--max-requests-jitter", "100"]

Dockerfile的几个设计要点:

多阶段构建。builder阶段安装编译依赖并编译Python包,生产镜像只包含运行时必需的文件。这将镜像大小从2.3GB缩减到680MB,加快了部署速度。

非root用户运行。安全最佳实践,避免容器内进程以root权限运行。

健康检查。K8s根据健康检查结果决定是否将流量路由到该Pod。

Gunicorn + Uvicorn Worker。Gunicorn作为进程管理器,每个worker运行一个Uvicorn实例处理异步请求。workers=4是在4核机器上经过测试的最优配置。

8.3 Kubernetes部署配置

以下是编排层的K8s部署配置:

apiVersion: apps/v1
kind: Deployment
metadata:
  name: orchestrator
  namespace: restaurant-ai
  labels:
    app: orchestrator
spec:
  replicas: 6
  selector:
    matchLabels:
      app: orchestrator
  strategy:
    type: RollingUpdate
    rollingUpdate:
      maxSurge: 2
      maxUnavailable: 1
  template:
    metadata:
      labels:
        app: orchestrator
    spec:
      containers:
      - name: orchestrator
        image: registry.cn-beijing.aliyuncs.com/weizhijing/orchestrator:v2.1.0
        ports:
        - containerPort: 8000
          name: http
        env:
        - name: LLM_API_KEY
          valueFrom:
            secretKeyRef:
              name: llm-secrets
              key: api-key
        - name: MILVUS_HOST
          value: "milvus.cluster.local"
        - name: REDIS_ADDR
          value: "redis-master.cache.svc:6379"
        - name: LOG_LEVEL
          value: "INFO"
        resources:
          requests:
            cpu: "2000m"
            memory: "4Gi"
          limits:
            cpu: "4000m"
            memory: "8Gi"
        livenessProbe:
          httpGet:
            path: /health
            port: 8000
          initialDelaySeconds: 30
          periodSeconds: 15
        readinessProbe:
          httpGet:
            path: /ready
            port: 8000
          initialDelaySeconds: 10
          periodSeconds: 10
        lifecycle:
          preStop:
            exec:
              command: ["/bin/sh", "-c", "sleep 15"]
      terminationGracePeriodSeconds: 60

---
apiVersion: v1
kind: HorizontalPodAutoscaler
metadata:
  name: orchestrator-hpa
  namespace: restaurant-ai
spec:
  scaleTargetRef:
    apiVersion: apps/v1
    kind: Deployment
    name: orchestrator
  minReplicas: 4
  maxReplicas: 12
  metrics:
  - type: Resource
    resource:
      name: cpu
      target:
        type: Utilization
        averageUtilization: 70
  - type: Resource
    resource:
      name: memory
      target:
        type: Utilization
        averageUtilization: 80

---
apiVersion: v1
kind: Service
metadata:
  name: orchestrator
  namespace: restaurant-ai
spec:
  selector:
    app: orchestrator
  ports:
  - port: 80
    targetPort: 8000
    name: http
  type: ClusterIP

K8s配置的几个关键点:

滚动更新策略。maxSurge=2和maxUnavailable=1的配置确保更新过程中始终有足够实例提供服务,不会出现服务中断。

HPA自动扩缩容。根据CPU和内存使用率自动调整Pod数量,在午市高峰期自动扩展到更多实例,在低谷期自动缩减以节省成本。minReplicas=4保证了基础可用性。

优雅终止。preStop钩子让Pod在收到终止信号后等待15秒,配合terminationGracePeriodSeconds=60,确保正在处理的请求能够完成。

就绪探针与存活探针分离。就绪探针检查/ready端点(包括依赖检查),存活探针检查/health端点(仅检查进程存活)。这样在依赖服务暂时不可用时,Pod不会被重启,只是从负载均衡中暂时摘除。

8.4 监控与告警

线上系统的可观测性是运维的生命线。我们建立了三层监控体系:

第一层:基础设施监控

使用Prometheus采集CPU、内存、磁盘、网络等基础指标,Grafana做可视化。关键告警规则包括:

  • Pod CPU使用率 > 85% 持续5分钟 → 告警
  • Pod内存使用率 > 90% 持续3分钟 → 告警
  • Pod重启次数 > 3次/小时 → 告警
  • Milvus连接失败 → 立即告警

第二层:应用性能监控(APM)

在应用层埋点采集关键业务指标:

监控指标 采集方式 告警阈值
请求QPS Prometheus Counter 骤降50% → 告警
P95响应延迟 Prometheus Histogram > 5s → 告警
意图识别准确率 日志采样 + 离线评估 < 80% → 告警
工具调用成功率 Prometheus Counter < 95% → 告警
大模型API错误率 Prometheus Counter > 5% → 告警
FAQ缓存命中率 Prometheus Counter < 20% → 告警(可能知识库有问题)
转人工率 Prometheus Counter > 15% → 告警(智能体可能有故障)

第三层:业务质量监控

这是最贴近业务体验的监控层。我们每天自动采样100条对话,用另一个大模型进行自动评分(1-5分),低于3分的对话自动推送到企业微信群,由人工review后决定是否需要优化知识库或调整提示词。

这个机制让我们能够在用户投诉之前发现潜在问题。有一次,自动评分发现多条关于"冬季新品菜单"的对话得分偏低,经检查发现是知识库中新菜品信息录入不完整。我们在2小时内修复了数据,避免了大规模的用户投诉。

8.5 灰度发布与A/B测试

智能体上线不是"一键切换"那么简单。对于一个每天服务数千名顾客的系统,任何变更都可能影响用户体验。我们采用了灰度发布策略,将新版本逐步推给一小部分流量,观察指标后再逐步扩大范围。

灰度发布通过K8s的Service和Istio流量分发实现。具体流程如下:

  1. 金丝雀部署:新版本先部署1个Pod,承接5%的流量。运行24小时,对比新旧版本的关键指标(回答准确率、P95延迟、错误率、用户满意度评分)。
  2. 指标对比:如果新版本的所有指标都不劣于旧版本(允许5%以内的波动),进入下一阶段;否则回滚。
  3. 逐步放量:从5% → 20% → 50% → 100%,每个阶段运行12-24小时,持续监控。
  4. 全量切换:确认100%流量稳定运行48小时后,下线旧版本。

这套灰度流程在第二次版本迭代中发挥了关键作用。那次我们升级了意图识别模块的规则集,金丝雀阶段就发现新规则集对"退款"类意图的识别准确率下降了8%——新规则把一些包含"退"字但不涉及退款的句子(如"退后两步拍照")误判为退款意图。我们在5%流量的阶段就发现了问题,1小时内修复并重新发布,影响范围极小。如果没有灰度发布,这个问题会影响所有用户。

8.6 知识库运营管理平台

技术团队交付的智能体只是一个"半成品"——真正让它持续发挥作用的是运营团队的日常维护。我们为运营团队开发了一个简单的知识库管理界面,让非技术人员也能高效地管理知识内容。

功能模块 核心功能 使用角色 日均操作量
FAQ管理 新增、编辑、删除FAQ条目,支持批量导入 客服主管 ~30条
菜单同步 查看菜品向量化状态,手动触发同步 运营专员 ~2次/周
门店信息 编辑门店基础信息,自动触发向量化更新 运营专员 ~5条/周
对话审查 查看低分对话,标注正确/错误,反馈给技术团队 客服专员 ~100条/天
效果看板 查看准确率、满意度、热门问题等统计 管理层 只读

这个管理平台的核心理念是"让运营人员5分钟内完成知识更新"。以新增一条FAQ为例:运营人员在界面上输入问题和答案,选择分类和关联门店,点击"发布"——系统自动完成向量化、入库和索引更新,整个过程不到30秒。在旧系统中,同样的操作需要技术人员手动修改规则文件、提交代码、等待部署,耗时至少30分钟。

"对话审查"模块是运营和技术之间的桥梁。客服专员每天查看自动评分低于3分的对话,判断是智能体的回答有误还是用户的问题超出了智能体的能力范围。对于回答有误的案例,客服专员可以直接在界面上标注"错误回答"并填写正确答案,这些标注数据每周汇总后用于优化评估集和微调提示词。这个闭环机制让系统的回答质量在上线后持续提升——上线第一个月的回答准确率是85%,第三个月自然增长到了91%,其中约一半的提升来自运营团队的标注反馈。

九、上线后的数据与反思

9.1 运营数据

系统上线运行3个月后,关键运营数据如下:

指标 上线前 上线后 变化
日均处理咨询量 3500条(人工处理) 4200条(智能体处理85%) +20%
平均响应时间 3-15分钟(人工) 2.8秒(智能体P95) 提升300倍
首次解决率 78% 91% +13%
人工客服工作量 3500条/天 525条/天 -85%
客服人力成本 ¥180,000/月 ¥85,000/月(含API成本) -53%
用户满意度评分 4.1/5.0 4.6/5.0 +12%
客服投诉率 3.2% 1.1% -66%

这些数据让张磊非常满意。但作为技术团队,我们更关注的是那些"数字背后的故事"。

9.2 踩过的坑和学到的教训

坑一:知识库更新不及时导致回答错误

上线第二周,有顾客反馈智能体告知的营业时间与实际不符。排查后发现,某门店因装修临时调整了营业时间,但运营人员只在内部OA系统中更新了信息,没有同步到知识库。这个问题暴露了我们的知识库更新流程存在断点。

解决方案:与门店管理系统打通,实现门店信息的自动同步。每天凌晨全量同步一次,关键信息变更(如营业状态)触发实时同步。同时在系统提示词中加入"如果不确定门店是否正常营业,请调用query_store_info工具查询最新状态"的指令。

坑二:大模型"自作主张"编造信息

有用户问"你们有亲子套餐吗?",知识库中没有相关信息,但大模型根据"亲子"+"套餐"的组合,编造了一个不存在的亲子套餐及其价格。这个幻觉问题在第一版中比较严重。

解决方案:在系统提示词中强化"不得编造"的规则,并在生成回答后增加一个验证步骤——检查回答中提到的菜品、价格、活动是否在知识库或工具返回结果中有依据。对于无法验证的内容,替换为"建议您咨询人工客服了解详情"。

坑三:高并发下的Embedding服务瓶颈

上线第一个月的某天午市高峰,系统突然变慢,P95延迟飙升至12秒。排查发现是Embedding服务的GPU显存不足,导致请求排队。由于所有向量化和检索请求都经过同一个Embedding服务,它成为了瓶颈。

解决方案:增加Embedding服务实例到2个(带负载均衡),同时引入Embedding缓存——对于完全相同的查询文本,直接返回缓存的向量结果,避免重复计算。缓存命中率达到约30%,有效缓解了GPU压力。

坑四:多轮对话中的上下文丢失

有用户进行了如下对话:"你们朝阳大悦城店在哪" → "营业到几点" → "有包间吗"。第三轮的回答出现了错误——智能体没有理解"有包间吗"问的是朝阳大悦城店的包间,而是检索到了其他门店的包间信息。

解决方案:在对话编排层增加上下文继承逻辑。当意图识别检测到省略主语的追问时,自动继承上一轮的门店实体,将其加入检索条件。同时在系统提示词中加入"如果用户的问题缺少上下文(如没有指定门店),请参考对话历史确定上下文"的指令。

9.3 架构反思与未来规划

回顾整个项目的架构设计,有一些决策在事后看来可以做得更好:

架构决策 当时的选择 反思 未来优化方向
对话状态管理 内存 + Redis 多实例间的状态一致性有挑战 引入分布式session store
知识库更新 定时全量 + 手动增量 实时性不够,依赖人工触发 事件驱动的增量更新
模型路由策略 基于意图的简单路由 粒度太粗,部分场景质量不稳定 基于问题复杂度的动态路由
评估体系 离线评估 + 日采样 反馈链路长,发现滞后 在线A/B测试框架
多模态能力 仅支持文本 无法处理图片(如菜品照片识别) 引入多模态大模型

未来3-6个月的规划包括:

  1. 引入多模态能力:支持用户发送菜品照片,智能体识别菜品并返回相关信息
  2. 个性化推荐:基于用户历史对话和消费记录,提供个性化菜品推荐
  3. 主动服务:在用户等待排队时主动推送预估等待时间和附近商圈推荐
  4. 多语言支持:增加英语和日语支持,服务外籍顾客
  5. 知识图谱:构建菜品-食材-过敏原-口味的知识图谱,提升复杂查询能力

十、总结:构建智能体的核心方法论

回顾整个项目,我想提炼出几条构建行业智能体的核心方法论,希望能对你有所启发:

第一,需求分析决定上限,技术实现决定下限。很多人沉迷于技术选型和模型调优,却忽略了最基础的需求分析。我们在需求分析上花了两整天,但这两整天的工作决定了整个项目的方向。如果需求分析做错了,后面的技术做得再好也是南辕北辙。

第二,数据质量 > 模型能力。在RAG架构中,知识库的质量远比大模型的能力重要。一个用GPT-4但知识库混乱的系统,表现不如一个用开源模型但知识库精心治理的系统。把80%的精力放在数据治理上,20%放在模型调优上,这是我们的经验之谈。

第三,分层解耦是可演进架构的基石。我们的五层架构让每一层都可以独立替换。上线后我们更换了大模型供应商、升级了Milvus版本、增加了Embedding缓存,每一次变更都只影响一个层,没有波及其他模块。如果没有这种解耦设计,每一次变更都会是一场噩梦。

第四,降级策略是生产可用的底线。在开发阶段,一切运行良好;到了生产环境,什么都有可能出问题。大模型API可能宕机、向量数据库可能OOM、网络可能抖动。每一层都要有降级方案,确保在最坏情况下系统仍然可用(哪怕体验下降)。

第五,评估体系是持续优化的引擎。没有度量就没有优化。从第一天就要建立评估集和评估流程,让每一次优化都有数据支撑。我们的评估集从最初的500条扩展到2000条,每一次优化都在评估集上验证效果后才上线。

第六,工程能力和AI能力同等重要。智能体不只是大模型API的调用——它是一个完整的系统工程,涉及网络编程、并发处理、缓存策略、监控告警、容器化部署等大量传统工程能力。一个只有AI能力但缺乏工程能力的团队,很难构建出生产可用的智能体系统。

最后,回到张磊的故事。系统上线3个月后,他请我吃了一顿饭。席间他说了一句话让我印象深刻:"这个智能体最大的价值不是省了多少客服人力,而是让我们的顾客在任何时间、任何渠道都能得到即时、准确的回应。这种体验提升带来的品牌好感度,是没法用钱衡量的。"

这也许就是AI落地餐饮行业的真正意义——不是替代人,而是让服务触手可及。当一位深夜加班的年轻人在小程序上询问"你们还有夜宵吗",在几秒内收到"三里屯店营业至凌晨两点,推荐您试试我们的招牌牛肉面,现在下单还有夜宵专属八八折优惠"这样的回复时,技术的温度便不再是抽象的概念,而是实实在在融入了人们的生活日常。


如果你觉得这篇文章对你有帮助,欢迎点赞、收藏和关注。也欢迎在评论区分享你的智能体构建经验。

推荐阅读:

Logo

有“AI”的1024 = 2048,欢迎大家加入2048 AI社区

更多推荐