RAG 2.0与Agentic检索增强生成系统架构深度实战:从多路召回排序到自查询与知识图谱融合的下一代企业知识引擎全解析

举报
江南清风起 发表于 2026/08/30 22:20:25 2026/08/30
【摘要】 RAG 2.0与Agentic检索增强生成系统架构深度实战:从多路召回排序到自查询与知识图谱融合的下一代企业知识引擎全解析 引言检索增强生成(RAG)从2023年的"向量库+LLM"简单拼装,已演进为2026年的Agentic RAG(RAG 2.0)体系:检索不再是一次性召回,而是Agent驱动的多轮迭代探索;数据源不再只有向量索引,而是混合了全文检索、知识图谱、SQL表格、结构化API...

RAG 2.0与Agentic检索增强生成系统架构深度实战:从多路召回排序到自查询与知识图谱融合的下一代企业知识引擎全解析

引言

检索增强生成(RAG)从2023年的"向量库+LLM"简单拼装,已演进为2026年的Agentic RAG(RAG 2.0)体系:检索不再是一次性召回,而是Agent驱动的多轮迭代探索;数据源不再只有向量索引,而是混合了全文检索、知识图谱、SQL表格、结构化API的多路召回;排序不再只靠相似度,而是引入交叉编码器重排与引用验证;生成不再只输出文本,而是附置信度与可追溯的引用链。企业落地的核心矛盾是:简单的向量RAG准确率天花板低(长尾问题召回率不足60%),而复杂Agentic RAG的延迟与成本又难以承受。本文将系统讲解RAG 2.0的完整工程体系,覆盖文档处理管道、多路混合检索、查询重写与自查询、重排序模型、引用验证与置信度、知识图谱融合、Agentic RAG编排、评估体系与成本优化,并提供Python/TypeScript可运行代码。

一、RAG系统架构总览与演进

1.1 从Naive RAG到Agentic RAG

Naive RAG(2023范式)的流程是:文档切块→嵌入→向量入库→查询嵌入→余弦相似度Top-K→拼接到Prompt→LLM生成。其缺陷已被工程实践充分暴露:切块破坏语义边界(表格、列表被切断)、单一向量召回长尾覆盖差、无重排导致相关但排序靠后的结果丢失、无法回答需要多跳推理的复杂问题、缺乏可验证性(幻觉难以追溯)。Advanced RAG引入了查询重写、混合检索(BM25+向量)、重排序(Cross-Encoder),显著提升召回质量但仍是单轮流水线。Agentic RAG(RAG 2.0)把检索本身建模为Agent的决策循环:模型自主决定何时检索、检索什么、是否需要多轮、何时停止并综合答案,并能调用工具(SQL查询、API调用、知识图谱遍历)而非仅向量搜索。

1.2 企业RAG架构分层

成熟的企业RAG系统分五层:数据层(文档摄入管道、知识图谱、结构化数据库)、索引层(向量索引、全文索引、图索引、列式索引的多路并存)、检索层(查询规划、多路召回、融合排序、重排)、生成层(上下文压缩、引用标注、置信度评估、自校验)、编排层(Agent决策循环、工具调用、会话记忆、成本控制)。以下定义系统核心类型与编排入口:

# rag/system.py - Agentic RAG系统核心
from __future__ import annotations
from dataclasses import dataclass, field
from typing import Any, Literal, Protocol
from enum import Enum
import asyncio

class RetrievalStrategy(str, Enum):
    VECTOR = "vector"
    FULLTEXT = "fulltext"
    HYBRID = "hybrid"
    GRAPH = "graph"
    SQL = "sql"
    TOOL = "tool"

@dataclass
class Document:
    id: str
    content: str
    metadata: dict[str, Any] = field(default_factory=dict)
    source: str = ""               # 文件路径或URL
    chunk_index: int = 0
    parent_id: str | None = None    # 父文档ID(用于上下文扩展)
    embedding: list[float] | None = None

@dataclass
class RetrievedChunk:
    document: Document
    score: float
    retrieval_strategy: RetrievalStrategy
    rerank_score: float | None = None    # 重排后分数
    citations: list[str] = field(default_factory=list)

@dataclass
class SearchQuery:
    original: str                     # 用户原始查询
    rewritten: str                     # 重写后的检索查询
    sub_queries: list[str] = field(default_factory=list)  # 子查询(多跳分解)
    filters: dict[str, Any] = field(default_factory=dict) # 元数据过滤
    strategy_hint: RetrievalStrategy | None = None

@dataclass
class RAGResponse:
    answer: str
    citations: list[Citation]
    confidence: float
    retrieval_trace: list[dict]       # 检索轨迹(可解释性)
    follow_up_suggestions: list[str] = field(default_factory=list)

@dataclass
class Citation:
    document_id: str
    source: str
    chunk_text: str
    relevance: float
    page: int | None = None

class Retriever(Protocol):
    async def retrieve(self, query: SearchQuery, top_k: int) -> list[RetrievedChunk]: ...

class Reranker(Protocol):
    async def rerank(self, query: str, chunks: list[RetrievedChunk]) -> list[RetrievedChunk]: ...

class QueryPlanner(Protocol):
    async def plan(self, user_query: str, conversation: list[dict]) -> SearchQuery: ...

class Generator(Protocol):
    async def generate(self, query: str, context: list[RetrievedChunk],
                       conversation: list[dict]) -> tuple[str, list[Citation], float]: ...

