0Pricing
AI Engineering Academy · 课时

实现核心 RAG 与智能体功能

按照整个学习路径中掌握的模式,构建文档摄取流程、向量存储索引、带重排序的检索以及智能体工具集成。

实现核心 RAG 与智能体功能 是 CoddyKit 上的免费 AI Engineering Academy 课时。 这是第 2 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 AI Engineering Academy 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 AI Engineering Academy 课程共包含 4 节课。

实现顺序与依赖关系

请自底向上构建系统:先从没有外部依赖的组件开始,再逐层添加依赖于它们的组件。对于 RAG + 智能体系统,顺序为:(1) 向量存储架构,(2) 摄取流水线,(3) 检索器,(4) 基础问答链,(5) 流式端点,(6) 带工具的智能体,(7) 缓存层,(8) 追踪插桩。在将每个组件集成到流水线之前,先单独测试它。

# Build order:
IMPL_ORDER = [
    'pgvector_schema',      # prerequisite for everything
    'document_ingestion',   # populate the vector store
    'hybrid_retriever',     # test retrieval in isolation
    'qa_chain_basic',       # integrate LLM with retrieval
    'streaming_endpoint',   # expose via API
    'function_calling',     # add agent tool calls
    'semantic_cache',       # reduce repeat API calls
    'langsmith_tracing',    # add after core works
    'injection_filter',     # harden before load testing
]

设置向量存储

在摄取文档之前,先创建具有正确架构的 pgvector 表。请包含向量、分块文本、所有元数据字段,以及用于选择性重新索引的 updated_at 时间戳。立即在嵌入列上创建向量索引(HNSW 或 IVFFlat)——在数百万行数据之后再添加索引,会比在空表上提前添加慢得多。

-- PostgreSQL schema with pgvector
CREATE EXTENSION IF NOT EXISTS vector;

CREATE TABLE document_chunks (
    id          BIGSERIAL PRIMARY KEY,
    doc_id      TEXT NOT NULL,
    chunk_text  TEXT NOT NULL,
    embedding   VECTOR(1536) NOT NULL,
    source_file TEXT,
    page_number INT,
    section     TEXT,
    doc_type    TEXT,
    tenant_id   TEXT NOT NULL,  -- for data isolation
    created_at  TIMESTAMPTZ DEFAULT NOW()
);

CREATE INDEX idx_chunks_hnsw ON document_chunks
USING hnsw (embedding vector_cosine_ops)
WITH (m = 16, ef_construction = 64);

CREATE INDEX idx_chunks_tenant ON document_chunks(tenant_id);

构建摄取流水线

将摄取实现为一个异步函数,接受文件路径或 URL,并返回已建立索引的分块数量。针对不同文件类型使用 LangChain 的文档加载器,并使用递归字符文本分割器进行分块。对嵌入调用进行批处理,以遵守 2048 个令牌的输入限制,并将 API 调用次数从数千次减少到几十次。

from langchain_community.document_loaders import PyPDFLoader
from langchain.text_splitter import RecursiveCharacterTextSplitter
from openai import AsyncOpenAI

async def ingest_document(file_path: str, doc_id: str, tenant_id: str) -> int:
    loader = PyPDFLoader(file_path)
    pages = loader.load()
    splitter = RecursiveCharacterTextSplitter(chunk_size=800, chunk_overlap=100)
    chunks = splitter.split_documents(pages)
    # Batch embed
    texts = [c.page_content for c in chunks]
    client = AsyncOpenAI()
    embeddings_response = await client.embeddings.create(
        model='text-embedding-3-small',
        input=texts
    )
    embeddings = [e.embedding for e in embeddings_response.data]
    await insert_chunks_to_pgvector(chunks, embeddings, doc_id, tenant_id)
    return len(chunks)

构建混合检索器

将稠密向量搜索与 BM25 关键词搜索结合起来,并使用倒数排名融合合并结果。将检索器实现为一个类,其中只包含一个 retrieve(query, tenant_id, top_k) 方法。在该方法中,使用 asyncio.gather 并发运行两种搜索,使用 RRF 合并排序列表,按照分块 ID 去重,并返回带有来源元数据的前 top_k 个结果。

import asyncio
from rank_bm25 import BM25Okapi

class HybridRetriever:
    def __init__(self, pool, k_rrf: int = 60):
        self.pool = pool
        self.k_rrf = k_rrf

    async def retrieve(self, query: str, tenant_id: str, top_k: int = 10) -> list:
        dense_results, sparse_results = await asyncio.gather(
            self._dense_search(query, tenant_id, top_k * 3),
            self._bm25_search(query, tenant_id, top_k * 3)
        )
        merged = self._rrf_merge(dense_results, sparse_results)
        return merged[:top_k]

    def _rrf_score(self, rank: int) -> float:
        return 1.0 / (self.k_rrf + rank + 1)

实现核心问答链

