Módulo 4: ChromaDB Setup y Configuración

Cápsula 05: Batch Ingestion en ChromaDB — del concepto al pipeline operable

Descripción de la cápsula

En M03/05 aprendiste por qué single inserts no escalan y los principios de un pipeline batch robusto: idempotencia, retry con backoff, observabilidad. Esta cápsula es la implementación concreta sobre ChromaDB. Vas a construir y benchmarkear un pipeline operable que ingerie 10K-100K documentos eficientemente, manejando los detalles específicos del backend: cuándo es seguro paralelizar (y cuándo no), cómo medir el sweet spot de batch_size empíricamente, qué errores típicos lanza ChromaDB y cómo recuperarte sin perder progreso.

Cuando termines, vas a tener un pipeline reutilizable que podés copiar a tu proyecto y adaptar — y vas a entender por qué cada componente está donde está.

Al finalizar esta cápsula serás capaz de:

  • ✅ Encontrar el batch_size óptimo para tu hardware con un benchmark de 5 minutos
  • ✅ Decidir cuándo paralelizar con ThreadPoolExecutor y cuándo no
  • ✅ Implementar el pipeline canónico con tqdm + retry + logging estructurado
  • ✅ Manejar los errores específicos de ChromaDB: ID duplicado, dimension mismatch, full disk
  • ✅ Diferenciar entre EphemeralClient (testing) y PersistentClient (producción) en el contexto de ingestion masiva
  • ✅ Anticipar el cuello de botella real: muchas veces no es ChromaDB, es la API de embeddings

Tiempo estimado: 30-35 minutos


Single insert vs batch — la demostración rápida

Antes de pasar al pipeline, una demostración mínima del impacto. Esto es lo que pasa con 1,000 documentos:

import chromadb
import time

client = chromadb.Client()  # ephemeral, in-memory
collection = client.create_collection("demo_single")

docs = [f"Document {i}" for i in range(1000)]
ids = [f"doc_{i}" for i in range(1000)]

# Single inserts (uno por uno)
start = time.perf_counter()
for doc, doc_id in zip(docs, ids):
    collection.add(documents=[doc], ids=[doc_id])
single_time = time.perf_counter() - start

print(f"Single inserts (1000 docs): {single_time:.2f}s")
# ~12 segundos


# Batch insert (todos juntos)
collection2 = client.create_collection("demo_batch")
start = time.perf_counter()
collection2.add(documents=docs, ids=ids)
batch_time = time.perf_counter() - start

print(f"Batch insert (1000 docs):    {batch_time:.2f}s")
print(f"Speedup: {single_time / batch_time:.0f}x")
# ~0.15 segundos, ~80x speedup

Por qué tan dramático: ya lo cubrimos en M03/05 — cada llamada a add() tiene overhead constante (validación, escritura, update del índice). Hacer 1,000 llamadas multiplica el overhead 1,000 veces; hacer una sola llamada con 1,000 docs lo paga una vez.

Pero hay un detalle importante para ChromaDB: el batch máximo razonable depende de si usás embeddings locales (default) o remotos (OpenAI vía embedding_function).


Encontrar tu sweet spot de batch_size

Las recomendaciones generales (cubiertas en M03/05) dicen 100-200 con OpenAI, 500-1000 local. Pero el sweet spot depende de tu hardware, tu modelo y tu dataset. Acá está el script para encontrarlo:

# benchmark_batch_size.py
import chromadb
from chromadb.utils import embedding_functions
import os
import time

# Setup con embeddings locales (más rápido para benchmark)
client = chromadb.PersistentClient(path="./chroma_bench_batch")

# Genera 10K docs sintéticos
total_docs = 10_000
docs = [
    f"This is document number {i} discussing technical topics like vector databases, "
    f"embedding models, and retrieval systems. Each document has approximately 30 words "
    f"to simulate typical chunk sizes in a real RAG application." for i in range(total_docs)
]
ids = [f"doc_{i:05d}" for i in range(total_docs)]