class AgenticRAG:
    """Agentic RAG编排器:驱动检索-评估-再检索的决策循环"""
    def __init__(
        self,
        planner: QueryPlanner,
        retrievers: dict[RetrievalStrategy, Retriever],
        reranker: Reranker,
        generator: Generator,
        max_iterations: int = 3,
        min_confidence: float = 0.65,
    ):
        self.planner = planner
        self.retrievers = retrievers
        self.reranker = reranker
        self.generator = generator
        self.max_iterations = max_iterations
        self.min_confidence = min_confidence

    async def answer(
        self, user_query: str, conversation: list[dict] | None = None,
    ) -> RAGResponse:
        conversation = conversation or []
        trace: list[dict] = []

        # 第一阶段:查询规划
        search_query = await self.planner.plan(user_query, conversation)
        trace.append({"phase": "planning", "query": search_query})

        all_chunks: list[RetrievedChunk] = []
        final_answer = ""
        final_citations: list[Citation] = []
        confidence = 0.0

        for iteration in range(self.max_iterations):
            # 第二阶段:多路并行检索
            strategies = self._select_strategies(search_query, iteration)
            tasks = [
                self.retrievers[s].retrieve(search_query, top_k=20)
                for s in strategies if s in self.retrievers
            ]
            results = await asyncio.gather(*tasks, return_exceptions=True)

            iteration_chunks: list[RetrievedChunk] = []
            for strategy, result in zip(strategies, results):
                if isinstance(result, Exception):
                    trace.append({"phase": "retrieve", "strategy": strategy,
                                  "error": str(result)})
                    continue
                iteration_chunks.extend(result)
                trace.append({"phase": "retrieve", "strategy": strategy,
                              "count": len(result)})

            # 去重与合并
            seen_ids = {c.document.id for c in all_chunks}
            iteration_chunks = [
                c for c in iteration_chunks if c.document.id not in seen_ids
            ]
            all_chunks.extend(iteration_chunks)

            # 第三阶段:重排序
            if all_chunks:
                all_chunks = await self.reranker.rerank(
                    search_query.rewritten, all_chunks,
                )
                all_chunks = all_chunks[:10]  # 截断到Top-10
                trace.append({"phase": "rerank", "kept": len(all_chunks),
                              "top_score": all_chunks[0].rerank_score})

            # 第四阶段:生成与置信度评估
            answer, citations, confidence = await self.generator.generate(
                search_query.rewritten, all_chunks, conversation,
            )
            trace.append({"phase": "generate", "confidence": confidence,
                          "citations": len(citations)})

            # 第五阶段:决策——是否需要再检索
            if confidence >= self.min_confidence or iteration == self.max_iterations - 1:
                final_answer = answer
                final_citations = citations
                break

            # 低置信度:让生成器给出需要补充检索的方向
            gap = await self.generator.identify_knowledge_gap(answer, all_chunks)
            trace.append({"phase": "gap_analysis", "gap": gap})
            search_query.sub_queries = gap.sub_queries
            search_query.rewritten = gap.refined_query

        return RAGResponse(
            answer=final_answer,
            citations=final_citations,
            confidence=confidence,
            retrieval_trace=trace,
            follow_up_suggestions=await self.generator.suggest_followups(
                final_answer, final_citations,
            ) if final_answer else [],
        )

    def _select_strategies(self, query: SearchQuery, iteration: int) -> list[RetrievalStrategy]:
        if query.strategy_hint:
            return [query.strategy_hint]
        if iteration == 0:
            return [RetrievalStrategy.HYBRID, RetrievalStrategy.GRAPH]
        return [RetrievalStrategy.VECTOR, RetrievalStrategy.FULLTEXT]

二、文档处理管道:从摄入到结构化索引

2.1 智能切块策略

切块质量直接决定召回上限。生产级管道采用层级切块(parent-child)+ 语义边界感知:

# rag/chunking.py - 层级语义切块
from __future__ import annotations
from dataclasses import dataclass
from typing import Iterator
import re

@dataclass
class ChunkConfig:
    target_size: int = 512        # 目标token数
    min_size: int = 128           # 最小块(太小不独立索引)
    max_size: int = 1024          # 超过则强制切分
    overlap: int = 64             # 重叠token数
    parent_window: int = 3        # 父文档窗口(子块前后各几个子块合并为父)

