课程0基础Agent开发课 / RAG与向量数据库 / RAG生产架构-从原型到企业级系统
— 24 min read

RAG生产架构-从原型到企业级系统

> **[进阶选读]** 本篇适合 RAG 原型已验证有效、准备推向生产环境的读者。涵盖并发处理、增量索引更新、语义缓存、监控与成本控制等工程化话题。

RAG 生产架构:从原型到企业级系统

[进阶选读] 本篇适合 RAG 原型已验证有效、准备推向生产环境的读者。涵盖并发处理、增量索引更新、语义缓存、监控与成本控制等工程化话题。

原型 RAG 搭起来不难:几十行代码,本地跑通,效果演示。但把它放到生产环境,面向真实用户流量时,问题才真正开始。

原型和生产之间有一条看不见的鸿沟,这条鸿沟不是技术深度,而是工程维度——并发、更新、监控、安全、成本。这篇文章专门讨论这条鸿沟里有什么,以及怎么过。


1.1 原型 RAG 和生产 RAG 的差距

RAG 生产架构图
RAG 生产架构演进——从简单向量检索的原型,到具备语义缓存、权限控制、监控和增量更新能力的企业级系统

先把两者的差距具体化:

维度 原型 RAG 生产 RAG
并发量 1 个用户,串行请求 数百并发,峰值冲击
文档更新 手动重跑索引脚本 增量更新,不停服
查询缓存 没有 语义缓存,减少重复调用
权限控制 所有人看所有文档 不同角色看不同内容
监控 控制台 print 召回率、准确率、响应时间指标
错误处理 报错就崩 降级、重试、告警
成本控制 不关心 每次调用都是钱

原型阶段可以忽略这七个维度。生产阶段,每一个都是潜在的事故点。


1.2 生产级 RAG 的组件清单

一个能承载真实业务的 RAG 系统,需要以下组件:

查询服务(在线)

文档处理流水线(离线)

命中

未命中

文档源
S3/数据库/CMS

解析服务
PDF/Word/HTML

分块服务
策略化分块

Embedding 服务
批量向量化

向量数据库
Milvus/PGVector

元数据库
权限/版本/来源

用户请求

鉴权层
用户权限校验

语义缓存
命中?

直接返回缓存

查询改写
可选

混合检索
向量+BM25

重排序
Reranker

LLM 生成答案
+ 引用来源

写入缓存

监控上报


1.3 文档处理流水线

1.3.1 全量重建 vs 增量更新

原型阶段通常是全量重建:清空向量库,重新 embedding 所有文档。文档少的时候没问题,几千份文档就不行了——耗时几小时,期间服务降级。

生产系统需要增量更新

python
# document_pipeline.py
# 增量更新文档索引的设计

import hashlib
from datetime import datetime
from typing import List, Dict, Optional
from dataclasses import dataclass


@dataclass
class DocumentRecord:
    """文档记录:追踪每份文档的版本状态。"""
    doc_id: str
    source_path: str
    content_hash: str           # 内容的 MD5(一种将任意文本转成固定长度指纹的算法,内容变化则指纹变化),用于判断是否变化
    chunk_ids: List[str]        # 该文档对应的所有 chunk ID
    indexed_at: datetime
    version: int