def benchmark_batch_size(batch_size: int) -> dict:
    """Mide throughput para un batch_size específico."""
    name = f"bench_bs_{batch_size}"
    try:
        client.delete_collection(name)
    except Exception:
        pass
    collection = client.create_collection(name)

    start = time.perf_counter()
    for i in range(0, total_docs, batch_size):
        end = min(i + batch_size, total_docs)
        collection.add(documents=docs[i:end], ids=ids[i:end])
    elapsed = time.perf_counter() - start

    throughput = total_docs / elapsed
    n_batches = (total_docs + batch_size - 1) // batch_size

    return {
        "batch_size": batch_size,
        "total_time_s": elapsed,
        "throughput_docs_per_s": throughput,
        "n_batches": n_batches,
    }


print(f"{'batch_size':>10} {'time (s)':>10} {'docs/s':>10} {'batches':>10}")
print("-" * 45)
for bs in [50, 100, 200, 500, 1000, 2000, 5000]:
    result = benchmark_batch_size(bs)
    print(
        f"{result['batch_size']:>10} "
        f"{result['total_time_s']:>10.2f} "
        f"{result['throughput_docs_per_s']:>10.0f} "
        f"{result['n_batches']:>10}"
    )

Output típico (Macbook M2, embeddings locales):

batch_size   time (s)     docs/s    batches
---------------------------------------------
        50      18.40       543        200
       100      11.20       893        100
       200       6.85      1460         50
       500       4.30      2326         20
      1000       3.85      2597         10
      2000       3.62      2762          5
      5000       4.10      2439          2

Lectura:

  • Throughput crece rápido entre 50 y 500 (overhead amortizado).
  • Sweet spot en batch_size=2000 — ~2,762 docs/segundo.
  • A batch_size=5000 el throughput cae ligeramente — probablemente memoria swap o GC pressure.

Tu output va a ser distinto. Puede que tu hardware prefiera 1000 (RAM ajustada) o 5000 (RAM holgada). Por eso el benchmark sobre tu máquina vale más que copiar números de un tutorial.

Con OpenAI embeddings, el patrón cambia

Si usás OpenAIEmbeddingFunction en lugar de embeddings locales, el cuello de botella deja de ser ChromaDB y pasa a ser la API de OpenAI:

openai_ef = embedding_functions.OpenAIEmbeddingFunction(
    api_key=os.getenv("OPENAI_API_KEY"),
    model_name="text-embedding-3-small"
)

Restricciones de OpenAI:

  • Máximo ~8191 tokens por request
  • Rate limits por tier (tier 1: ~3000 RPM, tier 2: ~5000 RPM)

Output típico con OpenAI:

batch_size   time (s)     docs/s    batches
---------------------------------------------
        50      32.10       312        200
       100      18.50       541        100
       200      12.20       820         50
       500       9.40      1064         20
      1000      11.20       893         10  ← ya hits rate limits
      2000     ERROR (token limit)

Con OpenAI, el sweet spot baja a ~200-500 docs/batch. Pasarse genera errores 429 (rate limit) o 400 (max tokens). Por eso la recomendación general "200 con OpenAI" — está calibrada para no romperse.


El pipeline canónico

Acá está la implementación que vas a copiar a tu proyecto. Combina batch + tqdm + retry + idempotencia + manejo de errores específicos de ChromaDB.

# ingestion_pipeline.py
import chromadb
from chromadb.utils import embedding_functions
import os
import time
from dataclasses import dataclass, field
from tqdm import tqdm

# Errores específicos que vamos a manejar
try:
    from openai import RateLimitError, APIError, APITimeoutError
except ImportError:
    RateLimitError = APIError = APITimeoutError = Exception