class SemanticChunker:
    """基于文档结构的语义切块:尊重标题、段落、列表、表格边界"""

    def __init__(self, config: ChunkConfig):
        self.config = config

    def chunk_document(self, content: str, source: str) -> list[Document]:
        # 1. 按Markdown结构分节
        sections = self._split_by_headers(content)
        # 2. 每节内按语义单元切分
        child_chunks: list[Document] = []
        for section in sections:
            units = self._split_to_semantic_units(section.text)
            for unit_text in self._merge_units(units):
                child_chunks.append(Document(
                    id=f"{source}#{len(child_chunks)}",
                    content=unit_text,
                    metadata={"section": section.title, "level": section.level},
                    source=source,
                    chunk_index=len(child_chunks),
                ))
        # 3. 生成父文档(上下文窗口)
        parent_map: dict[str, str] = {}
        for i, child in enumerate(child_chunks):
            start = max(0, i - self.config.parent_window)
            end = min(len(child_chunks), i + self.config.parent_window + 1)
            parent_text = "\n\n".join(
                child_chunks[j].content for j in range(start, end)
            )
            parent_id = f"{source}#parent-{start}-{end}"
            child.parent_id = parent_id
            parent_map[parent_id] = parent_text
        # 4. 同时索引父文档(检索命中子块时返回父块提供完整上下文)
        for pid, ptext in parent_map.items():
            child_chunks.append(Document(
                id=pid, content=ptext, source=source,
                metadata={"is_parent": True},
            ))
        return child_chunks

    def _split_by_headers(self, content: str) -> list[Section]:
        sections: list[Section] = []
        current = Section(title="", level=0, text="")
        for line in content.split("\n"):
            match = re.match(r"^(#{1,6})\s+(.+)$", line)
            if match:
                if current.text.strip():
                    sections.append(current)
                current = Section(
                    title=match.group(2), level=len(match.group(1)), text="",
                )
            else:
                current.text += line + "\n"
        if current.text.strip():
            sections.append(current)
        return sections

    def _split_to_semantic_units(self, text: str) -> list[str]:
        # 按段落、列表项、表格行拆分,保持完整单元
        units: list[str] = []
        buffer = ""
        in_table = False
        for line in text.split("\n"):
            if line.strip().startswith("|"):
                in_table = True
                buffer += line + "\n"
                continue
            elif in_table and not line.strip():
                if buffer.strip():
                    units.append(buffer)
                buffer = ""
                in_table = False
            if line.strip() == "":
                if buffer.strip():
                    units.append(buffer)
                buffer = ""
            elif re.match(r"^[-*]\s", line) or re.match(r"^\d+\.\s", line):
                if buffer.strip() and not re.match(r"^[-*]\s", buffer):
                    units.append(buffer)
                buffer = line + "\n"
            else:
                buffer += line + "\n"
        if buffer.strip():
            units.append(buffer)
        return units

    def _merge_units(self, units: list[str]) -> Iterator[str]:
        """将小单元合并到目标大小,大单元按句切分"""
        buffer = ""
        for unit in units:
            if self._estimate_tokens(unit) > self.config.max_size:
                if buffer:
                    yield buffer
                    buffer = ""
                yield from self._force_split(unit)
            elif self._estimate_tokens(buffer + unit) > self.config.target_size:
                if buffer:
                    yield buffer
                buffer = unit
            else:
                buffer += unit
        if buffer.strip():
            yield buffer

    def _force_split(self, text: str) -> Iterator[str]:
        sentences = re.split(r"(?<=[。!?.!?])\s*", text)
        buffer = ""
        for s in sentences:
            if self._estimate_tokens(buffer + s) > self.config.max_size:
                if buffer:
                    yield buffer
                buffer = s
            else:
                buffer += s
        if buffer:
            yield buffer

    def _estimate_tokens(self, text: str) -> int:
        # 粗略估算:中文按字数,英文按词数×1.3
        cjk = len(re.findall(r"[\u4e00-\u9fff]", text))
        en = len(text.split())
        return cjk + int(en * 1.3)

@dataclass
class Section:
    title: str
    level: int
    text: str

2.2 元数据提取与结构化

高质量元数据是精准过滤的前提。文档摄入时自动提取:来源类型(policy/manual/api_doc/issue)、部门、版本、生效日期、关键实体(产品名、模块名)、权限标签。以下用LLM辅助提取结构化元数据:

# rag/metadata_extractor.py - LLM辅助元数据抽取
from __future__ import annotations
import json
from pydantic import BaseModel, Field

class DocumentMetadata(BaseModel):
    doc_type: str = Field(description="policy|manual|api_doc|runbook|faq|spec")
    title: str
    summary: str = Field(description="一句话摘要")
    key_entities: list[str] = Field(default_factory=list,
        description="涉及的产品名、模块名、API名")
    effective_date: str | None = None
    expiry_date: str | None = None
    access_tags: list[str] = Field(default_factory=list,
        description="权限标签:public|internal|confidential")
    language: str = "zh"

EXTRACTION_PROMPT = """从以下文档内容提取结构化元数据。
仅基于文档明确信息,不要推测。日期格式YYYY-MM-DD。

文档内容:
{content}

输出JSON,字段:
- doc_type: 文档类型
- title: 文档标题
- summary: 不超过80字的摘要
- key_entities: 关键实体列表
- effective_date: 生效日期(如有)
- expiry_date: 失效日期(如有)
- access_tags: 权限标签
- language: 语言代码"""

async def extract_metadata(content: str, llm_client) -> DocumentMetadata:
    prompt = EXTRACTION_PROMPT.format(content=content[:4000])
    response = await llm_client.complete(
        prompt, response_format={"type": "json_object"},
        temperature=0,
    )
    data = json.loads(response)
    return DocumentMetadata(**data)

三、多路混合检索与融合排序

3.1 向量检索器

# rag/retrievers/vector.py - 向量检索(以Qdrant为例)
from qdrant_client import AsyncQdrantClient
from qdrant_client.models import (
    Distance, VectorParams, Filter, FieldCondition, MatchValue,
)
from ..system import RetrievedChunk, Document, SearchQuery, RetrievalStrategy

