Implementando recursos essenciais de RAG e agentes
Crie o pipeline de ingestão de documentos, a indexação do armazenamento vetorial, a recuperação com reclassificação e as integrações das ferramentas do agente, seguindo os padrões aprendidos ao longo da trilha.
Implementando recursos essenciais de RAG e agentes é uma aula grátis de AI Engineering Academy no CoddyKit. Esta é a aula 2 de 4. Você pode ler a aula completa abaixo gratuitamente — depois pratica ao vivo no navegador com um editor de código integrado e um tutor de IA 24/7. Faz parte do caminho de aprendizado de AI Engineering Academy, e seu progresso é sincronizado entre a web e o app CoddyKit. O curso de AI Engineering Academy inclui 4 aulas no total.
Ordem de implementação e dependências
Construa o sistema de baixo para cima: comece pelos componentes que não têm dependências externas e, depois, adicione camadas com os componentes que dependem deles. Para um sistema de RAG + agente, a ordem é: (1) esquema do armazenamento vetorial, (2) pipeline de ingestão, (3) recuperador, (4) cadeia básica de perguntas e respostas, (5) endpoint de transmissão, (6) agente com ferramentas, (7) camada de cache e (8) instrumentação de rastreamento. Teste cada componente isoladamente antes de integrá-lo à pipeline.
# 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
]Configurando o armazenamento vetorial
Crie a tabela do pgvector com o esquema correto antes de ingerir os documentos. Inclua colunas para o vetor, o texto da parte, todos os campos de metadados e um carimbo de data e hora updated_at para a reindexação seletiva. Crie imediatamente um índice vetorial (HNSW ou IVFFlat) na coluna de embeddings — adicionar o índice depois de milhões de linhas é muito mais lento do que adicioná-lo antecipadamente em uma tabela vazia.
-- 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);Construindo a pipeline de ingestão
Implemente a ingestão como uma única função assíncrona que aceite um caminho de arquivo ou uma URL e retorne o número de partes indexadas. Use os carregadores de documentos do LangChain para diferentes tipos de arquivo e o divisor recursivo de texto por caracteres para dividir os documentos. Faça chamadas de embeddings em lotes para permanecer dentro do limite de entrada de 2.048 tokens e reduzir as chamadas à API de milhares para dezenas.
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)Construindo o recuperador híbrido
Combine a busca vetorial densa com a busca por palavras-chave BM25 e mescle os resultados usando a fusão recíproca de posições. Implemente o recuperador como uma classe com um único método retrieve(query, tenant_id, top_k). Dentro dele, execute as duas buscas simultaneamente com asyncio.gather, mescle as listas classificadas com RRF, elimine duplicatas pelo ID da parte e retorne os resultados principais com seus metadados de origem.
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)Implementando a cadeia principal de perguntas e respostas
Construa a cadeia principal de perguntas e respostas usando LangChain LCEL. A cadeia recebe uma pergunta e as partes recuperadas, formata um prompt aumentado com instruções para usar apenas o contexto fornecido e transmite a resposta. Adicione instruções explícitas para que o modelo cite as fontes pelo nome do documento e pelo número da página e diga “Não sei” quando a resposta não estiver no contexto recuperado.
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})Adicionando chamadas de funções ao agente
Amplie o sistema básico de perguntas e respostas com chamadas de funções para lidar com consultas que exigem dados ou cálculos em tempo real. Defina ferramentas para: pesquisar informações atuais na web, executar consultas SQL em um banco de dados seguro somente para leitura e consultar registros específicos por ID. O agente decide quais ferramentas chamar com base na pergunta, executa-as e incorpora os resultados à resposta final.
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})Transmitindo a resposta do agente
Envolva o agente em uma StreamingResponse do FastAPI usando eventos enviados pelo servidor, para que o frontend exiba os tokens à medida que chegam. Para respostas do agente com chamadas de ferramentas, transmita um indicador de progresso enquanto a ferramenta é executada ('Pesquisando na base de conhecimento...') e, em seguida, transmita a resposta final token por token. Isso evita que a página pareça congelada durante a latência da execução da ferramenta.
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')Integração do Semantic Cache
Adicione o cache semântico como uma etapa anterior à recuperação no fluxo de consulta. Antes de acessar o recuperador e o LLM, gere a representação vetorial da pergunta do usuário e verifique o cache. Quando houver um acerto com similaridade acima do limite, retorne imediatamente a resposta armazenada, com uma sinalização cached: true. As consultas sem acerto seguem pelo fluxo completo, e a resposta resultante é armazenada no cache para futuras consultas semelhantes.
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}Adição do rastreamento com LangSmith
Instrumente o fluxo com LangSmith definindo duas variáveis de ambiente. Cada chamada de cadeia do LangChain é rastreada automaticamente com contagens de tokens, latência, entradas, saídas e quaisquer erros. Para código personalizado que não usa LangChain (recuperação e reordenação), envolva as chamadas em decoradores @traceable para incluí-las no rastreamento. Isso proporciona visibilidade completa, de ponta a ponta, de cada etapa do fluxo.
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 neededExecução de testes de integração
Escreva testes de integração que exercitem o fluxo completo de consulta, desde a entrada da pergunta até a resposta final. Use um pequeno armazenamento vetorial de teste com documentos conhecidos, para que possa escrever verificações determinísticas sobre quais trechos devem ser recuperados e o que a resposta deve conter. Execute os testes de integração em um ambiente de homologação que espelhe a infraestrutura de produção, mas use um armazenamento vetorial e uma chave de LLM separados.
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 queryMedição da qualidade de recuperação de referência
Antes de otimizar qualquer coisa, estabeleça uma referência de recuperação. Use seu conjunto de dados de avaliação para medir a taxa de acerto (o documento correto aparece nos 5 primeiros resultados?) e o MRR (em que posição ele aparece?). Execute essa medição depois de configurar o armazenamento vetorial inicial e antes de adicionar reordenação ou pesquisa híbrida. A referência mostra quais melhorias cada otimização realmente proporciona, fornecendo evidências sobre quais técnicas vale a pena manter.
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)
}Verificação rápida
Teste sua compreensão sobre a criação dos principais recursos de RAG e agentes.
Recapitulação da lição
Nesta lição, você aprendeu que a ordem de implementação de baixo para cima garante que cada componente seja testado isoladamente antes da integração; que a recuperação híbrida com pesquisa densa e BM25 concorrentes, combinada por meio de RRF, oferece a melhor qualidade de recuperação; e que o rastreamento do LangSmith com decoradores @traceable proporciona visibilidade completa do fluxo. A seguir, reforçaremos o sistema com padrões de segurança, cache e confiabilidade.
Perguntas Frequentes
A aula “Implementando recursos essenciais de RAG e agentes” é grátis?
Sim — o texto completo de “Implementando recursos essenciais de RAG e agentes” é grátis para ler aqui na web. Para praticá-la interativamente (um editor de código integrado e um tutor de IA 24/7) e desbloquear o restante do curso de AI Engineering Academy, atualize para CoddyKit PRO. O curso de AI Engineering Academy inclui 4 aulas no total.
O que vou aprender em “Implementando recursos essenciais de RAG e agentes”?
Crie o pipeline de ingestão de documentos, a indexação do armazenamento vetorial, a recuperação com reclassificação e as integrações das ferramentas do agente, seguindo os padrões aprendidos ao longo… Você pratica AI Engineering Academy com código prático que executa diretamente no navegador, e um tutor de IA 24/7 responde suas dúvidas enquanto trabalha na aula.
Preciso ter experiência prévia para começar AI Engineering Academy?
Nenhuma experiência prévia é necessária. AI Engineering Academy no CoddyKit é estruturado para alunos iniciantes até avançados, então você pode começar aqui ou desde o início e aprender no seu ritmo. Esta é a aula 2 de 4.
Quanto tempo leva a aula “Implementando recursos essenciais de RAG e agentes”?
A maioria das aulas CoddyKit leva cerca de 5–10 minutos. Cada uma é compacta e interativa, então você faz progresso constante e retoma exatamente de onde parou entre web e app.
Posso escrever e executar código nesta aula de AI Engineering Academy?
Sim. Cada aula de AI Engineering Academy inclui um editor de código integrado, então você escreve e executa código real direto no navegador e recebe feedback de IA instantaneamente — nenhuma configuração local necessária.
Todas as aulas deste curso
- Projetando a arquitetura de produção
- Implementando recursos essenciais de RAG e agentes
- Reforço: segurança, armazenamento em cache e confiabilidade
- Avaliação, implantação e retrospectiva