@dataclass
class IngestionStats:
    total_input: int = 0
    inserted: int = 0
    skipped_existing: int = 0
    failed: list[int] = field(default_factory=list)
    duration_seconds: float = 0.0

    @property
    def throughput_per_second(self) -> float:
        return self.inserted / self.duration_seconds if self.duration_seconds > 0 else 0

    def report(self):
        print(f"\nIngestion stats:")
        print(f"  Input total:       {self.total_input}")
        print(f"  Inserted:          {self.inserted}")
        print(f"  Skipped (existing): {self.skipped_existing}")
        print(f"  Failed batches:    {len(self.failed)}")
        print(f"  Duration:          {self.duration_seconds:.1f}s")
        print(f"  Throughput:        {self.throughput_per_second:.0f} docs/s")


def ingest_to_chromadb(
    collection: chromadb.Collection,
    documents: list[str],
    ids: list[str],
    metadatas: list[dict] | None = None,
    batch_size: int = 200,
    sleep_between_batches: float = 0.0,
    max_retries: int = 3,
    skip_existing: bool = True,
) -> IngestionStats:
    """
    Pipeline de ingestion robusto para ChromaDB.

    - Idempotente: con `skip_existing=True`, los IDs ya presentes se omiten.
    - Recuperable: cada batch tiene retry con exponential backoff.
    - Observable: progress bar + stats finales.
    - Tolera fallos parciales: si un batch falla todos los retries, se anota y continúa.
    """
    assert len(documents) == len(ids), "documents and ids must have same length"
    if metadatas is not None:
        assert len(metadatas) == len(ids), "metadatas length mismatch"

    stats = IngestionStats(total_input=len(ids))
    start = time.perf_counter()

    # Filtrar IDs existentes (idempotencia)
    if skip_existing:
        existing = collection.get(ids=ids, include=[])
        existing_ids = set(existing['ids'])
        if existing_ids:
            print(f"Skipping {len(existing_ids)} existing IDs")
        new_indices = [i for i, doc_id in enumerate(ids) if doc_id not in existing_ids]
        stats.skipped_existing = len(ids) - len(new_indices)
    else:
        new_indices = list(range(len(ids)))

    # Procesar en batches
    n_batches = (len(new_indices) + batch_size - 1) // batch_size

    for batch_num in tqdm(range(n_batches), desc="Ingesting", unit="batch"):
        b_start = batch_num * batch_size
        b_end = min(b_start + batch_size, len(new_indices))
        batch_indices = new_indices[b_start:b_end]

        batch_docs = [documents[i] for i in batch_indices]
        batch_ids = [ids[i] for i in batch_indices]
        batch_metas = [metadatas[i] for i in batch_indices] if metadatas else None

        # Retry con exponential backoff
        success = False
        for attempt in range(max_retries):
            try:
                kwargs = {"documents": batch_docs, "ids": batch_ids}
                if batch_metas is not None:
                    kwargs["metadatas"] = batch_metas
                collection.add(**kwargs)
                stats.inserted += len(batch_ids)
                success = True
                break

            except RateLimitError:
                wait = (2 ** attempt) * 5  # 5, 10, 20s
                print(f"\n  Rate limit (batch {batch_num+1}), waiting {wait}s...")
                time.sleep(wait)

            except (APIError, APITimeoutError) as e:
                wait = (2 ** attempt) * 2
                print(f"\n  API error (batch {batch_num+1}): {e}, waiting {wait}s...")
                time.sleep(wait)

            except ValueError as e:
                # ChromaDB lanza ValueError para problemas de schema (dimension mismatch, etc.)
                print(f"\n  Schema error (batch {batch_num+1}): {e}")
                # No tiene sentido reintentar errores de schema
                break

            except Exception as e:
                wait = 1
                print(f"\n  Unexpected error (batch {batch_num+1}): {e}, waiting {wait}s...")
                time.sleep(wait)

        if not success:
            stats.failed.append(batch_num + 1)

        if sleep_between_batches > 0:
            time.sleep(sleep_between_batches)

    stats.duration_seconds = time.perf_counter() - start
    return stats

Probarlo

# test_pipeline.py
import os
import chromadb
from chromadb.utils import embedding_functions

