Implementazione delle funzionalità fondamentali di RAG e degli agenti
Costruisca la pipeline di acquisizione dei documenti, l’indicizzazione nel vector store, il recupero con riordinamento e le integrazioni degli strumenti degli agenti, seguendo i pattern appresi durante il percorso.
Implementazione delle funzionalità fondamentali di RAG e degli agenti è una lezione AI Engineering Academy gratuita su CoddyKit. Questa è la lezione 2 di 4. Puoi leggere la lezione completa qui gratuitamente — poi esercitati direttamente nel browser con un editor di codice integrato e un tutor IA disponibile 24/7. Fa parte del percorso di apprendimento AI Engineering Academy, e i tuoi progressi si sincronizzano tra il web e l'app CoddyKit. Il corso AI Engineering Academy include 4 lezioni in totale.
Ordine di implementazione e dipendenze
Costruisca il sistema dal basso verso l'alto: inizi dai componenti che non hanno dipendenze esterne, quindi aggiunga i componenti che dipendono da essi. Per un sistema RAG + agent, l'ordine è: (1) schema del vector store, (2) pipeline di ingestione, (3) retriever, (4) catena Q&A di base, (5) endpoint di streaming, (6) agent con strumenti, (7) livello di caching, (8) strumentazione di tracing. Testi ogni componente in isolamento prima di integrarlo nella 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
]Configurare il vector store
Crei la tabella pgvector con lo schema corretto prima di acquisire i documenti. Includa colonne per il vettore, il testo del chunk, tutti i campi di metadati e un timestamp updated_at per la reindicizzazione selettiva. Crei immediatamente un indice vettoriale (HNSW o IVFFlat) sulla colonna degli embedding: aggiungere l'indice dopo milioni di righe è molto più lento che aggiungerlo inizialmente su una tabella vuota.
-- 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);Creare la pipeline di ingestione
Implementi l'ingestione come un'unica funzione async che accetta un percorso di file o un URL e restituisce il numero di chunk indicizzati. Utilizzi i document loader di LangChain per i diversi tipi di file e il recursive character text splitter per la suddivisione in chunk. Raggruppi le chiamate di embedding per rispettare il limite di input di 2048 token e ridurre le chiamate API da migliaia a poche decine.
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)Creare il retriever ibrido
Combini la ricerca vettoriale densa con la ricerca per parole chiave BM25 e unisca i risultati tramite reciprocal rank fusion. Implementi il retriever come una classe con un unico metodo retrieve(query, tenant_id, top_k). Al suo interno, esegua entrambe le ricerche in concorrenza con asyncio.gather, unisca le liste ordinate con RRF, elimini i duplicati in base all'ID del chunk e restituisca i risultati top_k con i relativi metadati della fonte.
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)Implementare la catena QA principale
Costruisca la catena Q&A principale utilizzando LangChain LCEL. La catena riceve una domanda e i chunk recuperati, formatta un prompt arricchito con istruzioni per usare esclusivamente il contesto fornito e trasmette la risposta in streaming. Aggiunga istruzioni esplicite affinché il modello citi le fonti indicando il nome del documento e il numero di pagina e dica «Non lo so» quando la risposta non è presente nel contesto recuperato.
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})Aggiungere il function calling all'agente
Estenda il sistema Q&A di base con il function calling per gestire le query che richiedono dati o calcoli in tempo reale. Definisca strumenti per: cercare informazioni aggiornate sul web, eseguire query SQL su un database sicuro di sola lettura e recuperare record specifici tramite ID. L'agente decide quali strumenti chiamare in base alla domanda, li esegue e integra i risultati nella risposta finale.
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})Trasmettere in streaming la risposta dell'agente
Avvolga l’agente in una StreamingResponse di FastAPI usando gli eventi inviati dal server, in modo che il frontend visualizzi i token man mano che arrivano. Per le risposte dell’agente con chiamate a strumenti, trasmetta un indicatore di avanzamento mentre lo strumento è in esecuzione («Searching knowledge base...»), quindi trasmetta la risposta finale token per token. In questo modo la pagina non sembra bloccata durante la latenza dovuta all’esecuzione dello strumento.
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')Collegamento della cache semantica
Aggiunga la cache semantica come passaggio di pre-recupero nella pipeline delle query. Prima di chiamare il retriever e l’LLM, trasformi la domanda dell’utente in un embedding e verifichi la cache. In caso di riscontro nella cache con una similarità superiore alla soglia, restituisca immediatamente la risposta memorizzata con il flag cached: true. In caso di mancata corrispondenza, la query prosegue attraverso l’intera pipeline e la risposta risultante viene memorizzata nella cache per query simili future.
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}Aggiunta del tracing con LangSmith
Strumenti la pipeline con LangSmith impostando due variabili d’ambiente. Ogni chiamata a una chain di LangChain viene tracciata automaticamente con il conteggio dei token, la latenza, gli input, gli output ed eventuali errori. Per il codice personalizzato non basato su LangChain, come il recupero e il reranking, racchiuda le chiamate in decorator @traceable per includerle nel trace. In questo modo otterrà una visibilità completa end-to-end su ogni passaggio della pipeline.
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 neededEsecuzione dei test di integrazione
Scriva test di integrazione che esercitino l’intera pipeline delle query, dall’input della domanda fino alla risposta finale. Utilizzi un piccolo vector store di test con documenti noti, così da poter definire asserzioni deterministiche sui chunk che dovrebbero essere recuperati e sul contenuto atteso della risposta. Esegua i test di integrazione in un ambiente di staging che rispecchi l’infrastruttura di produzione, ma utilizzi un vector store e una chiave LLM separati.
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 queryMisurazione della qualità di recupero di riferimento
Prima di ottimizzare qualsiasi elemento, stabilisca una baseline del recupero. Utilizzi il test set di valutazione per misurare il tasso di hit (il documento corretto compare tra i primi 5 risultati?) e l’MRR (quanto in alto si posiziona?). Esegua questa misurazione dopo aver configurato il vector store iniziale, prima di aggiungere il reranking o la ricerca ibrida. La baseline mostra quali miglioramenti apporta effettivamente ogni ottimizzazione, fornendo dati concreti sulle tecniche che vale la pena mantenere.
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 rapida
Verifichi la propria comprensione della creazione delle funzionalità fondamentali di RAG e degli agenti.
Riepilogo della lezione
In questa lezione ha appreso che l’ordine di implementazione dal basso verso l’alto garantisce che ogni componente venga testato in isolamento prima dell’integrazione, che il recupero ibrido con ricerca dense e BM25 concorrenti, combinato tramite RRF, offre la migliore qualità di recupero e che il tracing di LangSmith con decorator @traceable fornisce una visibilità completa sulla pipeline. Ora renderemo il sistema più robusto con pattern di sicurezza, caching e affidabilità.
Impara Python con un tutor IA — gratis
Scrivi ed esegui vero codice nel tuo browser, ricevi aiuto istantaneo da un tutor IA disponibile 24/7, e riprendi da dove hai lasciato sul web o nell'app.
- Corsi
- 30
- Lezioni
- 120
Domande Frequenti
La lezione «Implementazione delle funzionalità fondamentali di RAG e degli agenti» è gratuita?
Sì — il testo completo di «Implementazione delle funzionalità fondamentali di RAG e degli agenti» è gratuito qui sul web. Per esercitarvi in modo interattivo (un editor di codice integrato e un tutor IA 24/7) e sbloccare il resto del corso AI Engineering Academy, passa a CoddyKit PRO. Il corso AI Engineering Academy include 4 lezioni in totale.
Cosa imparerò in «Implementazione delle funzionalità fondamentali di RAG e degli agenti»?
Costruisca la pipeline di acquisizione dei documenti, l’indicizzazione nel vector store, il recupero con riordinamento e le integrazioni degli strumenti degli agenti, seguendo i pattern appresi duran… Eserciti AI Engineering Academy con codice pratico che esegui direttamente nel browser, e un tutor IA 24/7 risponde alle tue domande mentre lavori sulla lezione.
Ho bisogno di esperienza per iniziare AI Engineering Academy?
Non è richiesta alcuna esperienza precedente. AI Engineering Academy su CoddyKit è strutturato per principianti e studenti avanzati, quindi puoi iniziare da qui o dall'inizio e procedere al tuo ritmo. Questa è la lezione 2 di 4.
Quanto tempo richiede la lezione «Implementazione delle funzionalità fondamentali di RAG e degli agenti»?
La maggior parte delle lezioni CoddyKit richiede circa 5–10 minuti. Ogni lezione è breve e interattiva, quindi fai progressi costanti e riprendi esattamente da dove hai lasciato su web e app.
Posso scrivere ed eseguire codice in questa lezione AI Engineering Academy?
Sì. Ogni lezione AI Engineering Academy include un editor di codice integrato, quindi scrivi ed esegui codice reale direttamente nel tuo browser e ricevi feedback istantaneo dall'IA — nessuna configurazione locale necessaria.
Tutte le lezioni di questo corso
- Progettazione dell’architettura di produzione
- Implementazione delle funzionalità fondamentali di RAG e degli agenti
- Rafforzamento: sicurezza, caching e affidabilità
- Valutazione, deployment e retrospettiva