class IncrementalIndexer:
    """
    增量索引器:只处理新增或修改的文档。
    设计原则:
    1. 每份文档有唯一 ID 和内容 hash
    2. 处理前先对比 hash,相同则跳过
    3. 更新时先删除旧 chunk,再插入新 chunk
    """

    def __init__(self, vector_store, metadata_db):
        self.vector_store = vector_store    # 向量数据库客户端
        self.metadata_db = metadata_db      # 元数据库(记录文档版本信息)

    def compute_hash(self, content: str) -> str:
        """计算文档内容的 MD5 hash,用于变更检测。"""
        return hashlib.md5(content.encode()).hexdigest()

    def process_document(self, doc_id: str, content: str, metadata: dict) -> str:
        """
        处理单份文档。
        返回操作类型:'skipped'(未变化)、'created'(新增)、'updated'(更新)
        """
        new_hash = self.compute_hash(content)

        # 查询元数据库,看这份文档是否已经索引过
        existing_record: Optional[DocumentRecord] = self.metadata_db.get(doc_id)

        if existing_record and existing_record.content_hash == new_hash:
            # hash 相同,内容未变化,跳过处理
            # 这是增量更新的关键:避免对未变化文档重复 embedding
            return "skipped"

        # 文档有变化(新增或修改)
        if existing_record:
            # 更新:先删除旧的 chunk
            for old_chunk_id in existing_record.chunk_ids:
                self.vector_store.delete(old_chunk_id)

        # 分块并 embedding(这里省略具体实现,复用已有的 chunking 逻辑)
        chunks = self._chunk_document(content)
        new_chunk_ids = []

        for i, chunk in enumerate(chunks):
            chunk_id = f"{doc_id}_chunk_{i}"
            # 向量化并存储(实际调用 embedding 模型)
            # embedding = embedding_model.embed(chunk)
            # self.vector_store.upsert(chunk_id, embedding, {"text": chunk, **metadata})
            new_chunk_ids.append(chunk_id)

        # 更新元数据记录
        self.metadata_db.upsert(DocumentRecord(
            doc_id=doc_id,
            source_path=metadata.get("source", ""),
            content_hash=new_hash,
            chunk_ids=new_chunk_ids,
            indexed_at=datetime.now(),
            version=(existing_record.version + 1) if existing_record else 1,
        ))

        return "updated" if existing_record else "created"

    def _chunk_document(self, content: str) -> List[str]:
        """分块逻辑(简化版)。"""
        # 实际替换为 RecursiveCharacterTextSplitter 等
        chunk_size = 500
        return [content[i:i+chunk_size] for i in range(0, len(content), chunk_size)]

    def batch_process(self, documents: List[Dict]) -> Dict[str, int]:
        """
        批量处理文档,返回各操作类型的统计。
        生产环境通常作为定时任务或消息队列消费者运行。
        """
        stats = {"skipped": 0, "created": 0, "updated": 0, "failed": 0}

        for doc in documents:
            try:
                result = self.process_document(
                    doc_id=doc["id"],
                    content=doc["content"],
                    metadata=doc.get("metadata", {})
                )
                stats[result] += 1
            except Exception as e:
                stats["failed"] += 1
                print(f"处理文档 {doc['id']} 失败:{e}")

        return stats

1.4 语义缓存

RAG 查询的成本主要来自两部分:embedding 调用(便宜)和 LLM 生成(贵)。

如果同一个问题或高度相似的问题被反复问,每次都调用 LLM 是纯粹的浪费。语义缓存(Semantic Cache,一种将问题的向量表示与答案一起存储的缓存机制,新问题到来时先查相似历史问题,命中则直接返回缓存答案)在 embedding 空间里做相似度匹配:如果新问题和缓存中某个问题的向量相似度超过阈值,直接返回缓存答案,不调用 LLM。

这比关键词完全匹配的传统缓存覆盖率高得多——"怎么申请退货"和"如何办理退货手续"会命中同一个缓存条目。

python
# semantic_cache.py
# 语义缓存实现:利用向量相似度实现模糊命中

import time
import hashlib
from typing import Optional, List
from dataclasses import dataclass, field


@dataclass
class CacheEntry:
    query: str
    query_embedding: List[float]
    answer: str
    created_at: float
    hit_count: int = 0
    ttl_seconds: int = 3600  # 缓存 1 小时