openai_ef = embedding_functions.OpenAIEmbeddingFunction(
    api_key=os.getenv("OPENAI_API_KEY"),
    model_name="text-embedding-3-small"
)

client = chromadb.PersistentClient(path="./chroma_test")
collection = client.get_or_create_collection(
    name="pipeline_test",
    embedding_function=openai_ef,
    metadata={"hnsw:space": "cosine", "hnsw:M": 32, "hnsw:construction_ef": 200},
)

# Generar dataset
docs = [f"Document {i} about technical topics in AI engineering." for i in range(2000)]
ids = [f"doc_{i:05d}" for i in range(2000)]
metas = [{"index": i, "category": "demo"} for i in range(2000)]

stats = ingest_to_chromadb(
    collection,
    documents=docs,
    ids=ids,
    metadatas=metas,
    batch_size=200,
    sleep_between_batches=0.3,  # Margen para OpenAI rate limits
    max_retries=3,
)
stats.report()

Output esperado primer run:

Skipping 0 existing IDs
Ingesting: 100%|██████████| 10/10 [00:32<00:00, 3.27s/batch]

Ingestion stats:
  Input total:       2000
  Inserted:          2000
  Skipped (existing): 0
  Failed batches:    0
  Duration:          32.7s
  Throughput:        61 docs/s

Output al re-correr (idempotencia):

Skipping 2000 existing IDs
Ingesting: 0it [00:00, ?it/s]

Ingestion stats:
  Input total:       2000
  Inserted:          0
  Skipped (existing): 2000
  Failed batches:    0
  Duration:          0.4s
  Throughput:        0 docs/s

La segunda ejecución termina en 400ms porque verificó IDs existentes antes de insertar. Eso es idempotencia operativa.


¿Cuándo paralelizar con ThreadPoolExecutor?

La intuición dice "más threads = más rápido". Para ChromaDB, casi siempre es falso. Hay que entender el porqué.

El backend de persistencia limita la concurrencia

ChromaDB PersistentClient usa SQLite por debajo. SQLite serializa internamente las escrituras concurrentes — solo una transacción de escritura puede estar activa a la vez. Lanzar 8 threads que insertan concurrentemente da el mismo resultado que 1 thread, pero con overhead extra de coordinación.

Benchmark sobre PersistentClient:

from concurrent.futures import ThreadPoolExecutor

def ingest_one_batch(args):
    collection, batch_docs, batch_ids = args
    collection.add(documents=batch_docs, ids=batch_ids)

# Sequential
collection_seq = client.get_or_create_collection("seq")
batches = [(collection_seq, docs[i:i+200], ids[i:i+200]) for i in range(0, 2000, 200)]
start = time.perf_counter()
for batch in batches:
    ingest_one_batch(batch)
seq_time = time.perf_counter() - start

# Concurrent (4 workers)
collection_conc = client.get_or_create_collection("conc")
batches_conc = [(collection_conc, docs[i:i+200], ids[i:i+200]) for i in range(0, 2000, 200)]
start = time.perf_counter()
with ThreadPoolExecutor(max_workers=4) as executor:
    list(executor.map(ingest_one_batch, batches_conc))
conc_time = time.perf_counter() - start

print(f"Sequential: {seq_time:.2f}s")
print(f"Concurrent (4 workers): {conc_time:.2f}s")
print(f"Speedup: {seq_time/conc_time:.2f}x")

Output típico (PersistentClient):

Sequential: 8.40s
Concurrent (4 workers): 7.80s
Speedup: 1.08x   ← marginal, casi nada

¿Cuándo SÍ paraleliza con beneficio?

Cuando el cuello de botella es fuera de ChromaDB, específicamente la API de embeddings remota. Si cada batch llama a OpenAI por embeddings, esa llamada es I/O-bound (espera de red). Mientras un thread espera, otros pueden hacer trabajo.

