Módulo 8: Proyecto Integrador RAG con ChromaDB
Cápsula 04: Retrieval, Generation y API
Descripción de la cápsula
Vas a tomar el pipeline RAG end-to-end que construiste en M4/11 (función ask() mínima en ~150 líneas) y exponerlo como API REST profesional con FastAPI: endpoints documentados, validación de inputs, manejo de errores, citas estructuradas y trace_id para debugging.
Prerequisito clave (M4/11): ya tienes la lógica core del RAG — query → retrieve top-k → construir prompt con anti-alucinación → generar con GPT → respuesta con fuentes. Esta cápsula no reinventa esa lógica; la encapsula detrás de una API que terceros (frontend, mobile, otros servicios) pueden consumir.
Lo nuevo de esta cápsula es la capa de servicio: contratos de API estables, validación con Pydantic, observabilidad (latency por endpoint, error tracking, request tracing), y los detalles que separan un script Python de un servicio que un equipo puede operar.
Endpoints de la API
| Método | Ruta | Descripción |
|---|---|---|
POST | /ingest | Dispara ingestion de un directorio (async o sync según diseño) |
GET | /search | Búsqueda por similitud (query, top_k, filtros opcionales) |
POST | /ask | Pregunta RAG: retrieval + generation, respuesta con fuentes |
GET | /collections | Lista colecciones en ChromaDB |
DELETE | /collections/{name} | Borra una colección |
GET | /health | Health check para Docker/K8s |
Contrato del endpoint /ask
{
"answer": "Texto generado basado en el contexto recuperado.",
"sources": [
{
"doc_id": "doc_001",
"source": "/path/to/doc.txt",
"title": "Documento Ejemplo",
"score": 0.87
}
],
"confidence": 0.82,
"trace_id": "550e8400-e29b-41d4-a716-446655440000",
"fallback_reason": null
}
Si no hay evidencia suficiente: fallback_reason = "insufficient retrieval", confidence = 0, sources = [].
Estructura del proyecto
project/
├── api/
│ ├── __init__.py
│ ├── main.py # FastAPI app
│ ├── routes/
│ │ ├── ask.py
│ │ ├── search.py
│ │ └── ingest.py
│ └── services/
│ ├── retrieval.py
│ └── generation.py
├── ingestion/ # (de cápsula 03)
├── config.py
└── run_ingestion.py
Código completo
1. Configuración
# config.py
import os
from dotenv import load_dotenv
load_dotenv()
CHROMA_PATH = os.getenv("CHROMA_PATH", "./chroma_data")
COLLECTION_NAME = os.getenv("COLLECTION_NAME", "rag_docs")
EMBEDDING_MODEL = os.getenv("EMBEDDING_MODEL", "text-embedding-3-small")
LLM_MODEL = os.getenv("LLM_MODEL", "gpt-3.5-turbo")
TOP_K = int(os.getenv("TOP_K", "5"))
SCORE_THRESHOLD = float(os.getenv("SCORE_THRESHOLD", "0.5"))
MAX_QUESTION_LENGTH = int(os.getenv("MAX_QUESTION_LENGTH", "2000"))
2. Servicio de Retrieval
# api/services/retrieval.py
import chromadb
from chromadb.config import Settings
from openai import OpenAI
from typing import List, Optional
import os
client_openai = OpenAI()
chroma_client = chromadb.PersistentClient(path=os.getenv("CHROMA_PATH", "./chroma_data"))
def embed_text(text: str, model: str = "text-embedding-3-small") -> List[float]:
"""Genera embedding para un único texto (query)."""
response = client_openai.embeddings.create(model=model, input=[text])
return response.data[0].embedding
def retrieve(
query: str,
collection_name: str = "rag_docs",
top_k: int = 5,
score_threshold: float = 0.5,
where: Optional[dict] = None,
) -> dict:
"""
Retrieval: embed query, similarity search, filtrar por score.
Returns: documents, metadatas, ids, distances (lower = more similar).
"""
collection = chroma_client.get_collection(collection_name)
query_embedding = embed_text(query)
kwargs = {
"query_embeddings": [query_embedding],
"n_results": top_k,
"include": ["documents", "metadatas", "distances"],
}
if where:
kwargs["where"] = where
result = collection.query(**kwargs)
docs = result["documents"][0] if result["documents"] else []
metas = result["metadatas"][0] if result["metadatas"] else []
ids = result["ids"][0] if result["ids"] else []
distances = result["distances"][0] if result.get("distances") else []
# ChromaDB usa L2 por defecto; convertir a similarity score (1 / (1 + distance))
# O si usas cosine, distances puede ser cosine distance (1 - similarity)
scores = [1 / (1 + d) if d is not None else 0 for d in distances]
# Filtrar por umbral
filtered = [
(doc, meta, id_, score)
for doc, meta, id_, score in zip(docs, metas, ids, scores)
if score >= score_threshold
]
if not filtered:
return {"documents": [], "metadatas": [], "ids": [], "scores": []}
docs_f, metas_f, ids_f, scores_f = zip(*filtered)
return {
"documents": list(docs_f),
"metadatas": list(metas_f),
"ids": list(ids_f),
"scores": list(scores_f),
}
3. Servicio de Generation
# api/services/generation.py
from openai import OpenAI
from typing import List, Optional
client = OpenAI()
def build_context(documents: List[str], metadatas: List[dict]) -> str:
"""Ensambla el contexto para el prompt a partir de los chunks recuperados."""
parts = []
for i, (doc, meta) in enumerate(zip(documents, metadatas), 1):
source = meta.get("source", "unknown")
title = meta.get("title", meta.get("doc_id", "Doc"))
parts.append(f"[Fuente {i}: {title}]\n{doc}")
return "\n\n---\n\n".join(parts)
def generate_answer(
question: str,
documents: List[str],
metadatas: List[dict],
model: str = "gpt-3.5-turbo",
temperature: float = 0.2,
) -> str:
"""
Genera respuesta usando el contexto recuperado.
Si no hay documentos, devuelve mensaje de fallback sin llamar al LLM.
"""
if not documents or not metadatas:
return "No tengo suficiente información en mi base de conocimiento para responder esta pregunta."
context = build_context(documents, metadatas)
prompt = f"""Eres un asistente que responde preguntas basándote ÚNICAMENTE en el contexto proporcionado.
No inventes información. Si el contexto no contiene la respuesta, di que no tienes esa información.
Contexto:
{context}
Pregunta: {question}
Responde de forma clara y concisa. Cita las fuentes cuando sea relevante."""
response = client.chat.completions.create(
model=model,
messages=[
{"role": "system", "content": "Respondes basándote solo en el contexto dado. No inventes datos."},
{"role": "user", "content": prompt},
],
temperature=temperature,
)
return response.choices[0].message.content
def build_sources_response(metadatas: List[dict], scores: List[float]) -> List[dict]:
"""Formatea las fuentes para el contrato de respuesta."""
return [
{
"doc_id": m.get("doc_id", ""),
"source": m.get("source", ""),
"title": m.get("title", m.get("doc_id", "Doc")),
"score": round(s, 2),
}
for m, s in zip(metadatas, scores)
]
4. Rutas FastAPI
# api/routes/ask.py
from fastapi import APIRouter, HTTPException
from pydantic import BaseModel, Field
from uuid import uuid4
import os
from api.services.retrieval import retrieve
from api.services.generation import generate_answer, build_sources_response
router = APIRouter(prefix="/ask", tags=["RAG"])
MAX_QUESTION_LENGTH = int(os.getenv("MAX_QUESTION_LENGTH", "2000"))
class AskRequest(BaseModel):
question: str = Field(..., min_length=1, max_length=MAX_QUESTION_LENGTH)
class SourceItem(BaseModel):
doc_id: str
source: str
title: str
score: float
class AskResponse(BaseModel):
answer: str
sources: list[SourceItem]
confidence: float
trace_id: str
fallback_reason: str | None = None
@router.post("", response_model=AskResponse)
def ask(req: AskRequest):
trace_id = str(uuid4())
question = req.question.strip()
# Retrieval
result = retrieve(
query=question,
top_k=int(os.getenv("TOP_K", "5")),
score_threshold=float(os.getenv("SCORE_THRESHOLD", "0.5")),
)
docs = result["documents"]
metas = result["metadatas"]
scores = result["scores"]
# Fallback si no hay evidencia suficiente
if not docs or not metas:
return AskResponse(
answer="No encontré información relevante para responder tu pregunta.",
sources=[],
confidence=0.0,
trace_id=trace_id,
fallback_reason="insufficient retrieval",
)
# Calcular confidence (promedio de scores normalizado)
confidence = round(sum(scores) / len(scores), 2) if scores else 0.0
# Generation
answer = generate_answer(
question=question,
documents=docs,
metadatas=metas,
model=os.getenv("LLM_MODEL", "gpt-3.5-turbo"),
)
sources = build_sources_response(metas, scores)
return AskResponse(
answer=answer,
sources=[SourceItem(**s) for s in sources],
confidence=confidence,
trace_id=trace_id,
fallback_reason=None,
)
# api/routes/search.py
from fastapi import APIRouter, Query
from pydantic import BaseModel
import os
from api.services.retrieval import retrieve
router = APIRouter(prefix="/search", tags=["Search"])
class SearchResponse(BaseModel):
documents: list[str]
metadatas: list[dict]
ids: list[str]
scores: list[float]
@router.get("", response_model=SearchResponse)
def search(
q: str = Query(..., min_length=1),
top_k: int = Query(5, ge=1, le=20),
score_threshold: float = Query(0.0, ge=0, le=1),
):
result = retrieve(
query=q,
top_k=top_k,
score_threshold=score_threshold,
)
return SearchResponse(
documents=result["documents"],
metadatas=result["metadatas"],
ids=result["ids"],
scores=result["scores"],
)
# api/routes/ingest.py
from fastapi import APIRouter, HTTPException, BackgroundTasks
from pydantic import BaseModel
router = APIRouter(prefix="/ingest", tags=["Ingestion"])
class IngestRequest(BaseModel):
dir_path: str
@router.post("")
def trigger_ingest(req: IngestRequest, background_tasks: BackgroundTasks):
"""Dispara ingestion asíncrona de un directorio."""
from pathlib import Path
if not Path(req.dir_path).exists():
raise HTTPException(status_code=400, detail="Directorio no encontrado")
from ingestion.pipeline import ingest_directory
from ingestion.config import IngestionConfig
config = IngestionConfig()
def run():
ingest_directory(req.dir_path, config=config)
background_tasks.add_task(run)
return {"status": "ingestion started", "dir_path": req.dir_path}
# api/routes/collections.py
from fastapi import APIRouter
import chromadb
import os
router = APIRouter(prefix="/collections", tags=["Collections"])
chroma_client = chromadb.PersistentClient(path=os.getenv("CHROMA_PATH", "./chroma_data"))
@router.get("")
def list_collections():
cols = chroma_client.list_collections()
return [{"name": c.name, "count": c.count()} for c in cols]
@router.delete("/{name}")
def delete_collection(name: str):
try:
chroma_client.delete_collection(name)
return {"status": "deleted", "collection": name}
except Exception as e:
from fastapi import HTTPException
raise HTTPException(status_code=404, detail=str(e))
# api/main.py
from fastapi import FastAPI
from fastapi.middleware.cors import CORSMiddleware
from api.routes.ask import router as ask_router
from api.routes.search import router as search_router
from api.routes.ingest import router as ingest_router
from api.routes.collections import router as collections_router
app = FastAPI(
title="RAG API",
description="Sistema RAG con ChromaDB y OpenAI",
version="1.0.0",
docs_url="/docs", # Swagger
redoc_url="/redoc", # ReDoc
)
app.add_middleware(
CORSMiddleware,
allow_origins=["*"],
allow_credentials=True,
allow_methods=["*"],
allow_headers=["*"],
)
app.include_router(ask_router)
app.include_router(search_router)
app.include_router(collections_router)
@app.get("/health")
def health():
return {"status": "ok"}
@app.get("/")
def root():
return {
"message": "RAG API - Usa /docs para Swagger, /redoc para ReDoc",
"endpoints": ["/ask", "/search", "/ingest", "/collections", "/health"],
}
5. Ejecutar la API
# Instalar dependencias
pip install fastapi uvicorn chromadb openai python-dotenv pydantic
# Ejecutar
uvicorn api.main:app --reload --host 0.0.0.0 --port 8000
Acceso:
- Swagger: http://localhost:8000/docs
- ReDoc: http://localhost:8000/redoc
Reglas de calidad de respuesta
| Regla | Implementación |
|---|---|
| No inventar sin evidencia | Fallback explícito cuando retrieval vacío o scores bajos |
| Incluir fuentes | sources con doc_id, source, title, score |
| Evitar respuestas sin contexto | No llamar al LLM si no hay chunks; devolver mensaje estándar |
| Trazabilidad | trace_id en cada respuesta para correlación con logs |
| Validación de input | max_length en pregunta para evitar abuse |
Ejemplo de uso
POST /ask
curl -X POST "http://localhost:8000/ask" \
-H "Content-Type: application/json" \
-d '{"question": "¿Qué es un vector database?"}'
Respuesta esperada:
{
"answer": "Un vector database es un sistema de almacenamiento optimizado para...",
"sources": [
{
"doc_id": "intro_vectors",
"source": "/docs/intro.txt",
"title": "intro",
"score": 0.87
}
],
"confidence": 0.85,
"trace_id": "550e8400-e29b-41d4-a716-446655440000",
"fallback_reason": null
}
GET /search
curl "http://localhost:8000/search?q=vector%20database&top_k=3"
Ejercicios prácticos
Ejercicio 1: Extender /ask con metadata filtering
Añade un parámetro opcional doc_id o source a la request para filtrar la búsqueda solo en ciertos documentos.
Solución: En AskRequest añade doc_id: str | None = None. En la llamada a retrieve, si doc_id está presente, usa where={"doc_id": doc_id}.
Ejercicio 2: Respuesta "no tengo evidencia" cuando confidence < 0.5
Si el promedio de scores es menor a 0.5, no llames al LLM. Devuelve el mensaje de fallback y fallback_reason="low_confidence".
Solución: Antes de generate_answer, calcula confidence. Si confidence < 0.5, return AskResponse(..., fallback_reason="low_confidence") sin llamar al LLM.
Ejercicio 3: Incluir trace_id en logs
Configura un middleware o dependency que capture el trace_id y lo añada a cada log del request. Usa structlog o logging con contexto.
Solución: Middleware que genera trace_id al inicio del request y lo guarda en request.state.trace_id. En cada endpoint, pasas ese valor a la respuesta. Para logs, usa logging.LoggerAdapter con extra={"trace_id": trace_id}.
Ejercicio 4: Timeout y retry para OpenAI
Añade timeout de 30s a las llamadas de OpenAI y retry (1 intento) si falla por timeout.
Solución: client.chat.completions.create(..., timeout=30). Envuelve en try/except; en Timeout o APIConnectionError, retry una vez con time.sleep(2).
Ejercicio 5: Endpoint POST /ingest
Crea un endpoint que reciba {"dir_path": "/path/to/docs"} y ejecute el pipeline de ingestion de la cápsula 03. Para no bloquear, considera BackgroundTasks de FastAPI.
Solución:
from fastapi import BackgroundTasks
from ingestion.pipeline import ingest_directory
@router.post("/ingest")
def trigger_ingest(payload: dict, bg: BackgroundTasks):
dir_path = payload.get("dir_path")
if not dir_path:
raise HTTPException(400, "dir_path required")
def run():
ingest_directory(dir_path)
bg.add_task(run)
return {"status": "ingestion started", "dir": dir_path}
Ejercicio 6: Documentación de errores
Define un esquema de error estándar: {"error": str, "trace_id": str, "detail": str}. Usa un exception handler global para devolverlo en 4xx/5xx.
Solución:
from fastapi import Request
from fastapi.responses import JSONResponse
@app.exception_handler(Exception)
def global_exception_handler(request: Request, exc: Exception):
trace_id = getattr(request.state, "trace_id", "unknown")
return JSONResponse(
status_code=500,
content={"error": "internal_error", "trace_id": trace_id, "detail": str(exc)},
)
Troubleshooting
"Respuestas correctas pero sin trazabilidad"
Causa: No se incluyen fuentes ni trace_id.
Solución: Asegúrate de que AskResponse siempre incluya sources (aunque vacío) y trace_id. Genera trace_id al inicio del request y pásalo por todo el flujo.
"Respuestas inventadas cuando no hay contexto"
Causa: Se llama al LLM incluso con retrieval vacío o scores muy bajos.
Solución: Añade validación: si len(docs) == 0 o max(scores) < 0.5, no llames a generate_answer. Devuelve fallback explícito con fallback_reason.
"API inestable bajo carga"
Causa: Sin límites, timeouts ni control de concurrencia.
Solución: Limita longitud de pregunta (max_length=2000). Añade timeout a OpenAI. Considera rate limiting con slowapi o similar. Revisa connection pooling de ChromaDB.
"Swagger/ReDoc no muestran schemas correctos"
Causa: Modelos Pydantic sin response_model o con tipos complejos no documentados.
Solución: Usa response_model=AskResponse en el decorator. Define todos los modelos con BaseModel y campos tipados. FastAPI genera los schemas automáticamente.
"Errores 500 sin detalle útil"
Causa: Excepciones no capturadas, logs insuficientes.
Solución: Exception handler global que loguee traceback y devuelva mensaje genérico al cliente (no revelar internals). En desarrollo, debug=True puede ayudar; en producción, logs estructurados con trace_id.
Resumen
- La API expone
/ask,/search,/collections,/healthcon FastAPI. - El endpoint
/askintegra retrieval (ChromaDB + OpenAI embeddings) y generation (GPT-3.5-turbo). - Contrato de respuesta:
answer,sources,confidence,trace_id,fallback_reason. - Fallback explícito cuando no hay evidencia suficiente; no se inventan respuestas.
- Swagger (/docs) y ReDoc (/redoc) documentan la API automáticamente.
- Separar retrieval y generation en funciones independientes facilita tests y evolución.
- Validación de input, timeout y manejo de errores mejoran robustez.
Recursos adicionales
- FastAPI Tutorial
- FastAPI Background Tasks
- ChromaDB Query
- OpenAI Chat Completions
- Pydantic Models
- Uvicorn
- RAG Prompt Engineering
- Structured Logging (Python)
Tiempo estimado: 35-40 minutos
Siguiente: 05-testing-evaluacion.md