class VectorRetriever:
    def __init__(self, client: AsyncQdrantClient, collection: str,
                 embed_model, embed_dim: int = 1024):
        self.client = client
        self.collection = collection
        self.embed_model = embed_model
        self.embed_dim = embed_dim

    async def retrieve(self, query: SearchQuery, top_k: int = 20) -> list[RetrievedChunk]:
        # 查询嵌入
        query_vec = await self.embed_model.aembed(query.rewritten)
        # 元数据过滤
        must = []
        for key, value in query.filters.items():
            if isinstance(value, list):
                for v in value:
                    must.append(FieldCondition(key=f"meta.{key}", match=MatchValue(value=v)))
            else:
                must.append(FieldCondition(key=f"meta.{key}", match=MatchValue(value=value)))
        flt = Filter(must=must) if must else None

        results = await self.client.search(
            collection_name=self.collection,
            query_vector=query_vec,
            limit=top_k,
            query_filter=flt,
            with_payload=True,
        )
        return [
            RetrievedChunk(
                document=Document(
                    id=str(r.id),
                    content=r.payload["content"],
                    metadata=r.payload.get("meta", {}),
                    source=r.payload.get("source", ""),
                ),
                score=r.score,
                retrieval_strategy=RetrievalStrategy.VECTOR,
            )
            for r in results
        ]

3.2 全文检索与混合融合

# rag/retrievers/hybrid.py - BM25 + 向量融合(RRF算法)
import math
from collections import defaultdict
from ..system import RetrievedChunk, SearchQuery, RetrievalStrategy

class HybridRetriever:
    """组合BM25与向量检索,用Reciprocal Rank Fusion融合排序"""
    def __init__(self, bm25_retriever, vector_retriever, rrf_k: int = 60):
        self.bm25 = bm25_retriever
        self.vector = vector_retriever
        self.rrf_k = rrf_k

    async def retrieve(self, query: SearchQuery, top_k: int = 20) -> list[RetrievedChunk]:
        # 两路并行召回
        bm25_results, vec_results = await asyncio.gather(
            self.bm25.retrieve(query, top_k * 2),
            self.vector.retrieve(query, top_k * 2),
        )

        # RRF融合:score = sum(1 / (k + rank_in_each_list))
        rrf_scores: dict[str, float] = defaultdict(float)
        chunk_map: dict[str, RetrievedChunk] = {}

        for rank, chunk in enumerate(bm25_results):
            rrf_scores[chunk.document.id] += 1.0 / (self.rrf_k + rank + 1)
            chunk_map[chunk.document.id] = chunk
        for rank, chunk in enumerate(vec_results):
            rrf_scores[chunk.document.id] += 1.0 / (self.rrf_k + rank + 1)
            if chunk.document.id not in chunk_map:
                chunk_map[chunk.document.id] = chunk

        # 按RRF分数排序
        sorted_ids = sorted(rrf_scores, key=rrf_scores.get, reverse=True)
        return [
            RetrievedChunk(
                document=chunk_map[did].document,
                score=rrf_scores[did],
                retrieval_strategy=RetrievalStrategy.HYBRID,
            )
            for did in sorted_ids[:top_k]
        ]

import asyncio

四、查询重写与自查询

4.1 多跳查询分解

复杂问题需要分解为子查询分别检索再综合。以下实现基于LLM的查询规划器:

# rag/query_planner.py - 查询重写与多跳分解
from __future__ import annotations
import json
from ..system import SearchQuery, RetrievalStrategy

HYDE_PROMPT = """给定用户问题,生成一个假设性文档(Hypothetical Document),
这个文档如果存在,会直接回答该问题。用于检索增强。
仅输出文档内容,不要解释。

问题:{question}
假设性文档:"""

DECOMPOSE_PROMPT = """你是检索策略专家。分析用户问题,决定检索策略。

用户问题:{question}
对话历史(如有):{conversation}

输出JSON:
{{
  "rewritten": "重写后的检索查询(去除口语化,保留关键词,英文优先以提升向量召回)",
  "sub_queries": ["子查询1", "子查询2"],
  "filters": {{"meta_key": "value"}},
  "strategy": "hybrid|vector|fulltext|graph|sql",
  "needs_decomposition": true/false
}}

判断标准:
- needs_decomposition=true:问题包含多个独立子问题或需要多跳推理
- sub_queries:每个子查询应能独立检索回答
- filters:从问题中识别的元数据过滤条件(部门、文档类型等)
- strategy:graph适合关系型问题,sql适合数值统计,其余默认hybrid"""

class LLMQueryPlanner:
    def __init__(self, llm_client, embed_model):
        self.llm = llm_client
        self.embed = embed_model

    async def plan(self, user_query: str, conversation: list[dict]) -> SearchQuery:
        # HyDE:生成假设文档作为检索查询(提升语义召回)
        hyde_doc = await self.llm.complete(
            HYDE_PROMPT.format(question=user_query), temperature=0.3, max_tokens=200,
        )
        # 查询分解
        plan_raw = await self.llm.complete(
            DECOMPOSE_PROMPT.format(
                question=user_query,
                conversation=json.dumps(conversation[-4:], ensure_ascii=False),
            ),
            response_format={"type": "json_object"}, temperature=0,
        )
        plan = json.loads(plan_raw)

        return SearchQuery(
            original=user_query,
            rewritten=hyde_doc[:500] if plan.get("use_hyde", True) else plan["rewritten"],
            sub_queries=plan.get("sub_queries", []),
            filters=plan.get("filters", {}),
            strategy_hint=RetrievalStrategy(plan["strategy"]) if plan.get("strategy") else None,
        )

4.2 元数据自查询

让模型从自然语言中提取过滤条件,实现"用自然语言查结构化字段":

# rag/retrievers/self_query.py - 自查询过滤构建
SELF_QUERY_PROMPT = """从用户问题提取元数据过滤条件。

可用过滤字段:
- doc_type: policy|manual|api_doc|runbook|faq|spec
- department: engineering|product|design|ops|legal|finance
- product: 产品名(如Orders、Payments、Auth)
- effective_after: YYYY-MM-DD(生效日期下限)
- access_tag: public|internal|confidential

用户问题:{question}

输出JSON:{{"filters": {{...}}}},仅包含问题中明确提及的过滤条件。"""

