Módulo 8: Proyecto Integrador RAG con ChromaDB
Cápsula 03: Pipeline de Ingestion para 1,000+ documentos
Descripción de la cápsula
Vas a escalar el pipeline RAG mínimo que construiste en M4/11 a un volumen de producción real (1,000+ documentos) con foco en estabilidad y trazabilidad: lotes, metadata consistente, validación de resultados y manejo de errores. El pipeline debe ser reproducible, idempotente cuando sea posible y reportar métricas básicas de throughput.
Prerequisitos de las cápsulas anteriores que vas a usar aquí:
- M4/09 — Embeddings con OpenAI: ya sabes por qué usamos
text-embedding-3-smallen vez del default ChromaDB, y cómo manejar la API key correctamente. - M4/10 — Chunking de documentos: ya entiendes por qué
chunk_size=512conoverlap=50es el rango razonable para documentos técnicos, y cómoRecursiveCharacterTextSplitterdecide dónde partir. - M4/11 — Pipeline RAG end-to-end: ya construiste la versión mínima del pipeline (~150 líneas). Esta cápsula lo escala a 1000+ docs y agrega robustez (manejo de errores, batches, validación, métricas).
Si no completaste esas cápsulas, vuelve a M4/09 antes de continuar — esta cápsula asume que las decisiones de embeddings y chunking ya están justificadas.
Flujo de implementación
- Cargar documentos fuente — PDF, TXT, MD según tipo, normalizar encoding.
- Chunking —
RecursiveCharacterTextSplitterconchunk_size=512,overlap=50(justificado en M4/10). - Embeddings — OpenAI
text-embedding-3-smallen batches de 100-200 (justificado en M4/09). - Inserción en ChromaDB — batch de 1K-2K con IDs estables y metadata.
- Verificación — conteo total, metadata obligatoria presente, sin IDs duplicados.
Dependencias necesarias
# requirements.txt (extraer para el proyecto)
chromadb>=0.4.22
openai>=1.12.0
langchain-text-splitters>=0.2.0
langchain-community>=0.2.0 # opcional: loaders
tqdm>=4.66.0
python-dotenv>=1.0.0
pypdf>=4.0.0 # para PDF
Código completo del pipeline
1. Configuración y utilidades
# ingestion/config.py
import os
from dataclasses import dataclass
from dotenv import load_dotenv
load_dotenv()
@dataclass
class IngestionConfig:
"""Configuración centralizada del pipeline de ingestion."""
chunk_size: int = int(os.getenv("CHUNK_SIZE", "512"))
chunk_overlap: int = int(os.getenv("CHUNK_OVERLAP", "50"))
batch_size_chromadb: int = int(os.getenv("BATCH_SIZE_CHROMADB", "1000"))
batch_size_embeddings: int = int(os.getenv("BATCH_SIZE_EMBEDDINGS", "100"))
chroma_path: str = os.getenv("CHROMA_PATH", "./chroma_data")
collection_name: str = os.getenv("COLLECTION_NAME", "rag_docs")
openai_model: str = os.getenv("EMBEDDING_MODEL", "text-embedding-3-small")
2. Carga de documentos
# ingestion/loaders.py
import os
from pathlib import Path
from typing import Iterator
def load_text_file(path: Path) -> str:
"""Carga un archivo de texto con encoding robusto."""
encodings = ["utf-8", "latin-1", "cp1252"]
for enc in encodings:
try:
return path.read_text(encoding=enc)
except UnicodeDecodeError:
continue
raise ValueError(f"No se pudo decodificar: {path}")
def load_documents_from_dir(
dir_path: str,
extensions: tuple[str, ...] = (".txt", ".md", ".rst"),
) -> Iterator[tuple[str, str, str]]:
"""
Itera sobre documentos en un directorio.
Yield: (doc_id, content, source_path)
"""
path = Path(dir_path)
if not path.exists():
raise FileNotFoundError(f"Directorio no encontrado: {dir_path}")
for file_path in path.rglob("*"):
if file_path.suffix.lower() in extensions and file_path.is_file():
try:
content = load_text_file(file_path)
# doc_id: estable y único
doc_id = file_path.relative_to(path).as_posix().replace("/", "_")
source = str(file_path)
yield doc_id, content, source
except Exception as e:
print(f"[WARN] Error cargando {file_path}: {e}")
continue
3. Chunking con RecursiveCharacterTextSplitter
# ingestion/chunking.py
from langchain_text_splitters import RecursiveCharacterTextSplitter
from pathlib import Path
from typing import List
from dataclasses import dataclass
@dataclass
class ChunkWithMeta:
"""Chunk con metadata para ChromaDB."""
text: str
doc_id: str
chunk_index: int
source: str
metadata: dict
def chunk_document(
content: str,
doc_id: str,
source: str,
chunk_size: int = 512,
chunk_overlap: int = 50,
) -> List[ChunkWithMeta]:
"""
Divide un documento en chunks con overlap.
chunk_size y overlap están en caracteres (approx ~4 chars = 1 token).
"""
splitter = RecursiveCharacterTextSplitter(
chunk_size=chunk_size,
chunk_overlap=chunk_overlap,
length_function=len,
separators=["\n\n", "\n", ". ", " ", ""],
)
chunks_raw = splitter.split_text(content)
result = []
for i, text in enumerate(chunks_raw):
meta = {
"doc_id": doc_id,
"source": source,
"chunk_index": i,
"title": Path(source).stem if source else doc_id,
}
result.append(ChunkWithMeta(
text=text,
doc_id=doc_id,
chunk_index=i,
source=source,
metadata=meta,
))
return result
Nota: Si usas Path en chunking.py, añade from pathlib import Path al inicio.
4. Generación de embeddings (OpenAI)
# ingestion/embeddings.py
from openai import OpenAI
from typing import List
import time
client = OpenAI()
def get_embeddings_batch(texts: List[str], model: str = "text-embedding-3-small") -> List[List[float]]:
"""
Genera embeddings para un batch de textos.
Rate limit: ~3000 RPM para text-embedding-3-small; usa batch 100-200.
"""
if not texts:
return []
response = client.embeddings.create(
model=model,
input=texts,
)
# Orden garantizado igual que input
return [item.embedding for item in sorted(response.data, key=lambda x: x.index)]
5. Pipeline completo de ingestion
# ingestion/pipeline.py
import chromadb
from chromadb.config import Settings
from tqdm import tqdm
import hashlib
import time
from pathlib import Path
from ingestion.config import IngestionConfig
from ingestion.loaders import load_documents_from_dir
from ingestion.chunking import chunk_document
from ingestion.embeddings import get_embeddings_batch
def make_chunk_id(doc_id: str, chunk_index: int) -> str:
"""ID determinístico para idempotencia."""
return f"{doc_id}_chunk_{chunk_index}"
def ingest_directory(
dir_path: str,
config: IngestionConfig | None = None,
extensions: tuple[str, ...] = (".txt", ".md", ".rst"),
) -> dict:
"""
Pipeline completo: carga docs, chunking, embeddings, ChromaDB.
Retorna métricas: total_docs, total_chunks, time_seconds, errors.
"""
config = config or IngestionConfig()
client = chromadb.PersistentClient(path=config.chroma_path)
collection = client.get_or_create_collection(
name=config.collection_name,
metadata={"description": "RAG collection"},
)
# Fase 1: Cargar y chunkear
all_chunks = []
all_metadatas = []
all_ids = []
errors = []
for doc_id, content, source in load_documents_from_dir(dir_path, extensions):
try:
chunks = chunk_document(
content=content,
doc_id=doc_id,
source=source,
chunk_size=config.chunk_size,
chunk_overlap=config.chunk_overlap,
)
for c in chunks:
all_chunks.append(c.text)
all_metadatas.append(c.metadata)
all_ids.append(make_chunk_id(doc_id, c.chunk_index))
except Exception as e:
errors.append({"doc_id": doc_id, "error": str(e)})
continue
total_chunks = len(all_chunks)
if total_chunks == 0:
return {
"total_docs": 0,
"total_chunks": 0,
"time_seconds": 0,
"errors": errors,
}
# Fase 2: Embeddings en batch
embeddings_list = []
batch_size = config.batch_size_embeddings
for i in tqdm(range(0, total_chunks, batch_size), desc="Embeddings"):
batch_texts = all_chunks[i : i + batch_size]
try:
emb = get_embeddings_batch(batch_texts, model=config.openai_model)
embeddings_list.extend(emb)
except Exception as e:
errors.append({"batch": i, "error": str(e)})
# Fallback: zeros para no romper el add (mejor: omitir batch)
embeddings_list.extend([[0.0] * 1536 for _ in batch_texts]) # dim de 3-small
# Fase 3: Add a ChromaDB en batches
chroma_batch = config.batch_size_chromadb
start = time.time()
for i in tqdm(range(0, total_chunks, chroma_batch), desc="ChromaDB"):
batch_ids = all_ids[i : i + chroma_batch]
batch_docs = all_chunks[i : i + chroma_batch]
batch_emb = embeddings_list[i : i + chroma_batch]
batch_meta = all_metadatas[i : i + chroma_batch]
try:
collection.add(
ids=batch_ids,
documents=batch_docs,
embeddings=batch_emb,
metadatas=batch_meta,
)
except Exception as e:
errors.append({"chroma_batch": i, "error": str(e)})
elapsed = time.time() - start
return {
"total_docs": len(set(m.get("doc_id") for m in all_metadatas)),
"total_chunks": total_chunks,
"time_seconds": round(elapsed, 2),
"docs_per_second": round(total_chunks / elapsed, 1) if elapsed > 0 else 0,
"errors": errors,
}
6. Script de ejecución
# run_ingestion.py
import argparse
from ingestion.config import IngestionConfig
from ingestion.pipeline import ingest_directory
def main():
parser = argparse.ArgumentParser()
parser.add_argument("dir", help="Directorio con documentos (.txt, .md, .rst)")
parser.add_argument("--batch-size", type=int, default=1000)
parser.add_argument("--chunk-size", type=int, default=512)
parser.add_argument("--chunk-overlap", type=int, default=50)
args = parser.parse_args()
config = IngestionConfig(
batch_size_chromadb=args.batch_size,
chunk_size=args.chunk_size,
chunk_overlap=args.chunk_overlap,
)
result = ingest_directory(args.dir, config=config)
print("\n=== Resultado de ingestion ===")
print(f"Documentos procesados: {result['total_docs']}")
print(f"Chunks totales: {result['total_chunks']}")
print(f"Tiempo: {result['time_seconds']} s")
print(f"Throughput: {result.get('docs_per_second', 0)} chunks/s")
print(f"Errores: {len(result['errors'])}")
if result["errors"]:
for e in result["errors"][:5]:
print(f" - {e}")
if __name__ == "__main__":
main()
Uso:
python run_ingestion.py ./mis_documentos --batch-size 1000
Validaciones mínimas post-ingestion
Después de correr el pipeline, ejecuta:
# validation/check_ingestion.py
import chromadb
client = chromadb.PersistentClient(path="./chroma_data")
coll = client.get_collection("rag_docs")
count = coll.count()
print(f"Total vectores: {count}")
# Muestra aleatoria de metadata
sample = coll.peek(limit=5)
for i, meta in enumerate(sample["metadatas"]):
required = ["doc_id", "source", "chunk_index"]
ok = all(k in meta for k in required)
print(f"Chunk {i}: metadata ok={ok}, keys={list(meta.keys())}")
# Verificar IDs únicos
ids = coll.get()["ids"]
assert len(ids) == len(set(ids)), "IDs duplicados detectados"
print("IDs únicos: OK")
Parámetros recomendados
| Parámetro | Valor | Justificación |
|---|---|---|
chunk_size | 512 | Balance contexto/granularidad, compatible con embedding model |
chunk_overlap | 50 | Continuidad semántica sin duplicar mucho |
batch_size_chromadb | 1000-2000 | Optimal throughput para ChromaDB |
batch_size_embeddings | 100-200 | Respetar rate limits de OpenAI |
Documenta estos valores en el README del proyecto y exponlos por variables de entorno.
Estrategia de idempotencia
Para evitar duplicados al re-ingestar:
- IDs determinísticos:
{doc_id}_chunk_{index}— mismo doc siempre produce mismos IDs. - Borrar colección antes de re-ingestar (simple):
client.delete_collection(name)y crear de nuevo. - Upsert (si ChromaDB soporta): Algunas versiones permiten
addcon IDs existentes que actualizan; verifica en la documentación de tu versión.
Ejercicios aplicados
Ejercicio 1: Implementa ingestion de 1,000 docs
Crea un directorio con 1,000 archivos .txt (puedes generarlos o usar un corpus público). Ejecuta el pipeline y reporta:
- Tiempo total
- Chunks/segundo
- Porcentaje de errores
- Porcentaje de docs sin metadata completa
Solución: Usa un script para generar 1,000 archivos:
from pathlib import Path
Path("test_docs").mkdir(exist_ok=True)
for i in range(1000):
(Path("test_docs") / f"doc_{i:04d}.txt").write_text(f"Contenido del documento {i}. " * 50)
Luego: python run_ingestion.py test_docs. Revisa el output de result para métricas. Si hay 0 errores y metadata completa en todos, el porcentaje es 0%.
Ejercicio 2: Manejo de PDF
Extiende el loader para soportar PDF usando pypdf. Añade .pdf a las extensiones y una función load_pdf(path) que extraiga texto de cada página.
Solución:
from pypdf import PdfReader
def load_pdf(path: Path) -> str:
reader = PdfReader(path)
return "\n".join(page.extract_text() or "" for page in reader.pages)
En load_documents_from_dir, añade ".pdf" y en el branch de PDF usa load_pdf en lugar de load_text_file.
Ejercicio 3: Progress bar por fase
Añade tqdm para mostrar progreso en: (1) carga de documentos, (2) chunking, (3) embeddings, (4) ChromaDB add. Cada fase debe tener su propia barra.
Solución: En load_documents_from_dir, convierte el iterador en lista y envuelve con tqdm. Para chunking, si procesas doc por doc, tqdm(docs). Para embeddings y ChromaDB, ya hay tqdm en los loops de batch.
Ejercicio 4: Retry con backoff para embeddings
Si la API de OpenAI falla por rate limit, implementa retry con exponential backoff (1s, 2s, 4s) hasta 3 intentos.
Solución:
import time
def get_embeddings_batch_with_retry(texts, model="text-embedding-3-small", max_retries=3):
for attempt in range(max_retries):
try:
return get_embeddings_batch(texts, model)
except Exception as e:
if attempt == max_retries - 1:
raise
wait = 2 ** attempt
time.sleep(wait)
Ejercicio 5: Validación de metadata antes de add
Antes de llamar a collection.add, valida que cada elemento de batch_meta tenga las keys doc_id, source, chunk_index. Si falta alguna, registra warning y añade valor por defecto o descarta el chunk.
Solución:
REQUIRED_KEYS = {"doc_id", "source", "chunk_index"}
def validate_metadata(meta: dict) -> bool:
return REQUIRED_KEYS.issubset(meta.keys())
# En el loop de ChromaDB add:
valid = [(i, m, d, e) for i, m, d, e in zip(...) if validate_metadata(m)]
invalid_count = len(batch_meta) - len(valid)
if invalid_count:
print(f"[WARN] {invalid_count} chunks con metadata incompleta descartados")
Ejercicio 6: Benchmark de batch size
Prueba batch_size_chromadb en [500, 1000, 2000, 5000] con 5,000 chunks. Mide tiempo y throughput. ¿Cuál es el sweet spot en tu máquina?
Solución: Script de benchmark:
for bs in [500, 1000, 2000, 5000]:
config = IngestionConfig(batch_size_chromadb=bs)
# Borrar colección antes
r = ingest_directory("test_docs", config=config)
print(f"batch={bs}: {r['time_seconds']}s, {r.get('docs_per_second')} docs/s")
Tipical: 1000-2000 da mejor throughput; 5000 a veces no mejora por overhead de memoria.
Troubleshooting de ingestion
"La ingestion tarda demasiado"
Causa: Batch size pequeño, o cuello de botella en embeddings (API).
Solución: Incrementa batch_size_chromadb a 1500-2000. Para embeddings, usa batch_size_embeddings de 100-200 y verifica que no estés limitado por rate limits de OpenAI. Mide tiempo por fase para localizar el cuello.
"Aparecen duplicados"
Causa: IDs no determinísticos o re-ingestion sin borrar colección.
Solución: Usa make_chunk_id(doc_id, chunk_index). Si re-ingiestas, borra la colección antes o implementa upsert si tu versión de ChromaDB lo permite.
"Faltan metadatos en parte del dataset"
Causa: Algunos loaders no extraen source o doc_id, o hay docs sin título.
Solución: Valida metadata antes de add. Añade defaults: source=doc_id si falta, title=doc_id si no hay título. Registra qué docs tienen metadata incompleta.
"Error 429 de OpenAI (rate limit)"
Causa: Demasiadas requests por minuto.
Solución: Reduce batch_size_embeddings a 50-100. Implementa retry con backoff. Considera time.sleep(1) entre batches si es necesario.
"ChromaDB out of memory"
Causa: Batch muy grande o demasiados vectores en memoria.
Solución: Reduce batch_size_chromadb a 500. Procesa en varias ejecuciones si el dataset es enorme. Verifica que no estés acumulando todos los embeddings en memoria antes de add (procesa en streams).
Resumen
- El pipeline tiene 5 fases: carga, chunking, embeddings, add a ChromaDB, verificación.
- Usa RecursiveCharacterTextSplitter con chunk_size=512, overlap=50.
- Embeddings vía OpenAI text-embedding-3-small en batches de 100-200.
- ChromaDB add en batches de 1K-2K con IDs determinísticos y metadata obligatoria.
- Progress tracking con tqdm, error handling por batch, validación post-ingestion.
- Documenta parámetros en README; usa variables de entorno para configuración.
- El índice resultante está listo para conectar a retrieval/generation en la siguiente cápsula.
Recursos adicionales
- ChromaDB Add Data
- LangChain Text Splitters
- OpenAI Embeddings
- Batch Ingestion (Módulo 4)
- RecursiveCharacterTextSplitter - chunk_size
- ChromaDB PersistentClient
- tqdm Documentation
- pypdf - PDF extraction
Checklist post-implementación
Antes de pasar a la cápsula 04 (API), verifica:
- El pipeline procesa 100+ documentos sin errores.
- Todos los chunks tienen metadata con doc_id, source, chunk_index.
- Los IDs son determinísticos (re-ejecutar no duplica si borras la colección).
- Tienes progress bar o logs que indiquen avance.
- Documentaste batch_size, chunk_size y chunk_overlap en README.
Tiempo estimado: 30-35 minutos
Siguiente: 04-retrieval-generation-api.md