# Con embeddings remotos (OpenAI), paralelización SÍ ayuda
# 4 threads pueden tener 4 requests pendientes simultáneamente
# Cada uno espera ~150ms a OpenAI; con 4 paralelos efectivamente quedan 40ms

# Output típico con OpenAI + 4 workers:
# Sequential:  120s
# Concurrent:   45s
# Speedup:    2.7x

Regla práctica

Configuración¿Paralelizar?Workers recomendados
EphemeralClient + embeddings localesNo1
PersistentClient + embeddings localesNo (marginal)1-2
PersistentClient + OpenAI embeddings3-4
HTTP server (chromadb run) + OpenAI4-8

Cuidado con paralelizar mal:

  • Demasiados workers con OpenAI → rate limit cruzado entre threads, errores 429.
  • Workers con EphemeralClient distintos → cada uno tiene su DB en memoria; los datos no se comparten.
  • Workers escribiendo a la misma collection sin lock → ChromaDB lo serializa pero perdés tiempo en contención.

Errores específicos de ChromaDB y cómo manejarlos

Error 1: dimensión mismatch

ValueError: Embedding dimension 1536 does not match collection dimensionality 384

Causa: estás pasando un embedding (o usando un embedding_function) que produce vectores de tamaño distinto al que la collection espera.

Cuándo pasa:

  • Cambiaste embedding_function (ej: de default ChromaDB 384-dim a OpenAI 1536-dim) sin re-crear la collection.
  • Pasás embeddings=... directamente con vectores del tamaño equivocado.

Cómo prevenir: usá una sola fuente de truth para el modelo de embeddings. Si la collection se creó con OpenAIEmbeddingFunction, todas las queries y inserts deben usar el mismo embedding_function.

Cómo recuperar: crear collection nueva con la dimensionalidad correcta y re-insertar.

Error 2: ID duplicado

chromadb.errors.IDAlreadyExistsError: ID 'doc_42' already exists

Causa: estás insertando un ID que ya existe en la collection.

Cuándo pasa:

  • Re-ejecutás el pipeline sin idempotencia.
  • Generás IDs no deterministas (UUID aleatorio) y un retry inserta el mismo doc dos veces.

Cómo prevenir: el pipeline canónico arriba ya lo cubre con skip_existing=True.

Cómo recuperar: o filtrás los IDs existentes (idempotencia), o usás collection.upsert() en vez de collection.add(). upsert actualiza si el ID existe.

Error 3: full disk

sqlite3.OperationalError: database or disk is full

Causa: PersistentClient escribe a SQLite + parquet en disco. Cuando el disco se llena, falla.

Cuándo pasa: datasets grandes (>1M vectores) en máquinas con poco espacio.

Cómo prevenir: monitorear espacio en disco. Estimar storage:

storage ≈ N × D × 4 bytes (vectores) + overhead SQLite

Para 1M vectores × 1536 dim:
  vectores: 6 GB
  índice HNSW: ~3 GB
  metadata + IDs: ~500 MB
  total: ~10 GB en disco

Cómo recuperar: liberar espacio, o migrar a almacenamiento más grande, o pasar a vector DB distribuida.

Error 4: out of memory durante el build

Causa: insertás demasiados vectores con EphemeralClient (todo en RAM) o el grafo HNSW excede la RAM disponible.

Cuándo pasa:

  • 5M+ vectores con M=32 en máquina de 16 GB.
  • batch_size enorme con embeddings grandes.

Cómo prevenir: estimar RAM como vimos en M04/03. Si excede, bajar M o migrar a vector DB con almacenamiento persistente.

Error 5: ChromaDB server unavailable (HTTP mode)

chromadb.errors.ChromaError: HTTP request failed: connection refused

Causa: estás usando chromadb.HttpClient(host=..., port=...) y el server está caído o no acepta conexiones.

Cómo manejar: retry con backoff, healthcheck antes de empezar pipeline, alerta si el server no responde después de N reintentos.


Trampas y errores comunes

Trampa 1: usar EphemeralClient para ingestion masiva