使用 LangChain LCEL 构建核心问答链。该链接受问题和检索到的分块,根据“只能使用所提供上下文”的指令格式化增强提示,并以流式方式返回响应。向模型添加明确指令,要求它通过文档名称和页码引用来源;如果答案不在检索到的上下文中,则回答“我不知道”。

from langchain.prompts import ChatPromptTemplate
from langchain_openai import ChatOpenAI
from langchain.schema.output_parser import StrOutputParser

RAG_PROMPT = ChatPromptTemplate.from_messages([
    ('system', 'You are a precise assistant. Answer ONLY using the provided context. '
               'Cite sources as [DocName, p.N]. If the answer is not in the context, say "I don\'t have information about that."'),
    ('user', 'Context:\n{context}\n\nQuestion: {question}')
])

llm = ChatOpenAI(model='gpt-4o', temperature=0, streaming=True)
qa_chain = RAG_PROMPT | llm | StrOutputParser()

async def answer_question(question: str, chunks: list) -> str:
    context = '\n\n'.join(f'[{c["source"]}]\n{c["text"]}' for c in chunks)
    return await qa_chain.ainvoke({'context': context, 'question': question})

为智能体添加函数调用

为基础问答系统扩展函数调用,以处理需要实时数据或计算的查询。为以下功能定义工具:搜索网络获取最新信息、针对安全的只读数据库运行 SQL 查询,以及通过 ID 查询特定记录。智能体根据问题决定调用哪些工具,执行这些工具,并将结果整合到最终答案中。

from openai import AsyncOpenAI
import json

TOOLS = [
    {
        'type': 'function',
        'function': {
            'name': 'search_knowledge_base',
            'description': 'Search the internal document knowledge base for relevant information',
            'parameters': {
                'type': 'object',
                'properties': {
                    'query': {'type': 'string', 'description': 'Search query'},
                    'top_k': {'type': 'integer', 'default': 5}
                },
                'required': ['query']
            }
        }
    }
]

async def agent_with_tools(question: str, tenant_id: str) -> str:
    client = AsyncOpenAI()
    messages = [{'role': 'user', 'content': question}]
    while True:
        resp = await client.chat.completions.create(
            model='gpt-4o', messages=messages, tools=TOOLS)
        if resp.choices[0].finish_reason != 'tool_calls':
            return resp.choices[0].message.content
        tool_call = resp.choices[0].message.tool_calls[0]
        args = json.loads(tool_call.function.arguments)
        result = await dispatch_tool(tool_call.function.name, args, tenant_id)
        messages.append({'role': 'tool', 'tool_call_id': tool_call.id, 'content': result})

流式传输智能体响应

使用服务器发送事件将代理封装在 FastAPI 的 StreamingResponse 中,以便前端在令牌到达时立即显示。对于包含工具调用的代理响应,在工具执行期间流式传输进度指示器(“正在搜索知识库……”),然后逐个令牌地流式传输最终答案。这样可以避免工具执行延迟导致页面看起来像被冻结一样。

from fastapi import FastAPI
from fastapi.responses import StreamingResponse
import asyncio

app = FastAPI()

async def event_stream(question: str, tenant_id: str):
    yield 'data: {"type": "start"}\n\n'
    chunks = await retriever.retrieve(question, tenant_id)
    yield 'data: {"type": "retrieving", "count": ' + str(len(chunks)) + '}\n\n'
    async for token in qa_chain.astream({'context': format_context(chunks), 'question': question}):
        yield f'data: {{"type": "token", "content": {repr(token)}}}\n\n'
    yield 'data: {"type": "done"}\n\n'

@app.post('/query/stream')
async def stream_query(question: str, tenant_id: str):
    return StreamingResponse(event_stream(question, tenant_id), media_type='text/event-stream')

接入语义缓存

将语义缓存作为查询流水线中的检索前步骤。在访问检索器和 LLM 之前,先对用户问题进行嵌入并检查缓存。如果缓存命中且相似度高于阈值,则立即返回缓存的答案,并附带 cached: true 标志。缓存未命中时,继续执行完整流水线,并将得到的答案存入缓存,以供未来的相似查询使用。

async def query_pipeline(question: str, tenant_id: str) -> dict:
    # 1. Check semantic cache
    cache_hit = await semantic_cache.lookup(question, tenant_id, threshold=0.92)
    if cache_hit:
        return {'answer': cache_hit.answer, 'cached': True, 'sources': cache_hit.sources}

    # 2. Retrieve
    chunks = await retriever.retrieve(question, tenant_id, top_k=5)
    chunks = await reranker.rerank(question, chunks, top_n=3)

    # 3. Generate
    answer = await answer_question(question, chunks)
    sources = [c['source'] for c in chunks]

    # 4. Cache result
    await semantic_cache.store(question, tenant_id, answer, sources)

    return {'answer': answer, 'cached': False, 'sources': sources}