class SemanticCache:
    """
    语义缓存:相似查询共享同一个缓存条目。
    核心思路:把历史查询的 embedding 存起来,
    新查询来时先算相似度,高于阈值就返回缓存。
    """

    def __init__(self, similarity_threshold: float = 0.92, max_size: int = 1000):
        # similarity_threshold:相似度阈值,越高要求越严格
        # 0.92 是经验值:能覆盖语义相同但表述不同的问题,同时避免误命中
        self.threshold = similarity_threshold
        self.max_size = max_size
        self.entries: List[CacheEntry] = []

    def _cosine_similarity(self, vec1: List[float], vec2: List[float]) -> float:
        """计算余弦相似度。生产环境通常用 numpy 或向量数据库的内置函数。"""
        dot_product = sum(a * b for a, b in zip(vec1, vec2))
        norm1 = sum(a ** 2 for a in vec1) ** 0.5
        norm2 = sum(b ** 2 for b in vec2) ** 0.5
        if norm1 == 0 or norm2 == 0:
            return 0.0
        return dot_product / (norm1 * norm2)

    def get(self, query_embedding: List[float]) -> Optional[str]:
        """
        查询缓存。
        返回缓存答案(命中时)或 None(未命中时)。
        """
        now = time.time()

        best_similarity = 0.0
        best_entry = None

        for entry in self.entries:
            # 检查 TTL
            if now - entry.created_at > entry.ttl_seconds:
                continue  # 过期,跳过(惰性删除)

            similarity = self._cosine_similarity(query_embedding, entry.query_embedding)
            if similarity > best_similarity:
                best_similarity = similarity
                best_entry = entry

        if best_entry and best_similarity >= self.threshold:
            best_entry.hit_count += 1
            return best_entry.answer

        return None  # 未命中

    def set(self, query: str, query_embedding: List[float], answer: str):
        """写入缓存条目。超过最大容量时淘汰命中次数最少的条目(LFU 策略)。"""
        if len(self.entries) >= self.max_size:
            # LFU(Least Frequently Used,最少使用淘汰策略——缓存满时删掉使用次数最少的条目)
            self.entries.sort(key=lambda e: e.hit_count)
            self.entries.pop(0)

        self.entries.append(CacheEntry(
            query=query,
            query_embedding=query_embedding,
            answer=answer,
            created_at=time.time(),
        ))

    def get_stats(self) -> dict:
        """缓存统计信息,用于监控。"""
        now = time.time()
        active = [e for e in self.entries if now - e.created_at <= e.ttl_seconds]
        return {
            "total_entries": len(self.entries),
            "active_entries": len(active),
            "total_hits": sum(e.hit_count for e in active),
        }

1.5 权限控制:不同用户看不同文档

企业场景里,员工不应该看到不属于自己权限范围的文档。RAG 系统的权限控制需要在检索层而非答案层实现——不是"查出来了但不展示",而是"根本查不出来没有权限的内容"。

python
# access_control.py
# 基于元数据过滤的 RAG 权限控制

from typing import List, Dict


class PermissionAwareRetriever:
    """
    带权限控制的检索器。
    核心:在向量检索时附加元数据过滤条件,
    确保只检索当前用户有权访问的文档。
    """

    def __init__(self, vector_store, user_permission_service):
        self.vector_store = vector_store
        self.permission_service = user_permission_service

    def get_user_accessible_filter(self, user_id: str) -> dict:
        """
        获取当前用户的访问权限过滤条件。
        实际场景从权限服务或数据库查询用户的角色和可访问的文档集合。
        """
        user_roles = self.permission_service.get_roles(user_id)
        accessible_doc_ids = self.permission_service.get_accessible_docs(user_id)

        # 返回向量数据库的元数据过滤条件(各数据库语法不同)
        # Chroma 格式:
        return {
            "$or": [
                {"access_level": "public"},                          # 公开文档
                {"owner_id": user_id},                               # 用户自己的文档
                {"doc_id": {"$in": accessible_doc_ids}},             # 授权文档
                {"required_roles": {"$in": user_roles}},             # 角色匹配
            ]
        }

    def retrieve(self, query_embedding: List[float], user_id: str, k: int = 5) -> List[dict]:
        """
        带权限过滤的检索。
        向量数据库在执行相似度搜索的同时应用元数据过滤,
        这比"先检索、再过滤"的方式性能好得多。
        """
        filter_condition = self.get_user_accessible_filter(user_id)

        # 向量数据库支持在查询时附加过滤条件
        # Chroma 示例:
        # results = self.vector_store.similarity_search_by_vector(
        #     embedding=query_embedding,
        #     k=k,
        #     filter=filter_condition
        # )
        #
        # Milvus 示例(使用 expr 参数):
        # results = self.vector_store.query(
        #     data=[query_embedding],
        #     limit=k,
        #     filter=f"access_level == 'public' or owner_id == '{user_id}'"
        # )

        # 模拟返回(实际替换为真实向量数据库调用)
        return [
            {"text": "公开文档内容", "metadata": {"access_level": "public"}},
        ]