El error:

client = chromadb.Client()  # ephemeral, all in RAM
collection = client.create_collection("docs")
collection.add(documents=[...] * 1_000_000, ids=[...])  # OOM!

Síntoma: OOM kill cuando el dataset crece.

Cómo prevenir: EphemeralClient es solo para tests y prototipos pequeños. Para cualquier ingestion real, PersistentClient con path en disco.

client = chromadb.PersistentClient(path="./chroma_db")

Trampa 2: ignorar el costo de re-ejecutar el pipeline

El error: pipeline corre 30 minutos. Falla en el 95%. Re-ejecutás desde cero. Pagás 30 minutos más + el costo de embeddings duplicados.

Cómo prevenir: idempotencia con skip_existing=True. Si re-corres, los 95% ya insertados se saltean en segundos.

Trampa 3: paralelizar con PersistentClient sin medir

El error: copiás ThreadPoolExecutor(max_workers=8) de un tutorial. Ves el código "elegante" y deployás.

Síntoma: throughput igual o peor que sequential. Logs muestran contención de SQLite.

Cómo prevenir: benchmarkear sequential vs concurrent con tu setup específico antes de decidir. La regla práctica de la tabla arriba.

Trampa 4: olvidar el embedding_function al recuperar la collection

El error:

# Día 1: crear con OpenAI ef
collection = client.create_collection("docs", embedding_function=openai_ef)
collection.add(documents=docs, ids=ids)

# Día 2: recuperar SIN el embedding_function
collection_v2 = client.get_collection("docs")  # ❌
collection_v2.add(documents=more_docs, ids=more_ids)
# Usa el embedding default de ChromaDB → dimension mismatch

Cómo prevenir: siempre pasar el mismo embedding_function al recuperar la collection:

collection_v2 = client.get_collection("docs", embedding_function=openai_ef)

Trampa 5: ingestion sin tracking de progreso → matar procesos sanos

El error: script corre sin tqdm ni logs. 25 minutos sin output. Operador piensa que se colgó. kill -9.

Síntoma: trabajo perdido.

Cómo prevenir: siempre tqdm o logging por batch. Es 1 línea de código y previene este escenario.

Trampa 6: no separar el pipeline de embeddings y el de inserción

El error: una sola función hace download del documento, chunking, embedding y inserción. Cuando algo falla, no sabés dónde.

Cómo prevenir: separar fases en pipeline:

[descargar docs] → [chunkear] → [embeber] → [insertar a ChromaDB]
   pipeline 1      pipeline 2    pipeline 3       pipeline 4

Cada fase tiene su propia tolerancia a errores y se puede ejecutar/reintentar independientemente. Si embedding falla, no rehacés el chunking.


Ejercicio aplicado

Escenario: te dan un dataset de 50,000 PDFs. Especificaciones:

  • Cada PDF tiene 5-30 páginas. Ya están chunkeados (en M4/10 vimos chunking) → ~150,000 chunks totales.
  • Modelo de embeddings: OpenAI text-embedding-3-small.
  • Hardware: laptop con 16 GB RAM, MacOS, conexión estable.
  • OpenAI tier 2 (5000 RPM, 5M tokens/min).
  • Budget: $20 USD para embeddings.
  • SLA operativo: el ingestion debe terminar en ≤2 horas (te dieron una mañana de trabajo).
  • Robustez: si el script crashea (network drop, laptop sleep), debe poder reanudar sin duplicar.

Tu trabajo: implementá el pipeline (o adaptá el canónico de arriba) con configuración justificada. Calculá el costo estimado y el tiempo, y ajustá si no cabe en presupuesto.

Solución

Análisis de viabilidad:

total_chunks = 150_000
avg_tokens_per_chunk = 500  # estimación para chunks de ~400 chars
total_tokens = total_chunks * avg_tokens_per_chunk  # 75M tokens

# Costo
cost_per_million_tokens = 0.02
cost = total_tokens / 1_000_000 * cost_per_million_tokens
# = $1.50 — cabe holgadamente en $20