async def build_self_query(question: str, llm_client) -> dict:
    raw = await llm_client.complete(
        SELF_QUERY_PROMPT.format(question=question),
        response_format={"type": "json_object"}, temperature=0,
    )
    return json.loads(raw).get("filters", {})

五、重排序与引用验证

5.1 交叉编码器重排

向量召回的Bi-Encoder把查询与文档独立编码,精度有限。Cross-Encoder把查询-文档拼接后联合编码,精度显著更高但计算量大,适合对召回的Top-20重排取Top-10:

# rag/reranker.py - 交叉编码器重排
from sentence_transformers import CrossEncoder
import asyncio

class CrossEncoderReranker:
    def __init__(self, model_name: str = "BAAI/bge-reranker-v2-m3",
                 device: str = "cuda", batch_size: int = 32):
        self.model = CrossEncoder(model_name, device=device)
        self.batch_size = batch_size

    async def rerank(self, query: str, chunks: list[RetrievedChunk]) -> list[RetrievedChunk]:
        if not chunks:
            return []
        pairs = [(query, c.document.content) for c in chunks]

        # 分批推理
        scores: list[float] = []
        for i in range(0, len(pairs), self.batch_size):
            batch = pairs[i:i + self.batch_size]
            batch_scores = self.model.predict(batch, show_progress_bar=False)
            scores.extend(batch_scores.tolist())

        for chunk, score in zip(chunks, scores):
            chunk.rerank_score = float(score)

        chunks.sort(key=lambda c: c.rerank_score, reverse=True)
        return chunks

class LLMReranker:
    """LLM-as-judge重排:用强模型评分,适合对Top候选精排"""
    def __init__(self, llm_client, model: str = "gpt-4o-mini"):
        self.llm = llm_client
        self.model = model

    async def rerank(self, query: str, chunks: list[RetrievedChunk]) -> list[RetrievedChunk]:
        if not chunks:
            return []
        # 让LLM对每个chunk打0-10的相关性分
        scoring_tasks = [
            self._score_chunk(query, c) for c in chunks
        ]
        scores = await asyncio.gather(*scoring_tasks)
        for chunk, score in zip(chunks, scores):
            chunk.rerank_score = score
        chunks.sort(key=lambda c: c.rerank_score, reverse=True)
        return chunks

    async def _score_chunk(self, query: str, chunk: RetrievedChunk) -> float:
        prompt = f"""对以下文档片段与查询的相关性打分(0.0-1.0):
查询:{query}
文档:{chunk.document.content[:500]}

仅输出一个浮点数。"""
        raw = await self.llm.complete(prompt, temperature=0, max_tokens=10)
        try:
            return max(0.0, min(1.0, float(raw.strip())))
        except ValueError:
            return 0.5

5.2 引用验证与幻觉检测

生成阶段必须标注引用并验证:每个事实陈述关联到具体检索块,无法关联的陈述标记为低置信:

# rag/generator.py - 带引用验证的生成
from __future__ import annotations
import re

CITATION_PROMPT = """基于以下检索到的文档片段回答用户问题。

要求:
1. 仅使用提供的文档内容,不要编造
2. 每个事实陈述后用[数字]标注来源,如"系统支持SSO[1]"
3. 如果文档不足以回答,明确说"根据现有资料无法确定"
4. 如有多份文档冲突,指出差异

检索文档:
{context}

用户问题:{question}

回答(附引用):"""

CONFIDENCE_PROMPT = """评估你对以下回答的置信度(0.0-1.0)。

回答:{answer}
引用来源数:{citation_count}
检索文档总数:{retrieved_count}
是否有"无法确定"声明:{has_uncertain}

输出JSON:{{"confidence": 0.0-1.0, "reason": "简短理由"}}"""