# 文档索引时附加权限元数据
def index_document_with_permissions(
    content: str,
    doc_id: str,
    owner_id: str,
    access_level: str,          # "public" | "private" | "role_based"
    required_roles: List[str],  # 当 access_level="role_based" 时生效
    vector_store,
    embedding_model,
):
    """
    带权限元数据的文档索引。
    关键:把权限信息作为元数据和向量一起存储,检索时用于过滤。
    """
    embedding = embedding_model.embed(content)  # 替换为真实 embedding 调用

    vector_store.upsert(
        id=doc_id,
        vector=embedding,
        metadata={
            "text": content,
            "doc_id": doc_id,
            "owner_id": owner_id,
            "access_level": access_level,
            "required_roles": required_roles,
        }
    )

1.6 监控指标

没有监控的 RAG 系统是黑盒。以下是最关键的监控指标:

python
# rag_monitoring.py
# RAG 系统关键监控指标

import time
from dataclasses import dataclass, field
from typing import List, Optional
from collections import defaultdict


@dataclass
class RAGQueryMetrics:
    """单次查询的指标记录。"""
    query_id: str
    query_text: str
    user_id: str

    # 耗时指标(毫秒)
    retrieval_latency_ms: float = 0    # 检索耗时
    rerank_latency_ms: float = 0       # 重排序耗时
    llm_latency_ms: float = 0          # LLM 生成耗时
    total_latency_ms: float = 0        # 总耗时

    # 检索质量指标
    retrieved_chunks: int = 0          # 检索到的 chunk 数量
    avg_retrieval_score: float = 0     # 平均检索相似度分
    cache_hit: bool = False            # 是否命中缓存

    # 答案质量指标(需要后处理或用户反馈)
    has_answer: bool = True            # LLM 是否给出了答案(非"我不知道")
    user_feedback: Optional[int] = None  # 用户评分(1-5),有则记录


class RAGMonitor:
    """
    RAG 监控聚合器。
    实际生产中接入 Prometheus(开源监控系统,采集指标数据)+ Grafana(开源可视化仪表盘,将指标数据展示成图表),
    或 Datadog、New Relic 等 APM(Application Performance Monitoring,应用性能监控)工具。
    这里演示核心指标的计算逻辑。
    """

    def __init__(self):
        self.metrics_buffer: List[RAGQueryMetrics] = []

    def record(self, metrics: RAGQueryMetrics):
        """记录一次查询的指标。"""
        self.metrics_buffer.append(metrics)

        # 实时告警:响应时间超过 5 秒
        if metrics.total_latency_ms > 5000:
            self._alert(f"慢查询告警:{metrics.query_id} 耗时 {metrics.total_latency_ms:.0f}ms")

    def get_summary(self, last_n: int = 100) -> dict:
        """
        获取最近 N 次查询的汇总指标。
        这些指标应该定期上报到监控系统。
        """
        recent = self.metrics_buffer[-last_n:]
        if not recent:
            return {}

        total = len(recent)
        cache_hits = sum(1 for m in recent if m.cache_hit)
        has_answers = sum(1 for m in recent if m.has_answer)
        latencies = [m.total_latency_ms for m in recent]
        feedback = [m.user_feedback for m in recent if m.user_feedback is not None]

        # P99 延迟(第99百分位延迟):99% 的请求都在这个时间内完成,是衡量系统最差情况响应速度的指标
        sorted_latencies = sorted(latencies)
        p50 = sorted_latencies[int(total * 0.5)]
        p95 = sorted_latencies[int(total * 0.95)]
        p99 = sorted_latencies[int(total * 0.99)] if total >= 100 else sorted_latencies[-1]

        return {
            "query_count": total,
            "cache_hit_rate": cache_hits / total,           # 缓存命中率,越高越省钱
            "answer_rate": has_answers / total,              # 回答率,过低说明检索质量差
            "avg_latency_ms": sum(latencies) / total,
            "p50_latency_ms": p50,
            "p95_latency_ms": p95,
            "p99_latency_ms": p99,
            "avg_user_rating": sum(feedback) / len(feedback) if feedback else None,
        }

    def _alert(self, message: str):
        """发送告警(实际接入钉钉/企微/PagerDuty)。"""
        print(f"[ALERT] {message}")


