Módulo 3: Features Esenciales para RAG
Cápsula 05: Batch Operations — el costo oculto del ingestion ingenuo
Descripción de la cápsula
El primer instinto cuando tienes 1 millón de documentos para indexar es escribir un loop:
for doc in documents:
collection.add(documents=[doc.text], ids=[doc.id], metadatas=[doc.metadata])
Eso funciona perfectamente para 100 documentos. Falla catastróficamente para 100,000. No por un bug — por arquitectura. Cada llamada add() con un solo documento dispara overhead que es independiente de cuánto datos estés insertando: serialización, validación, escritura al storage, actualización del índice HNSW, fsync. Si haces 1 millón de llamadas, multiplicas ese overhead 1 millón de veces.
Batch operations es la diferencia entre que tu ingestion termine en 2 minutos o en 3 horas. No es optimización prematura — es la diferencia entre un pipeline operable y uno que rompe el flujo de trabajo del equipo cada vez que llegan documentos nuevos.
Esta cápsula te da el modelo mental para entender por qué single inserts fallan a escala, los criterios para elegir el batch size correcto según tu caso, y los patrones de error más caros (el que destruye datos en silencio cuando un batch falla a la mitad).
Al finalizar esta cápsula serás capaz de:
- ✅ Calcular el costo real (en tiempo y dinero) de ingerir N documentos con single inserts vs batches
- ✅ Elegir un
batch_sizecon criterio basado en RAM disponible, costo de embeddings y latencia aceptable - ✅ Identificar los tres puntos de falla en pipelines de batch ingestion: rate limits, OOM, crash a mitad de batch
- ✅ Diseñar idempotencia para que retry de batches no duplique documentos
- ✅ Anticipar el modo de falla más caro: batch que parece exitoso pero pierde documentos en silencio
Tiempo estimado: 30-40 minutos
Por qué single inserts no escalan
Vamos a hacer el cálculo que casi nadie hace antes de descubrirlo en producción.
Anatomía de un single insert
Cuando llamas collection.add(documents=[doc], ids=[id]), ChromaDB hace estas operaciones:
- Validación del input (~0.5ms): tipos, IDs únicos, metadata válida.
- Generación de embedding vía la
embedding_functionconfigurada:- Default ChromaDB local: ~5-15ms por documento.
- OpenAI API: ~80-150ms por llamada (incluye latencia de red + procesamiento del API).
- Escritura al storage (SQLite o backend persistente): ~3-8ms.
- Actualización del índice HNSW (recálculo de conexiones del grafo): ~1-3ms.
fsyncopcional para garantizar durabilidad: ~5-10ms.
Total: ~10-25ms con embedding local, ~85-180ms con OpenAI.
Multiplicado por volumen:
| Documentos | Single insert + local | Single insert + OpenAI | Batch + OpenAI (batch=200) |
|---|---|---|---|
| 1,000 | 15-25 segundos | 1.5-3 minutos | ~6 segundos |
| 10,000 | 2.5-4 minutos | 14-30 minutos | ~60 segundos |
| 100,000 | 25-40 minutos | 2.4-5 horas | ~10 minutos |
| 1,000,000 | 4-7 horas | 24-50 horas ❌ | ~100 minutos |
El insight: lo que hace explotar el tiempo no es el procesamiento real de los datos — es el overhead constante por llamada. Generar el embedding de un documento cuesta 100ms. Generar embeddings de 100 documentos en una sola llamada API cuesta ~150ms total (200ms con la latencia de red, una vez). Es decir, en batch obtienes 100 embeddings por casi el mismo tiempo que cuesta uno solo.
El cálculo de costo en API
Para OpenAI específicamente, hay un costo monetario adicional al hacer single inserts:
# Calculadora de costo por API calls
COST_PER_REQUEST_OVERHEAD_SECONDS = 0.1 # latencia + procesamiento por llamada
COST_PER_TOKEN = 0.02 / 1_000_000 # USD por token (text-embedding-3-small)
TOKENS_PER_DOC_AVG = 500
scenarios = {
"single_insert": {"requests": 100_000, "tokens_per_request": 500},
"batch_size_10": {"requests": 10_000, "tokens_per_request": 5_000},
"batch_size_100": {"requests": 1_000, "tokens_per_request": 50_000},
"batch_size_500": {"requests": 200, "tokens_per_request": 250_000},
}
for name, s in scenarios.items():
total_tokens = s["requests"] * s["tokens_per_request"]
monetary_cost = total_tokens * COST_PER_TOKEN
time_overhead = s["requests"] * COST_PER_REQUEST_OVERHEAD_SECONDS
print(f"\n=== {name} ===")
print(f" Requests al API: {s['requests']:,}")
print(f" Costo monetario: ${monetary_cost:.2f}")
print(f" Tiempo de overhead: {time_overhead/60:.1f} minutos")
Output:
=== single_insert ===
Requests al API: 100,000
Costo monetario: $1.00
Tiempo de overhead: 166.7 minutos
=== batch_size_10 ===
Requests al API: 10,000
Costo monetario: $1.00
Tiempo de overhead: 16.7 minutos
=== batch_size_100 ===
Requests al API: 1,000
Costo monetario: $1.00
Tiempo de overhead: 1.7 minutos
=== batch_size_500 ===
Requests al API: 200
Costo monetario: $1.00
Tiempo de overhead: 0.3 minutos
Lectura clave: el costo monetario es el mismo porque depende de tokens totales, no de número de requests. Pero el tiempo de overhead se reduce de 2.7 horas (single insert) a 20 segundos (batch=500). Y aún hay otro costo oculto: cada request puede fallar (rate limit, timeout) — más requests = más oportunidades de fallar.
Cómo elegir el batch_size correcto
No hay un valor mágico. Hay tres restricciones que se cruzan.
Restricción 1: API rate limits
OpenAI permite hasta 8191 tokens por request en text-embedding-3-small, y un máximo de elementos en el array input por request (~2048 para texto). Si tus documentos promedian 500 tokens:
- Batch teórico máximo por tokens: 8191 / 500 ≈ 16 docs
- Batch teórico máximo por elementos: 2048 docs
El bottleneck son los tokens, no los elementos. Para text-embedding-3-small, batches de 100-200 documentos son seguros.
Para otros modelos (Cohere, Sentence Transformers locales), revisar la documentación específica.
Restricción 2: RAM disponible
El batch entero vive en memoria mientras se procesa:
# Aproximación de RAM consumida por batch
batch_size = 1000
embedding_dim = 1536
bytes_per_float = 4 # float32
ram_per_batch = batch_size * embedding_dim * bytes_per_float
# = 1000 * 1536 * 4 = 6.144 MB por batch
# Si procesas en paralelo (4 workers):
total_ram = ram_per_batch * 4 # = 24.6 MB
Eso es trivial para batches de 1000. Pero si haces batch_size=100_000 con vectores de 4096 dim (modelos grandes) y 16 workers paralelos, consumes ~26 GB solo en embeddings — y eso sin contar metadata, IDs, y el overhead de Python.
Regla práctica: mantén el RAM por batch (incluyendo workers paralelos) bajo el 30% del RAM disponible de tu máquina.
Restricción 3: latencia de fallo
Si un batch falla a mitad de procesamiento, ¿cuánto trabajo pierdes? Con batch_size=10000, un timeout te hace reiniciar 10000 docs. Con batch_size=200, solo 200.
Trade-off claro:
- Batch grande (1000+): mejor throughput, peor recuperación de errores.
- Batch pequeño (50-200): mejor recuperación, más overhead.
Recomendaciones por escenario
| Escenario | batch_size sugerido | Razón |
|---|---|---|
| Ingestion inicial de dataset estático | 500-1000 | Throughput sobre todo, rate limit no es problema en una sola tirada |
| Pipeline de producción con docs llegando continuamente | 50-200 | Latencia importa, fallos son más caros si afectan otros docs |
| API de OpenAI con embeddings | 100-200 | Sweet spot: amortiza overhead sin riesgo de rate limit |
| Embeddings locales (CPU) | 200-500 | Sin rate limit, limita por RAM |
| Embeddings locales (GPU) | 500-2000 | GPU ama batches grandes |
| Documentos enormes (>5000 tokens) | 20-50 | Reduce batch para no exceder límite de tokens por request |
Default razonable si no sabes: batch_size=200 con OpenAI, batch_size=500 local.
El pipeline correcto: idempotente, observable, recuperable
# pipeline_batch_ingestion.py
import os
import time
from dataclasses import dataclass
from tqdm import tqdm
from openai import RateLimitError, APIError
import chromadb
from chromadb.utils import embedding_functions
@dataclass
class IngestionResult:
total_docs: int
inserted: int
skipped: int
failed_batches: list[int]
duration_seconds: float
def ingest_documents_batched(
collection,
documents: list[str],
metadatas: list[dict],
ids: list[str],
batch_size: int = 200,
sleep_between: float = 0.0,
max_retries: int = 3,
) -> IngestionResult:
"""
Pipeline de ingestion con tres garantías:
- Idempotente: si un id ya existe, no se inserta (sin error).
- Observable: progress bar + logging por batch.
- Recuperable: si un batch falla, se reintenta antes de abandonarlo.
"""
assert len(documents) == len(metadatas) == len(ids), (
"documents, metadatas, ids deben tener el mismo largo"
)
start = time.time()
inserted = 0
skipped = 0
failed_batches = []
# Verificar IDs ya existentes (idempotencia)
existing_response = collection.get(ids=ids, include=[])
existing_ids = set(existing_response['ids'])
if existing_ids:
print(f" {len(existing_ids)} IDs ya existen, se omitirán")
new_indices = [i for i, doc_id in enumerate(ids) if doc_id not in existing_ids]
skipped = len(ids) - len(new_indices)
# Procesar en batches con barra de progreso
n_batches = (len(new_indices) + batch_size - 1) // batch_size
for batch_num in tqdm(range(n_batches), desc="Ingesting batches"):
batch_start = batch_num * batch_size
batch_end = min(batch_start + batch_size, len(new_indices))
batch_indices = new_indices[batch_start:batch_end]
batch_docs = [documents[i] for i in batch_indices]
batch_metas = [metadatas[i] for i in batch_indices]
batch_ids = [ids[i] for i in batch_indices]
# Retry con exponential backoff
success = False
for attempt in range(max_retries):
try:
collection.add(
documents=batch_docs,
metadatas=batch_metas,
ids=batch_ids,
)
inserted += len(batch_ids)
success = True
break
except RateLimitError:
wait = (2 ** attempt) * 5 # 5, 10, 20 segundos
print(f"\n Rate limit en batch {batch_num+1}, esperando {wait}s...")
time.sleep(wait)
except APIError as e:
wait = (2 ** attempt) * 2
print(f"\n API error en batch {batch_num+1}: {e}, esperando {wait}s...")
time.sleep(wait)
except Exception as e:
print(f"\n Error inesperado en batch {batch_num+1}: {e}")
break
if not success:
failed_batches.append(batch_num + 1)
print(f"\n ❌ Batch {batch_num+1} falló después de {max_retries} intentos")
if sleep_between > 0:
time.sleep(sleep_between)
duration = time.time() - start
return IngestionResult(
total_docs=len(ids),
inserted=inserted,
skipped=skipped,
failed_batches=failed_batches,
duration_seconds=duration,
)
Probar el pipeline
# test_pipeline.py
openai_ef = embedding_functions.OpenAIEmbeddingFunction(
api_key=os.getenv("OPENAI_API_KEY"),
model_name="text-embedding-3-small"
)
client = chromadb.PersistentClient(path="./chroma_batch_test")
collection = client.get_or_create_collection(
name="batch_demo",
embedding_function=openai_ef
)
# Generar dataset sintético
docs = [f"Document number {i} about vector databases and RAG systems." for i in range(1000)]
metas = [{"index": i, "category": "test"} for i in range(1000)]
ids = [f"doc_{i:04d}" for i in range(1000)]
result = ingest_documents_batched(
collection=collection,
documents=docs,
metadatas=metas,
ids=ids,
batch_size=200,
sleep_between=0.5,
max_retries=3,
)
print(f"\nResultado:")
print(f" Total: {result.total_docs}")
print(f" Insertados: {result.inserted}")
print(f" Omitidos (ya existían): {result.skipped}")
print(f" Batches fallidos: {result.failed_batches}")
print(f" Duración: {result.duration_seconds:.1f}s ({result.inserted/result.duration_seconds:.0f} docs/s)")
Output esperado (primera ejecución):
0 IDs ya existen, se omitirán
Ingesting batches: 100%|█████████| 5/5 [00:14<00:00, 2.96s/batch]
Resultado:
Total: 1000
Insertados: 1000
Omitidos (ya existían): 0
Batches fallidos: []
Duración: 14.8s (68 docs/s)
Output al reejecutar (idempotencia):
1000 IDs ya existen, se omitirán
Ingesting batches: 0it [00:00, ?it/s]
Resultado:
Total: 1000
Insertados: 0
Omitidos (ya existían): 1000
Batches fallidos: []
Duración: 0.3s (0 docs/s)
La segunda ejecución no duplica nada y termina en 300ms porque verificó IDs existentes antes de insertar. Eso es idempotencia: ejecutar el pipeline dos veces produce el mismo estado final que ejecutarlo una vez.
Trampas y errores comunes
Trampa 1: batch_size demasiado grande golpea rate limits o OOM
El error: ves benchmarks que dicen "más batch = más throughput" y subes batch_size=5000 con OpenAI.
Síntoma A: errores RateLimitError o BadRequestError: max tokens exceeded. El batch nunca se inserta.
Síntoma B (más sutil): el batch funciona localmente pero en CI con menos RAM, OOM kills el proceso a mitad de ingestion.
Por qué pasa: tu pipeline excede los límites del API (8191 tokens por request) o del entorno de ejecución.
Cómo prevenir:
- Calcula tokens por batch antes de ejecutar:
import tiktoken
enc = tiktoken.encoding_for_model("text-embedding-3-small")
batch = documents[0:200]
total_tokens = sum(len(enc.encode(d)) for d in batch)
print(f"Batch tokens: {total_tokens}")
# Si > 8000, reducir batch_size
- Mide RAM usado en local antes de deployar a CI.
Trampa 2: pipeline sin idempotencia → duplicación masiva
El error: un batch falla a mitad de ingestion. El operador relanza el script. Los documentos exitosos del primer intento + todos los del segundo intento → duplicados.
Síntoma: queries devuelven resultados con scores idénticos casi consecutivos (mismos documentos con IDs distintos). collection.count() es mayor de lo esperado.
Por qué pasa: el pipeline asume "estado vacío al inicio". Si no se cumple, falla.
Cómo prevenir: la implementación de arriba lo cubre — verificar collection.get(ids=...) antes de insertar y solo insertar los IDs nuevos.
Trampa 3: no manejar failures de batch parciales
El error: un batch de 200 docs falla con timeout. El código asume "todo el batch falló" y los marca para retry. Pero ChromaDB internamente ya escribió 150 de 200 antes del timeout.
Síntoma: retry inserta los 200 de nuevo, generando 150 duplicados.
Por qué pasa: ChromaDB no tiene transacciones — cada add() puede ser parcialmente exitoso.
Cómo prevenir:
- Si tienes idempotencia por ID (recomendado), el retry es seguro.
- Si tu pipeline genera IDs nuevos en cada run (mal patrón), arreglarlo: los IDs deben ser deterministas a partir del documento (hash, doc_id estable).
# ❌ IDs no deterministas
import uuid
ids = [str(uuid.uuid4()) for _ in documents]
# ✅ IDs deterministas a partir del contenido
import hashlib
ids = [
hashlib.md5(doc.encode()).hexdigest() for doc in documents
]
# O usando metadata estable
ids = [f"doc_{meta['source']}_{meta['chunk_index']}" for meta in metadatas]
Trampa 4: ignorar el progreso, sin visibilidad
El error: loop con 50,000 iteraciones sin progress bar ni logs. El script corre 30 minutos, parece colgado, alguien lo mata pensando que está en deadlock.
Síntoma: trabajo perdido, frustración.
Cómo prevenir: siempre tqdm o equivalente en pipelines de batch ingestion. Es 1 línea de código y previene matar procesos sanos.
Trampa 5: sleep_between=0 con OpenAI en producción
El error: procesar 100 batches consecutivos sin pausa. OpenAI tiene rate limits por minuto (RPM, request per minute). Si haces 100 requests en 10 segundos, vas a hitting el rate limit.
Síntoma: primer 30% del pipeline funciona, después empiezan errores 429.
Cómo prevenir: agregar time.sleep(0.5-1.0) entre batches grandes. Pierdes <5% de throughput total y evitas los retries innecesarios.
Trampa 6: batch ingestion paralelo sin control de concurrencia
El error: "más rápido = más threads". Ejecutas 16 workers con batch_size=1000 cada uno.
Síntoma: rate limit 429 en TODOS los workers simultáneamente. Algunos workers entran en backoff, otros siguen, el sistema es impredecible. ChromaDB locks por escritura concurrente al SQLite.
Cómo prevenir:
- Para OpenAI: respetar el rate limit del tier. Tier 1 permite ~3000 RPM, lo que con
batch_size=200son ~15 batches/segundo. Con 4 workers paralelos estás en el límite. - Para ChromaDB persistent: las escrituras concurrentes a SQLite están serializadas internamente, así que más workers no acelera más allá de 2-3.
# Concurrencia controlada
from concurrent.futures import ThreadPoolExecutor
with ThreadPoolExecutor(max_workers=4) as executor:
futures = [
executor.submit(ingest_batch, batch)
for batch in chunks(documents, batch_size=200)
]
for future in futures:
future.result() # awaits and re-raises
Ejercicio aplicado
Escenario: trabajas en una empresa de noticias que necesita ingerir 200,000 artículos al sistema RAG. Características:
- Cada artículo: ~800 tokens promedio (algunos hasta 4000)
- Modelo de embeddings: OpenAI
text-embedding-3-small - Tier de OpenAI: Tier 2 (5000 RPM, 5M tokens/minuto)
- RAM de la máquina de ingestion: 16 GB
- Restricción operativa: el ingestion no puede tardar más de 90 minutos (ventana de mantenimiento)
- Restricción de robustez: si el script crashea, se debe poder reanudar sin duplicar nada
Pregunta: diseña la configuración del pipeline (batch_size, paralelismo, sleep, idempotencia) y justifica con números cada decisión.
Solución
Cálculo de viabilidad:
total_docs = 200_000
avg_tokens_per_doc = 800
total_tokens = total_docs * avg_tokens_per_doc # 160M tokens
# Costo monetario
cost = (total_tokens / 1_000_000) * 0.02 # = $3.20
# Throughput requerido
target_minutes = 90
required_docs_per_minute = total_docs / target_minutes # 2,222 docs/min
required_tokens_per_minute = total_tokens / target_minutes # 1.78M tokens/min
print(f"Costo: ${cost}")
print(f"Throughput requerido: {required_docs_per_minute} docs/min")
print(f"Tokens/min requeridos: {required_tokens_per_minute:.0f}")
print(f"Tier 2 permite: 5,000,000 tokens/min → margen 2.8x")
Resultado:
- Costo: $3.20 (despreciable)
- Throughput requerido: 2,222 docs/min
- Tokens/min requeridos: 1.78M (tier 2 permite 5M, margen suficiente)
Diseño del pipeline:
1. batch_size = 100
Justificación:
- Algunos artículos tienen 4000 tokens. Con batch=100, el peor caso es 400K tokens por request, dentro del límite de 8191 por elemento individual (no por batch).
- Margen para hits de rate limit: 5000 RPM permite 50 batches/segundo, suficiente.
- En caso de fallo, perder 100 docs es manejable.
Verificación de tokens por batch:
avg_tokens_per_batch = 100 * 800 # 80,000 tokens — OK
worst_tokens_per_batch = 100 * 4000 # 400,000 tokens
# Por documento individual: 4000 < 8191 → OK
Si algún documento individual excede 8191 tokens, hay que truncarlo o chunkearlo (cubierto en M4/10) antes de la ingestion.
2. Paralelismo: 3 workers
Justificación:
- 3 workers × 100 docs/batch × ~2 batches/segundo (con OpenAI) ≈ 600 docs/segundo = 36,000 docs/minuto.
- Con 200,000 docs → ~5.5 minutos de procesamiento puro de embeddings.
- Tiempo total esperado: 6-10 minutos (incluyendo overhead, retries, escrituras a ChromaDB).
- Tier 2 permite ~83 RPS, 3 workers × 2 RPS = 6 RPS, dentro del límite con margen 13x.
3. sleep_between = 0.2 segundos
Justificación:
- A 2 batches/segundo por worker, con sleep=0.2 quedamos en ~1.7 batches/segundo. Dentro del rate limit cómodamente.
- Permite que ChromaDB persista las escrituras al disco sin saturar.
4. Idempotencia
IDs deterministas a partir del artículo:
def article_to_id(article):
"""ID estable basado en el ID del artículo en la fuente."""
return f"news_{article['source_id']}_{article['version']}"
Pipeline verifica collection.get(ids=batch_ids) antes de cada batch. Si todos los IDs ya existen (re-run después de crash), se omite el batch instantáneamente.
5. Manejo de fallos
config = {
"batch_size": 100,
"max_workers": 3,
"sleep_between": 0.2,
"max_retries": 3,
"retry_backoff": "exponential", # 5s, 10s, 20s
"log_failed_batches_to": "./failed_batches.json",
}
Failed batches se guardan a JSON con sus indices. Si fallan más de 5 batches, abortar pipeline para investigación. Si fallan 1-5 batches, dejar que termine y procesar los failed después manualmente.
6. Memoria
ram_per_batch = 100 * 1536 * 4 # 614 KB por batch
ram_total = ram_per_batch * 3 # 1.8 MB con 3 workers
# Trivial para máquina de 16 GB
Resumen de decisiones:
| Parámetro | Valor | Justificación |
|---|---|---|
batch_size | 100 | Margen para docs de 4000 tokens; recovery fácil |
max_workers | 3 | Throughput de 36K docs/min; 13x bajo el rate limit |
sleep_between | 0.2s | Evita saturar OpenAI y SQLite |
max_retries | 3 | 5s/10s/20s backoff cubre rate limits transitorios |
| Idempotencia | IDs deterministas | Reanudación segura sin duplicados |
| Tiempo estimado | 8-12 minutos | Bajo el límite de 90 min con 7-10x margen |
| Costo | $3.20 | Trivial vs valor del pipeline |
Bonus — observabilidad: instrumentar el pipeline con logs estructurados (docs_per_second, failed_batches, total_duration). Si los números futuros se alejan del baseline, sabes que algo cambió (modelo nuevo, dataset distinto, network).
Resumen y siguiente paso
Lo que aprendiste:
- Single inserts no escalan. El overhead constante por llamada multiplica linealmente con el número de documentos, transformando 100 docs (factible) en 1M docs (días).
- Batch ingestion amortiza el overhead. Con
batch_size=200y OpenAI, el throughput salta de ~10 docs/segundo a ~150 docs/segundo. - El
batch_sizecorrecto depende de tres restricciones: límites del API (tokens por request), RAM disponible y latencia de fallo aceptable. - Defaults razonables: 100-200 con OpenAI, 500 local con CPU, 1000+ con GPU. Pero siempre verifica los límites de tu modelo específico.
- Pipeline robusto requiere idempotencia (IDs deterministas + verificación previa) + retry con backoff + observabilidad (progress bar + logs por batch).
- Modos de falla más caros: duplicación por re-run sin idempotencia, OOM por batch demasiado grande, rate limits por concurrencia descontrolada.
Checkpoint: antes de avanzar, deberías poder:
- Calcular el tiempo total esperado de ingestion para N documentos con un
batch_sizedado. - Justificar la elección de
batch_sizepara un escenario específico citando los tres trade-offs (API, RAM, fallo). - Diseñar idempotencia para que un pipeline pueda reanudarse después de un crash sin duplicar.
Siguiente cápsula: 06 — Distance Metrics.
Acabas de aprender a meter datos al sistema eficientemente. Pero hay una decisión que tomaste sin pensar al crear las collections: la distance metric. Cosine, L2 o dot product — cada una asume una geometría distinta del espacio de embeddings y elegir mal degrada el retrieval en silencio. La cápsula 06 te da el modelo mental geométrico para elegir con criterio.
Recursos
- OpenAI — Rate Limits Guide — Límites por tier y estrategias de manejo
- OpenAI — Batching Embeddings — Recomendaciones oficiales de batching
- ChromaDB — Adding Documents — Operaciones de inserción y batch
- tqdm — Progress Bars in Python — Barras de progreso para pipelines
- Tenacity — Retry Library for Python — Alternativa robusta a try/except con retry manual
- Idempotency Patterns in Distributed Systems — Stripe explica el concepto general aplicable a APIs
Tiempo estimado: 30-40 minutos Siguiente: 06-distance-metrics.md