# Throughput requerido
target_minutes = 120
required_chunks_per_minute = total_chunks / target_minutes  # 1,250 chunks/min
required_tokens_per_minute = total_tokens / target_minutes  # 625K tokens/min

print(f"Costo estimado: ${cost}")
print(f"Throughput requerido: {required_chunks_per_minute} chunks/min")
print(f"Tokens/min requeridos: {required_tokens_per_minute:.0f}")
print(f"Tier 2 permite 5,000,000 tokens/min → margen 8x")

Resultado: $1.50 (8% del presupuesto), throughput requerido bien por debajo del rate limit. Es viable.

Configuración del pipeline:

import chromadb
from chromadb.utils import embedding_functions
import os
from concurrent.futures import ThreadPoolExecutor
import time

openai_ef = embedding_functions.OpenAIEmbeddingFunction(
    api_key=os.getenv("OPENAI_API_KEY"),
    model_name="text-embedding-3-small"
)

client = chromadb.PersistentClient(path="./chroma_pdfs")
collection = client.get_or_create_collection(
    name="pdf_chunks",
    embedding_function=openai_ef,
    metadata={
        "hnsw:space": "cosine",
        "hnsw:M": 32,
        "hnsw:construction_ef": 200,
        "hnsw:search_ef": 50,
    },
)

# Pipeline con paralelismo controlado (OpenAI es I/O bound → paraleliza bien)
def ingest_chunks_parallel(chunks, ids, metadatas, batch_size=200, workers=3):
    """Pipeline con 3 workers paralelos para amortizar latencia de OpenAI."""

    # Idempotencia: filtrar IDs existentes
    existing = collection.get(ids=ids, include=[])
    existing_set = set(existing['ids'])
    new_indices = [i for i, doc_id in enumerate(ids) if doc_id not in existing_set]

    print(f"Total: {len(ids)}, ya existen: {len(existing_set)}, a insertar: {len(new_indices)}")

    # Crear batches
    batches = []
    for batch_start in range(0, len(new_indices), batch_size):
        batch_idx = new_indices[batch_start:batch_start + batch_size]
        batches.append({
            "documents": [chunks[i] for i in batch_idx],
            "ids": [ids[i] for i in batch_idx],
            "metadatas": [metadatas[i] for i in batch_idx],
        })

    def process_batch(batch_data):
        for attempt in range(3):
            try:
                collection.add(**batch_data)
                return True
            except Exception as e:
                if attempt < 2:
                    time.sleep((2 ** attempt) * 5)  # 5, 10s
                else:
                    print(f"FAILED batch after retries: {e}")
                    return False
        return False

    start = time.perf_counter()
    success = 0
    failed = 0

    with ThreadPoolExecutor(max_workers=workers) as executor:
        results = list(executor.map(process_batch, batches))

    success = sum(1 for r in results if r)
    failed = sum(1 for r in results if not r)
    elapsed = time.perf_counter() - start

    print(f"\nCompleted in {elapsed/60:.1f} minutes")
    print(f"Successful batches: {success}/{len(batches)}")
    print(f"Failed batches: {failed}")
    print(f"Throughput: {len(new_indices)/elapsed:.0f} docs/s")

Decisiones justificadas:

  1. PersistentClient porque 150K vectores × 1536 dim × 4 bytes = 920 MB solo en vectores; con HNSW overhead total ~1.5 GB. Cabe en 16 GB pero el resto de la app necesita RAM también — mejor en disco.

  2. batch_size=200 porque OpenAI text-embedding-3-small permite hasta ~16 docs por request por límite de tokens (500 × 16 = 8000). En la práctica, el SDK de OpenAI maneja batches mayores splitting internamente, pero batch_size=200 es seguro y conocido.

  3. 3 workers paralelos porque:

    • Cada worker espera 150-300ms por request a OpenAI (I/O bound, no CPU)
    • Tier 2 (5000 RPM) → ~83 RPS — con 3 workers a 2 RPS cada uno = 6 RPS, dentro del límite con margen 14x
    • Más workers (4+) no ayudaría porque el cuello pasa a SQLite serialización
  4. Retry con backoff (5s, 10s) para manejar rate limits transitorios sin perder progreso.

  5. Idempotencia con skip_existing porque la laptop puede entrar en sleep, network puede fallar — el pipeline debe reanudar sin duplicar.

  6. Metadata HNSW para producción (M=32, construction_ef=200) — vamos a usar este índice por meses, vale la pena el build de mejor calidad.