class CitingGenerator:
    def __init__(self, llm_client, model: str = "claude-sonnet-4-20250514"):
        self.llm = llm_client
        self.model = model

    async def generate(self, query: str, context: list[RetrievedChunk],
                       conversation: list[dict]) -> tuple[str, list[Citation], float]:
        # 构建带编号的上下文
        context_text = self._build_context(context)
        prompt = CITATION_PROMPT.format(
            context=context_text, question=query,
        )
        answer = await self.llm.complete(
            prompt, temperature=0.1, model=self.model,
        )

        # 解析引用标记
        citations = self._extract_citations(answer, context)
        # 置信度评估
        confidence = await self._assess_confidence(answer, citations, len(context))
        return answer, citations, confidence

    def _build_context(self, chunks: list[RetrievedChunk]) -> str:
        parts = []
        for i, chunk in enumerate(chunks, 1):
            parts.append(
                f"[{i}] 来源:{chunk.document.source}\n"
                f"内容:{chunk.document.content[:800]}\n"
            )
        return "\n---\n".join(parts)

    def _extract_citations(self, answer: str,
                           chunks: list[RetrievedChunk]) -> list[Citation]:
        # 匹配 [1] [2] 等引用标记
        refs = set(int(m) for m in re.findall(r"\[(\d+)\]", answer))
        citations = []
        for ref_num in sorted(refs):
            if 1 <= ref_num <= len(chunks):
                chunk = chunks[ref_num - 1]
                citations.append(Citation(
                    document_id=chunk.document.id,
                    source=chunk.document.source,
                    chunk_text=chunk.document.content[:200],
                    relevance=chunk.rerank_score or chunk.score,
                ))
        return citations

    async def _assess_confidence(self, answer: str,
                                 citations: list[Citation],
                                 retrieved_count: int) -> float:
        prompt = CONFIDENCE_PROMPT.format(
            answer=answer[:500],
            citation_count=len(citations),
            retrieved_count=retrieved_count,
            has_uncertain="无法确定" in answer or "不足以" in answer,
        )
        raw = await self.llm.complete(
            prompt, response_format={"type": "json_object"}, temperature=0,
        )
        import json
        try:
            data = json.loads(raw)
            return float(data["confidence"])
        except (json.JSONDecodeError, KeyError):
            # 回退:有引用且无"无法确定"则0.7,否则0.3
            return 0.7 if citations and "无法确定" not in answer else 0.3

    async def identify_knowledge_gap(self, answer: str,
                                      chunks: list[RetrievedChunk]):
        """低置信时识别知识缺口,指导下一轮检索"""
        prompt = f"""回答置信度低。分析缺失了什么信息。

回答:{answer}
已有文档主题:{[c.document.metadata.get('section', '') for c in chunks]}

输出JSON:
{{"refined_query": "改进的检索查询", "sub_queries": ["需要补充检索的子问题"]}}"""
        raw = await self.llm.complete(
            prompt, response_format={"type": "json_object"}, temperature=0,
        )
        import json
        return json.loads(raw)

    async def suggest_followups(self, answer: str,
                                 citations: list[Citation]) -> list[str]:
        prompt = f"""基于以下回答,生成3个用户可能想追问的问题。

回答:{answer[:500]}

输出JSON数组:["问题1", "问题2", "问题3"]"""
        raw = await self.llm.complete(
            prompt, response_format={"type": "json_object"}, temperature=0.3,
        )
        import json
        try:
            return json.loads(raw)
        except json.JSONDecodeError:
            return []

六、知识图谱融合检索

6.1 图谱构建与混合检索

知识图谱捕获实体关系,弥补向量检索在"多跳关系推理"上的短板。以企业组织架构与产品依赖为例:

# rag/graph/kg_retriever.py - 知识图谱检索
from neo4j import AsyncGraphDatabase

class KnowledgeGraphRetriever:
    def __init__(self, uri: str, user: str, password: str):
        self.driver = AsyncGraphDatabase.driver(uri, auth=(user, password))

    async def retrieve(self, query: SearchQuery, top_k: int = 10) -> list[RetrievedChunk]:
        # 1. 实体抽取(简化:用LLM或NER模型)
        entities = await self._extract_entities(query.rewritten)
        if not entities:
            return []

        # 2. 图谱遍历:1-2跳邻居
        cypher = """
        MATCH (e:Entity)
        WHERE e.name IN $entities
        CALL {
            WITH e
            MATCH (e)-[r*1..2]-(neighbor:Entity)
            RETURN neighbor, r
        }
        WITH neighbor, collect(r) as relations
        RETURN neighbor.name AS name,
               neighbor.type AS type,
               [rel IN relations | type(rel[0])] AS relation_types,
               neighbor.description AS description
        LIMIT $limit
        """
        async with self.driver.session() as session:
            result = await session.run(
                cypher, entities=entities, limit=top_k,
            )
            records = await result.data()

        # 3. 将图谱结果转为文档块
        chunks = []
        for record in records:
            content = (
                f"实体:{record['name']}{record['type']})\n"
                f"关系:{', '.join(record['relation_types'])}\n"
                f"描述:{record.get('description', '')}"
            )
            chunks.append(RetrievedChunk(
                document=Document(
                    id=f"kg:{record['name']}",
                    content=content,
                    metadata={"source": "knowledge_graph",
                              "entity": record["name"]},
                ),
                score=1.0,
                retrieval_strategy=RetrievalStrategy.GRAPH,
            ))
        return chunks

    async def _extract_entities(self, text: str) -> list[str]:
        # 简化:基于词典的实体匹配(生产用NER模型或LLM)
        known = await self._load_entity_names()
        found = [name for name in known if name in text]
        return found

    async def _load_entity_names(self) -> list[str]:
        async with self.driver.session() as session:
            result = await session.run("MATCH (e:Entity) RETURN e.name AS name")
            records = await result.data()
            return [r["name"] for r in records]

6.2 GraphRAG:社区摘要与全局检索

微软GraphRAG范式:把实体图聚类为社区,用LLM为每个社区生成摘要,全局问题检索社区摘要而非原始节点,显著降低上下文长度并改善全局问答:

# rag/graph/graph_rag.py - GraphRAG社区检索
class CommunitySummarizer:
    """离线管道:图聚类 -> 社区摘要 -> 向量索引"""
    async def build_communities(self, driver):
        # Leiden算法聚类
        async with driver.session() as session:
            await session.run("""
                CALL gds.leiden.write(
                    'entityGraph',
                    { writeProperty: 'community' }
                )
            """)
            # 每社区生成摘要
            result = await session.run("""
                MATCH (e:Entity)
                WITH e.community AS cid, collect(e) AS members
                WITH cid, members,
                     [m IN members | m.name + ': ' + coalesce(m.description, '')] AS texts
                RETURN cid, texts
            """)
            communities = await result.data()

        summaries = []
        for community in communities:
            text = "\n".join(community["texts"])
            summary = await self.llm.complete(
                f"总结以下实体群组的关键信息(200字内):\n{text[:3000]}",
                temperature=0.3,
            )
            summaries.append({
                "id": community["cid"],
                "summary": summary,
                "embedding": await self.embed_model.aembed(summary),
            })
        # 存入向量库
        await self.vector_store.upsert("communities", summaries)
        return summaries

