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-small en vez del default ChromaDB, y cómo manejar la API key correctamente.
  • M4/10 — Chunking de documentos: ya entiendes por qué chunk_size=512 con overlap=50 es el rango razonable para documentos técnicos, y cómo RecursiveCharacterTextSplitter decide 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

  1. Cargar documentos fuente — PDF, TXT, MD según tipo, normalizar encoding.
  2. ChunkingRecursiveCharacterTextSplitter con chunk_size=512, overlap=50 (justificado en M4/10).
  3. Embeddings — OpenAI text-embedding-3-small en batches de 100-200 (justificado en M4/09).
  4. Inserción en ChromaDB — batch de 1K-2K con IDs estables y metadata.
  5. 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ámetroValorJustificación
chunk_size512Balance contexto/granularidad, compatible con embedding model
chunk_overlap50Continuidad semántica sin duplicar mucho
batch_size_chromadb1000-2000Optimal throughput para ChromaDB
batch_size_embeddings100-200Respetar 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:

  1. IDs determinísticos: {doc_id}_chunk_{index} — mismo doc siempre produce mismos IDs.
  2. Borrar colección antes de re-ingestar (simple): client.delete_collection(name) y crear de nuevo.
  3. Upsert (si ChromaDB soporta): Algunas versiones permiten add con 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


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