Tiempo estimado de ejecución:

  • 150K chunks / 6 RPS efectivo = 25,000 segundos = ~7 minutos (sí, mucho menos de 2 horas)
  • Pero más realista, considerando overhead: 15-30 minutos.
  • Cabe holgadamente en la mañana de trabajo.

Plan B si algo se rompe:

  • Si laptop entra en sleep durante el run: el pipeline está corriendo desde laptop, sleep mata el proceso. Al reanudar, ejecutar de nuevo — idempotencia salta los ya insertados.
  • Si rate limit dispara repetidamente: bajar workers a 2 o subir sleep_between_batches. Reanudar.
  • Si batches específicos fallan después de 3 retries: anotarlos en log, intentar manualmente al final con batches más chicos (50 en vez de 200).

Verificación post-ingestion:

# Sanity check
print(f"Total docs en collection: {collection.count()}")
# Esperado: 150,000

# Test query
result = collection.query(
    query_texts=["¿cómo configurar HNSW?"],
    n_results=5,
)
for doc, meta in zip(result['documents'][0], result['metadatas'][0]):
    print(f"  [{meta['source']}] {doc[:80]}...")

Resumen y siguiente paso

Lo que aprendiste:

  • Single inserts en ChromaDB son ~80x más lentos que batch — el sweet spot es 1000-2000 docs con embeddings locales, 100-500 con OpenAI.
  • Encontrar el batch_size óptimo para tu hardware vale 5 minutos de benchmark — no copies números de tutoriales.
  • Paralelizar con ThreadPoolExecutor solo ayuda cuando el cuello de botella es I/O remoto (OpenAI). Con PersistentClient + embeddings locales, no aporta.
  • El pipeline canónico combina: idempotencia (skip existing IDs), retry con backoff, manejo de errores específicos de ChromaDB, observabilidad con tqdm.
  • Errores típicos de ChromaDB y cómo manejarlos: dimension mismatch, ID duplicado, full disk, OOM.
  • EphemeralClient es para tests; producción usa PersistentClient.
  • Separar pipeline en fases (descarga → chunking → embedding → inserción) mejora la recuperación ante fallos.

Checkpoint: antes de avanzar, deberías poder:

  • Encontrar el batch_size óptimo con un benchmark de 5 minutos sobre tu hardware.
  • Decidir cuántos workers usar según embeddings local vs remoto.
  • Implementar el pipeline canónico con idempotencia y manejo de errores.

Siguiente cápsula: 06 — Query Optimization.

Acabás de meter datos al sistema eficientemente. Ahora vas a optimizar el otro lado: las queries. Vas a aprender a medir latencia con percentiles (no promedios), a entender qué parámetros mueven la aguja (n_results, ef_search, metadata filtering), y a benchmarkear cambios antes de deployarlos.


Recursos

  1. ChromaDB — Adding Data — Operaciones de inserción
  2. ChromaDB — Persistent vs Ephemeral Client — Diferencias y cuándo usar cada uno
  3. OpenAI — Rate Limits Guide — Límites por tier
  4. tqdm Documentation — Progress bars
  5. Tenacity — Retry Library — Alternativa más robusta a retry manual
  6. SQLite Concurrency Internals — Por qué paralelizar PersistentClient da poco speedup

Tiempo estimado: 30-35 minutos Siguiente: 06-query-optimization.md