七、评估体系与成本优化

7.1 RAG评估指标

RAG系统评估分检索与生成两端。检索端:召回率(Recall@K,相关文档是否被召回)、精确率(Precision@K,召回中有多少相关)、MRR(相关结果的平均排名)。生成端:忠实度(Faithfulness,答案是否仅基于检索内容,无幻觉)、答案相关性(Answer Relevance,是否回答了问题)、引用准确率(Citation Accuracy,引用标记是否正确对应)。以下实现一个评估框架:

# eval/rag_evaluator.py - RAG评估框架
from __future__ import annotations
from dataclasses import dataclass

@dataclass
class EvalCase:
    question: str
    expected_answer: str
    relevant_doc_ids: list[str]
    category: str  # factual|comparative|multi_hop|temporal|aggregation

class RAGEvaluator:
    def __init__(self, rag_system, llm_client):
        self.rag = rag_system
        self.llm = llm_client

    async def evaluate(self, cases: list[EvalCase]) -> dict:
        results = []
        for case in cases:
            response = await self.rag.answer(case.question)
            metrics = await self._compute_metrics(case, response)
            results.append({"question": case.question, **metrics})

        # 汇总
        import numpy as np
        summary = {
            "total": len(results),
            "recall_at_10": np.mean([r["recall"] for r in results]),
            "precision_at_10": np.mean([r["precision"] for r in results]),
            "faithfulness": np.mean([r["faithfulness"] for r in results]),
            "answer_relevance": np.mean([r["answer_relevance"] for r in results]),
            "citation_accuracy": np.mean([r["citation_accuracy"] for r in results]),
            "confidence_avg": np.mean([r["confidence"] for r in results]),
        }
        by_category = {}
        for r in results:
            cat = r.get("category", "unknown")
            by_category.setdefault(cat, []).append(r)
        summary["by_category"] = {
            cat: {k: np.mean([r[k] for r in rs])
                  for k in ["recall", "faithfulness", "answer_relevance"]}
            for cat, rs in by_category.items()
        }
        return summary

    async def _compute_metrics(self, case: EvalCase, response: RAGResponse) -> dict:
        # 检索端指标
        retrieved_ids = [c.document_id for c in response.citations]
        relevant_set = set(case.relevant_doc_ids)
        hits = len(set(retrieved_ids) & relevant_set)
        recall = hits / len(relevant_set) if relevant_set else 1.0
        precision = hits / len(retrieved_ids) if retrieved_ids else 0.0

        # 生成端指标(LLM-as-judge)
        faithfulness = await self._score_faithfulness(
            response.answer, [c.chunk_text for c in response.citations],
        )
        answer_relevance = await self._score_answer_relevance(
            case.question, response.answer,
        )
        citation_accuracy = await self._score_citation_accuracy(
            response.answer, response.citations,
        )
        return {
            "recall": recall, "precision": precision,
            "faithfulness": faithfulness,
            "answer_relevance": answer_relevance,
            "citation_accuracy": citation_accuracy,
            "confidence": response.confidence,
            "category": case.category,
        }

    async def _score_faithfulness(self, answer: str, sources: list[str]) -> float:
        prompt = f"""评估回答的忠实度(0.0-1.0)。
回答中的事实陈述是否都能从来源中找到支撑?
1.0=完全忠实,0.0=全部编造。

回答:{answer[:800]}
来源:{chr(10).join(s[:200] for s in sources[:3])}

仅输出浮点数。"""
        raw = await self.llm.complete(prompt, temperature=0, max_tokens=10)
        try:
            return float(raw.strip())
        except ValueError:
            return 0.5

    async def _score_answer_relevance(self, question: str, answer: str) -> float:
        prompt = f"""评估回答与问题的相关性(0.0-1.0)。
1.0=完美回答,0.0=完全偏题。

问题:{question}
回答:{answer[:500]}

仅输出浮点数。"""
        raw = await self.llm.complete(prompt, temperature=0, max_tokens=10)
        try:
            return float(raw.strip())
        except ValueError:
            return 0.5

    async def _score_citation_accuracy(self, answer: str,
                                        citations: list[Citation]) -> float:
        if not citations:
            return 0.0
        prompt = f"""检查回答中的引用标记[数字]是否正确指向了相应来源。

回答:{answer[:600]}
引用来源:
{chr(10).join(f'[{i+1}] {c.chunk_text[:150]}' for i, c in enumerate(citations))}

正确引用的比例(0.0-1.0):"""
        raw = await self.llm.complete(prompt, temperature=0, max_tokens=10)
        try:
            return float(raw.strip())
        except ValueError:
            return 0.5

7.2 成本优化策略

Agentic RAG的成本来自多轮LLM调用(查询规划、HyDE、重排、生成、置信度评估)。优化策略分四档:第一档,模型分级——规划与生成用强模型,HyDE与重排评分用快模型,嵌入用本地模型;第二档,缓存——查询嵌入与常见query的检索结果缓存(语义缓存:相似query复用),命中率30%可节省显著延迟与成本;第三档,早停——首轮检索若Top-1重排分>0.9直接生成跳过迭代;第四档,异步管道——文档摄入的元数据提取与嵌入离线批处理,在线仅做检索与生成。以下实现语义缓存:

