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
ThreadPoolExecutory 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=5000el 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 locales | No | 1 |
| PersistentClient + embeddings locales | No (marginal) | 1-2 |
| PersistentClient + OpenAI embeddings | Sí | 3-4 |
| HTTP server (chromadb run) + OpenAI | Sí | 4-8 |
Cuidado con paralelizar mal:
- Demasiados workers con OpenAI → rate limit cruzado entre threads, errores 429.
- Workers con
EphemeralClientdistintos → 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=32en máquina de 16 GB. batch_sizeenorme 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:
-
PersistentClientporque 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. -
batch_size=200porque 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, perobatch_size=200es seguro y conocido. -
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
-
Retry con backoff (5s, 10s) para manejar rate limits transitorios sin perder progreso.
-
Idempotencia con
skip_existingporque la laptop puede entrar en sleep, network puede fallar — el pipeline debe reanudar sin duplicar. -
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
ThreadPoolExecutorsolo ayuda cuando el cuello de botella es I/O remoto (OpenAI). ConPersistentClient+ 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.
EphemeralClientes para tests; producción usaPersistentClient.- 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
- ChromaDB — Adding Data — Operaciones de inserción
- ChromaDB — Persistent vs Ephemeral Client — Diferencias y cuándo usar cada uno
- OpenAI — Rate Limits Guide — Límites por tier
- tqdm Documentation — Progress bars
- Tenacity — Retry Library — Alternativa más robusta a retry manual
- SQLite Concurrency Internals — Por qué paralelizar PersistentClient da poco speedup
Tiempo estimado: 30-35 minutos Siguiente: 06-query-optimization.md