# 在 RAG 查询流程中埋点
def rag_query_with_monitoring(
    query: str,
    user_id: str,
    retriever,
    llm,
    monitor: RAGMonitor,
    semantic_cache: SemanticCache,
    embedding_model,
) -> str:
    """
    带监控埋点的 RAG 查询函数。
    关键:在每个阶段记录时间,方便定位性能瓶颈。
    """
    import uuid
    query_id = str(uuid.uuid4())[:8]

    metrics = RAGQueryMetrics(
        query_id=query_id,
        query_text=query,
        user_id=user_id,
    )

    total_start = time.time()

    # 1. 语义缓存查询
    query_embedding = embedding_model.embed(query)  # 替换为真实 embedding 调用
    cached_answer = semantic_cache.get(query_embedding)
    if cached_answer:
        metrics.cache_hit = True
        metrics.total_latency_ms = (time.time() - total_start) * 1000
        monitor.record(metrics)
        return cached_answer

    # 2. 检索阶段
    retrieval_start = time.time()
    # chunks = retriever.retrieve(query_embedding, user_id=user_id)  # 真实检索调用
    chunks = [{"text": f"相关文档片段关于{query}"}]  # 模拟
    metrics.retrieval_latency_ms = (time.time() - retrieval_start) * 1000
    metrics.retrieved_chunks = len(chunks)

    # 3. LLM 生成阶段
    llm_start = time.time()
    context = "\n".join(c["text"] for c in chunks)
    # answer = llm.invoke(f"基于以下信息回答:{context}\n\n问题:{query}")  # 真实 LLM 调用
    answer = f"基于检索内容,关于「{query}」的回答:..."  # 模拟
    metrics.llm_latency_ms = (time.time() - llm_start) * 1000
    metrics.has_answer = not answer.startswith("我不知道")

    # 4. 写入缓存
    semantic_cache.set(query, query_embedding, answer)

    metrics.total_latency_ms = (time.time() - total_start) * 1000
    monitor.record(metrics)

    return answer

1.7 常见性能问题与优化

问题一:首次查询慢(冷启动)

原因:embedding 模型首次加载需要时间,向量数据库连接池未预热。
优化:服务启动时预加载模型,建立数据库连接池,对 embedding 模型做 warmup 请求。

问题二:大批量文档更新时服务不稳定

原因:批量 embedding 请求冲击 API 限流,或者向量库写入压力过大。
优化:限制批量处理的并发数,使用令牌桶算法(一种流量控制方法,把 API 调用配额想象成令牌,有令牌才能发请求,按固定速率补充令牌,防止突发流量打垮 API)控制 API 调用频率;向量库写入走异步队列,与在线查询服务隔离。

问题三:LLM 生成延迟高

原因:context 过长(塞了太多 chunk),或者 LLM 服务本身繁忙。
优化:控制送入 LLM 的 chunk 数量(3-5 个通常足够),启用流式输出让用户感知更快,接入备用 LLM 服务。

问题四:检索召回率随文档量增加而下降

原因:文档量大时,向量空间密度增加,相似度阈值需要调整。
优化:定期评估检索质量(RAGAS 或人工评估),动态调整 top-k 参数,考虑加重排序层。


生产级 RAG 和原型 RAG 的差距不在算法,在工程:

  • 增量更新:内容 hash 比对,只处理变化的文档
  • 语义缓存:相似问题共享答案,降低 LLM 调用成本
  • 权限控制:在检索层过滤,不是在答案层过滤
  • 全链路监控:每个阶段埋点,定位性能瓶颈

一段跑通的代码不是系统。系统要有观测、有告警、有更新机制,才能长期稳定运行。

本页目录