# rag/semantic_cache.py - 语义缓存
import hashlib
import json
from dataclasses import dataclass

@dataclass
class CacheEntry:
    query: str
    query_embedding: list[float]
    response: str
    citations: list[dict]
    timestamp: float

class SemanticCache:
    def __init__(self, embed_model, vector_store, ttl: int = 3600,
                 similarity_threshold: float = 0.92):
        self.embed = embed_model
        self.store = vector_store
        self.ttl = ttl
        self.threshold = similarity_threshold

    async def get(self, query: str) -> CacheEntry | None:
        query_vec = await self.embed.aembed(query)
        results = await self.store.search(
            "semantic_cache", query_vec, limit=1,
        )
        if not results:
            return None
        hit = results[0]
        if hit.score < self.threshold:
            return None
        import time
        if time.time() - hit.payload["timestamp"] > self.ttl:
            return None
        return CacheEntry(
            query=hit.payload["query"],
            query_embedding=query_vec,
            response=hit.payload["response"],
            citations=hit.payload["citations"],
            timestamp=hit.payload["timestamp"],
        )

    async def put(self, query: str, response: str, citations: list[dict]):
        query_vec = await self.embed.aembed(query)
        import time
        await self.store.upsert(
            "semantic_cache",
            [{
                "id": hashlib.md5(query.encode()).hexdigest(),
                "vector": query_vec,
                "payload": {
                    "query": query, "response": response,
                    "citations": citations, "timestamp": time.time(),
                },
            }],
        )

八、部署架构与运维

8.1 生产部署拓扑

企业RAG系统典型部署:摄入管道(异步Worker批处理文档,写入向量库+图库+全文索引),在线服务(API网关→查询规划→多路检索→重排→生成),缓存层(Redis做会话缓存,Qdrant/Pgvector做语义缓存),可观测性(Langfuse/OpenTelemetry追踪每轮检索-生成链路)。容器化示例:

# docker-compose.yml 核心服务
version: "3.9"
services:
  qdrant:
    image: qdrant/qdrant:v1.12.0
    volumes: ["qdrant-data:/qdrant/storage"]
    ports: ["6333:6333"]

  neo4j:
    image: neo4j:5.25-community
    environment:
      NEO4J_AUTH: neo4j/password
    volumes: ["neo4j-data:/data"]
    ports: ["7474:7474", "7687:7687"]

  elasticsearch:
    image: docker.elastic.co/elasticsearch/elasticsearch:8.15.0
    environment:
      discovery.type: single-node
      xpack.security.enabled: "false"
    volumes: ["es-data:/usr/share/elasticsearch/data"]
    ports: ["9200:9200"]

  rag-api:
    build: ./rag-service
    environment:
      QDRANT_URL: http://qdrant:6333
      NEO4J_URL: bolt://neo4j:7687
      ES_URL: http://elasticsearch:9200
      LLM_API_KEY: ${LLM_API_KEY}
    ports: ["8000:8000"]
    depends_on: [qdrant, neo4j, elasticsearch]

  ingest-worker:
    build: ./rag-service
    command: ["python", "-m", "rag.ingest_worker"]
    environment:
      QDRANT_URL: http://qdrant:6333
      NEO4J_URL: bolt://neo4j:7687
    depends_on: [qdrant, neo4j]

volumes:
  qdrant-data:
  neo4j-data:
  es-data:

总结

RAG 2.0的本质是把检索从一次性操作升级为Agent驱动的迭代探索:查询规划器通过HyDE与多跳分解将用户意图转化为可检索的形式,多路混合检索(向量+BM25+图+SQL)以互补方式覆盖不同信息需求,RRF融合与交叉编码器重排分层提升精度,引用验证与置信度评估将幻觉抑制从"希望"变为"可度量",知识图谱融合弥补向量检索的关系推理短板,Agentic决策循环让系统在低置信时自主再检索直至满足阈值。工程落地的核心权衡在质量与成本之间:强模型+多轮迭代+图检索带来最高质量但延迟与token消耗显著,生产部署需要模型分级、语义缓存、早停策略与异步摄入管道的组合优化。评估必须闭环:召回率、忠实度、引用准确率、分品类表现构成多维体检,任何改动都应通过评估集回归验证。当RAG系统从问答扩展为企业知识中枢——连接文档库、代码库、数据库、工单系统、设计稿,用统一的检索-生成范式服务所有知识消费场景——其工程深度直接决定了组织"把信息转化为决策"的效率上限,这正是Agentic RAG在AI落地浪潮中的真正价值坐标。

【声明】本内容来自华为云开发者社区博主,不代表华为云及华为云开发者社区的观点和立场。转载时必须标注文章的来源(华为云社区)、文章链接、文章作者等基本信息,否则作者和本社区有权追究责任。如果您发现本社区中有涉嫌抄袭的内容,欢迎发送邮件进行举报,并提供相关证据,一经查实,本社区将立刻删除涉嫌侵权内容,举报邮箱: cloudbbs@huaweicloud.com
  • 点赞
  • 收藏
  • 关注作者

评论(0

0/1000
抱歉,系统识别当前为高风险访问,暂不支持该操作

全部回复

上滑加载中

设置昵称

在此一键设置昵称,即可参与社区互动!

*长度不超过10个汉字或20个英文字符,设置后3个月内不可修改。

*长度不超过10个汉字或20个英文字符,设置后3个月内不可修改。