添加 LangSmith 跟踪

通过设置两个环境变量,使用 LangSmith 为流水线添加监测。每次 LangChain 链调用都会自动记录跟踪信息,包括令牌数量、延迟、输入、输出以及所有错误。对于自定义的非 LangChain 代码(检索、重排序),请使用 @traceable 装饰器封装调用,将其纳入跟踪范围。这样便可完整了解流水线端到端的每个步骤。

import os
from langsmith import traceable

os.environ['LANGCHAIN_TRACING_V2'] = 'true'
os.environ['LANGCHAIN_API_KEY'] = os.environ['LANGSMITH_API_KEY']
os.environ['LANGCHAIN_PROJECT'] = 'document-qa-production'

# Wrap non-LangChain steps with @traceable
@traceable(name='hybrid_retrieval')
async def traced_retrieval(question: str, tenant_id: str, top_k: int) -> list:
    return await retriever.retrieve(question, tenant_id, top_k)

@traceable(name='cohere_reranking')
async def traced_reranking(question: str, chunks: list) -> list:
    return await reranker.rerank(question, chunks)

# LangChain LCEL chains are automatically traced — no extra code needed

运行集成测试

编写集成测试,从问题输入到最终答案,完整执行查询流水线。使用包含已知文档的小型测试向量存储,以便对应该检索哪些文本块以及答案应包含什么内容编写确定性的断言。在与生产基础设施相似、但使用独立向量存储和 LLM 密钥的预发布环境中运行集成测试。

import pytest

@pytest.mark.asyncio
async def test_full_pipeline_returns_grounded_answer():
    # Setup: ingest known document
    await ingest_document('tests/fixtures/policy.pdf', 'policy_v1', 'test_tenant')

    # Query with a question that has a known answer in the document
    result = await query_pipeline(
        question='What is the cancellation policy?',
        tenant_id='test_tenant'
    )

    assert result['answer'] is not None
    assert len(result['answer']) > 50
    assert '24 hours' in result['answer']  # known fact in document
    assert 'policy.pdf' in str(result['sources'])
    assert result['cached'] is False  # fresh query

衡量基准检索质量

在优化任何内容之前,先建立检索基线。使用评估测试集衡量命中率(正确文档是否出现在前 5 个结果中?)和 MRR(排名有多高?)。在设置初始向量存储后、添加重排序或混合搜索之前,运行此基线。基线可以告诉您每项优化实际带来了哪些改进,从而为值得保留的技术提供依据。

async def measure_retrieval_baseline(test_cases: list) -> dict:
    hits = 0
    reciprocal_ranks = []
    for case in test_cases:
        results = await retriever.retrieve(case['question'], case['tenant_id'], top_k=5)
        result_docs = [r['doc_id'] for r in results]
        if case['relevant_doc'] in result_docs:
            hits += 1
            rank = result_docs.index(case['relevant_doc']) + 1
            reciprocal_ranks.append(1.0 / rank)
        else:
            reciprocal_ranks.append(0.0)
    return {
        'hit_rate_at_5': hits / len(test_cases),
        'mrr': sum(reciprocal_ranks) / len(reciprocal_ranks)
    }

快速检查

测试您对构建核心 RAG 和代理功能的理解。

课程回顾

在本课中,您学到了:自底向上的实现顺序可确保每个组件在集成前都经过单独测试;通过 RRF 合并并发的稠密搜索与 BM25 搜索组成的混合检索,可获得最佳检索质量;而使用 @traceable 装饰器进行 LangSmith 跟踪,则能完整了解流水线。接下来,我们将通过安全性、缓存和可靠性模式来增强系统。

常见问题解答

「实现核心 RAG 与智能体功能」课时是免费的吗?

是的 — 「实现核心 RAG 与智能体功能」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 AI Engineering Academy 课程的其余内容,请升级到 CoddyKit PRO。 AI Engineering Academy 课程共包含 4 节课。

「实现核心 RAG 与智能体功能」这节课中我会学到什么?

按照整个学习路径中掌握的模式,构建文档摄取流程、向量存储索引、带重排序的检索以及智能体工具集成。 你通过在浏览器中直接运行的动手代码来练习 AI Engineering Academy,全天候 AI 导师会在你学习这节课的过程中回答你的问题。

学习 AI Engineering Academy 需要有经验吗?

无需任何先前经验。CoddyKit 上的 AI Engineering Academy 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 2 节课,共 4 节。

「实现核心 RAG 与智能体功能」课时需要多长时间?

大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。

我能在这节 AI Engineering Academy 课中编写并运行代码吗?

能。每节 AI Engineering Academy 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。

此课程中的所有课时

  1. 设计生产环境架构
  2. 实现核心 RAG 与智能体功能
  3. 加固:安全性、缓存与可靠性
  4. 评估、部署与复盘
← 返回 AI Engineering Academy