Patrones de Procesamiento por Lotes
Diseña pipelines robustos de procesamiento por lotes para grandes datasets con retry, idempotencia y observabilidad.
Visión General
El procesamiento por lotes es la columna vertebral de pipelines de datos, flujos de trabajo ETL y generación de reportes. A diferencia del procesamiento de streams, los trabajos por lotes procesan conjuntos de datos acotados en chunks, lo que los hace más simples de razonar pero requieren atención cuidadosa a la idempotencia, tolerancia a fallos y observabilidad.
Cuándo Usar
Usa este recurso cuando:
- Procesas grandes datasets que no caben en memoria. Consulta Retry Logic para manejar fallos transitorios.
- Construyes pipelines ETL para data warehouses
- Generas reportes o agregaciones nocturnas
- Migras datos entre sistemas con ventanas de mantenimiento
Solución
Pipeline Resiliente de Procesamiento por Lotes (Python)
import logging
from typing import Callable, List, Iterator
class BatchProcessor:
def __init__(self, batch_size: int = 1000, max_retries: int = 3):
self.batch_size = batch_size
self.max_retries = max_retries
self.processed = 0
self.failed = []
def process(
self,
items: Iterator[dict],
handler: Callable[[List[dict]], None]
) -> dict:
batch = []
for item in items:
batch.append(item)
if len(batch) >= self.batch_size:
self._execute(batch, handler)
batch = []
if batch:
self._execute(batch, handler)
return {"processed": self.processed, "failed": len(self.failed)}
def _execute(self, batch: List[dict], handler: Callable):
for attempt in range(self.max_retries):
try:
handler(batch)
self.processed += len(batch)
return
except Exception as e:
logging.warning(f"Batch fallido (intento {attempt + 1}): {e}")
if attempt == self.max_retries - 1:
self.failed.extend(batch)
Seguimiento Idempotente de Trabajos (SQL)
CREATE TABLE job_runs (
job_id VARCHAR(64) PRIMARY KEY,
started_at TIMESTAMP NOT NULL DEFAULT NOW(),
completed_at TIMESTAMP,
status VARCHAR(20) CHECK (status IN ('running', 'completed', 'failed')),
checksum VARCHAR(64)
);
-- Antes de comenzar, verifica si ya está completado
SELECT * FROM job_runs WHERE job_id = 'daily_report_2025_01_15' AND status = 'completed';
Explicación
Un pipeline de producción por lotes necesita tres propiedades:
- Idempotencia: Ejecutar el mismo trabajo dos veces debe producir el mismo resultado. Usa IDs de trabajo y checksums para saltar trabajo ya procesado. Consulta Endpoints Idempotentes para patrones de deduplicación.
- Tolerancia a fallos: Fallos individuales de batch no deben crashear todo el trabajo. Implementa reintentos con backoff exponencial y una cola de mensajes fallidos.
- Observabilidad: Rastrea progreso, throughput y errores. Registra métricas para items procesados, latencia y tasas de fallo.
Estrategia de chunking: Ajusta el tamaño de batches para balancear uso de memoria y throughput. Demasiado pequeño = overhead; demasiado grande = riesgo de OOM.
Variantes
| Patrón | Caso de Uso | Compromiso |
|---|---|---|
| Procesamiento por chunks | Archivos grandes, límites de memoria | Más simple, mayor latencia |
| Workers paralelos | Transformaciones CPU-bound | Complejo, necesita coordinación |
| MapReduce | Agregación distribuida | Escala horizontalmente |
| Change Data Capture | Sincronización incremental | Requiere soporte de la fuente |
Lo que funciona
- Diseña para idempotencia: Cada trabajo debe ser seguro de reintentar
- Registra todo: Inicio de trabajo, fin, y resultado de cada batch
- Usa transacciones: Envuelve escrituras de batch en transacciones de base de datos
- Monitorea profundidad de cola: Alerta cuando batches pendientes excedan umbrales
- Implementa circuit breakers: Detén reintentos si el downstream está unhealthy
Errores Comunes
- No manejar fallos parciales: Un batch de 1000 donde 1 falla necesita reintento individual
- Ignorar límites de memoria: Cargar datasets enteros en RAM crashea el proceso
- Faltar checkpointing: Un trabajo de 6 horas que falla a las 5:55 debe reiniciar desde cero
- Pérdida silenciosa de datos: Errores logueados pero no visibles para operadores
- Sin estrategia de rollback: Trabajos fallidos dejan la base de datos en estado inconsistente
Referencia Rápida
- Comando principal: ejecuta la solución base del artículo y verifica el resultado esperado.
- Validación: confirma que los tests pasan y que las métricas clave no se degradaron.
- Rollback: si algo falla, revierte el cambio y consulta la sección de Troubleshooting.
Lectura Adicional
- Documentación oficial: consulta la referencia actualizada del framework o herramienta utilizada.
- Guías relacionadas: explora las guías de batch-processing y data para profundizar.
- Patrones complementarios: revisa los patrones de diseño aplicables a tu stack tecnológico.
- Postmortems públicos: estudia incidentes reales de equipos que enfrentaron problemas similares en producción.
Notas de Producción
- Despliega gradualmente usando canary o blue-green para detectar regresiones temprano.
- Configura alertas para errores, latencia p99 y tasa de fallos antes de habilitar en producción.
- Documenta el rollback en el runbook; prueba el procedimiento en staging al menos una vez por trimestre.
- Revisa logs estructurados con correlation IDs para trazar requests end-to-end en incidentes.
Puntos Clave
- Aplica patrones de procesamiento por lotes cuando necesites una solución práctica para tu caso de uso.
- Monitorea el rendimiento después de implementar; mide latencia, errores y uso de recursos antes y después.
- Revisa la sección de Troubleshooting ante errores comunes; la mayoría tienen causa raíz documentada con solución.
- Mantén dependencias actualizadas y ejecuta tests en CI para prevenir regresiones en producción.
Ver También
- Database Indexing — optimización del rendimiento de queries para reads en batch
- Rendimiento Web — técnicas de rendimiento frontend y backend
- Load Testing — validación del rendimiento de batch jobs bajo carga
Troubleshooting
- Pipeline output does not match expectations: validate input schemas, intermediate states, and row counts at each step.
- Data quality degrades over time: add data validation checks and anomaly detection. Define SLIs for freshness, completeness, and accuracy.
- Job fails intermittently: look for race conditions, external dependencies, and resource contention. Retry with idempotency and bounded backoff.
- Schema changes break consumers: use schema registries and backward-compatible evolution.
- Storage costs grow unexpectedly: audit partition retention, compression, and duplicate copies. Archive cold data and set lifecycle policies.
Errores Comunes en Producción
- Copiar el ejemplo sin adaptarlo a volúmenes y modos de fallo reales.
- Saltar tests de carga e inyección de errores antes del primer despliegue productivo.
- Codificar valores fijos que deberían ser configurables por entorno.
- Olvidar agregar logging y monitoreo en cada paso.
- Desplegar sin plan de rollback ni estrategia de backup probada.
- Asumir que el ejemplo mínimo escalará sin agregar caché o procesamiento por lotes.
- No documentar la versión y configuración usadas en producción.
- Dejar la receta sin cambios cuando evolucionan las dependencias o la escala.
Preguntas frecuentes
¿Qué tan grande debería ser cada batch?
Comienza con 100-1000 items por batch. Haz benchmark con tus datos y restricciones de memoria. Para inserts en base de datos, 500-2000 rows por batch balancea throughput y tamaño de transacción. Para procesamiento de archivos, 10-50 MB por chunk evita OOM en contenedores de 512 MB. Monitorea el uso de memoria con resource.getrusage(resource.RUSAGE_SELF).ru_maxrss en Python. Si los batches exceden la memoria disponible, reduce el batch size o cambia a streaming. Para APIs network-bound, batches más grandes (500-5000) amortizan latencia. Para transforms CPU-bound, batches más pequeños (100-500) permiten mejor paralelismo. Siempre testea con volúmenes de datos similares a producción — los datos sintéticos raramente revelan patrones de presión de memoria.
¿Debería usar una cola de trabajos como Celery o un cron job?
Usa Celery/Redis para sistemas distribuidos con múltiples workers, retry logic y monitoreo. Celery provee task routing, priority queues y dead-letter handling out of the box. Para pipelines simples de un solo nodo, cron jobs bastan — pero añade flock para prevenir runs overlapping: flock -n /tmp/batch.lock python batch_job.py. Para Kubernetes, usa recursos CronJob con concurrencyPolicy: Forbid. Consulta Rate Limiting para controlar throughput. Para setups cloud-native, AWS Step Functions o Google Cloud Workflows proveen retry y state management manejados sin overhead de infraestructura.
¿Cómo manejo cambios de schema en medio del pipeline?
Versiona tu lógica de trabajo y schemas de datos. Ejecuta versiones viejas y nuevas en paralelo durante la migración. Usa una columna schema_version en las tablas target: ALTER TABLE orders ADD COLUMN schema_version INT DEFAULT 1. El batch job verifica schema_version y aplica la transformación correcta. Para Avro/Protobuf, usa schema registry para manejar cambios backward-compatible. Para JSON, valida con JSON Schema versionado por $id. Despliega nuevas versiones del job con un feature flag: if config.use_new_schema: transform_v2(record) else: transform_v1(record). Monitorea la salida de ambas versiones para discrepancias. Una vez confiado, descomisiona la versión vieja.
¿Cómo implemento checkpointing para batch jobs de larga duración?
El checkpointing guarda el progreso para que los jobs fallidos resuman desde el último batch exitoso. Almacena checkpoints en una tabla durable: CREATE TABLE job_checkpoints (job_id VARCHAR(64), batch_number INT, status VARCHAR(20), PRIMARY KEY (job_id, batch_number)). Después de cada batch, escribe: INSERT INTO job_checkpoints VALUES ('daily_etl', 42, 'completed'). En restart, query el último batch completado: SELECT MAX(batch_number) FROM job_checkpoints WHERE job_id = 'daily_etl' AND status = 'completed'. Resume desde el batch 43. Para checkpoints basados en archivos, escribe un archivo JSON a S3 o disco local después de cada batch. Usa writes atómicos para prevenir corrupción: escribe a un temp file, luego rename. Para jobs distribuidos, usa Redis o etcd para coordinación distribuida de checkpoints.
¿Cómo manejo dead-letter queues para items fallidos de batch?
Las dead-letter queues (DLQ) aislan items fallidos para inspección manual y retry. En Python, mantén una DLQ list: dead_letter_queue = [] y appendea items fallidos con contexto: dead_letter_queue.append({'item': item, 'error': str(e), 'batch_id': batch_id, 'timestamp': datetime.utcnow()}). Después de que el job completa, escribe la DLQ a una tabla de base de datos o message queue. En Celery, configura task_queues con un DLQ exchange dedicado. En AWS SQS, setea RedrivePolicy para mover mensajes después de maxReceiveCount intentos. Para DLQs basadas en base de datos: CREATE TABLE dead_letters (id SERIAL PRIMARY KEY, job_id VARCHAR(64), item JSONB, error TEXT, created_at TIMESTAMP DEFAULT NOW()). Procesa items de DLQ separadamente con un retry job que corre cada hora.
¿Cómo monitoro throughput y progreso de batch jobs?
Emite métricas en cada boundary de batch. En Python, usa prometheus_client: from prometheus_client import Counter, Histogram; processed = Counter('batch_processed_total', 'Items processed'); latency = Histogram('batch_duration_seconds', 'Batch duration'). Push métricas a Prometheus Pushgateway para batch jobs: from prometheus_client import push_to_gateway; push_to_gateway('localhost:9091', job='batch_etl', registry=registry). Para setups cloud, emite CloudWatch custom metrics o Datadog statsd. Trackea: items procesados, items fallidos, duración del batch, queue depth y uso de memoria. Setea dashboards de Grafana con alerts para: throughput abajo del 50% del average, failure rate arriba del 5% y duración del job excediendo SLA. Loguea JSON estructurado para cada batch: {"batch_id": 42, "processed": 1000, "failed": 3, "duration_ms": 1250}.
¿Cómo manejo backpressure en pipelines de procesamiento por lotes?
El backpressure ocurre cuando los sistemas downstream no pueden mantener el ritmo del productor de batches. Implementa rate limiting con un semaphore: from threading import Semaphore; rate_limiter = Semaphore(10) — acquire antes de cada batch write y release después. Para writes a base de datos, usa connection pooling con un max pool size para limitar inserts concurrentes. Para API calls, implementa un token bucket: import time; time.sleep(1 / max_rps). Monitorea la latencia downstream — si aumenta más de 2x el baseline, reduce el batch size o pausa el procesamiento. En Celery, setea worker_prefetch_multiplier = 1 para prevenir que los workers over-fetchen. Para pipelines basados en Kafka, setea max.poll.records para controlar el batch size. Usa concurrent.futures.ThreadPoolExecutor(max_workers=N) para limitar paralelismo en Python.
¿Cómo testeo batch jobs para correctitud?
Testea batch jobs con data fixtures determinísticos. Crea un dataset de test: test_data = [{'id': i, 'value': i * 2} for i in range(10000)]. Testea idempotencia corriendo el job dos veces y verificando que la salida sea idéntica: assert run_job(test_data) == run_job(test_data). Testea fault tolerance inyectando fallos: def failing_handler(batch): if len(batch) > 500: raise Exception('simulated failure') y verifica que el job retiene y registra fallos. Testea checkpointing matando el job mid-run y verificando que resume correctamente. Usa property-based testing con Hypothesis: @given(st.lists(st.dictionaries(st.text(), st.integers()))) para generar inputs edge-case. Para jobs basados en SQL, usa testcontainers para levantar una base de datos real: from testcontainers.postgres import PostgresContainer.
¿Cómo manejo procesamiento por lotes para datos de series temporales?
Los batch jobs de series temporales procesan datos en ventanas de tiempo fijas. Usa windowed batching: agrupa records por window_start y window_end timestamps. Por ejemplo, SELECT date_trunc('hour', timestamp) AS window, COUNT(*) FROM events GROUP BY 1 procesa batches por hora. Para datos que llegan tarde, usa watermarking: permite datos hasta 5 minutos tarde seteando watermark = current_time - 5 minutes. Procesa ventanas solo después de que el watermark pase. Para retención, particiona por tiempo: CREATE TABLE events_2025_01 PARTITION OF events FOR VALUES FROM ('2025-01-01') TO ('2025-02-01'). Dropea particiones viejas en lugar de borrar rows: DROP TABLE events_2024_01. Para downsampling, agrega datos crudos en tablas resumen de 1-minuto, 1-hora y 1-día en batch jobs separados.
¿Cómo manejo procesamiento por lotes con semántica exactly-once?
El procesamiento exactly-once requiere writes idempotentes y checkpoints transaccionales. Usa una transacción para escribir tanto la salida del batch como el checkpoint atómicamente: BEGIN; INSERT INTO results SELECT * FROM staging; INSERT INTO job_checkpoints VALUES ('job_1', 42, 'completed'); COMMIT;. Si la transacción falla, tanto los datos como el checkpoint se rollbackean — el job retiene el batch. Para consumers de Kafka, usa enable.idempotence=true y transactional.id para exactly-once del lado del producer. Para writes a base de datos, usa INSERT ... ON CONFLICT DO NOTHING para manejar batches duplicados safe. Para calls a APIs externas, usa idempotency keys: Idempotency-Key: batch_42_run_3 en headers de request. Acepta que exactly-once tiene un coste de performance — mide si at-least-once con writes idempotentes basta para tu use case.
Recursos Relacionados
Caché y Memoización en Python, JavaScript y Java
Cómo cachear computaciones costosas y respuestas de API usando caches en memoria, LRU, TTL y distribuidos en Python, JavaScript y Java.
RecipeValidar y Sanitizar Datos de Input de Usuario
Cómo validar, sanitizar y restringir datos de input de usuario en el boundary de aplicación usando schemas, type checking y librerías de validación.
RecipeFormateo de Fechas
Cómo parsear, formatear y manipular fechas a través de timezones usando Python, JavaScript y Java.
RecipeDeep Clone en JavaScript: structuredClone vs lodash vs JSON
Compará métodos de clonación profunda en JavaScript, Python y Java. Creá copias independientes de objetos y arrays, manejá referencias circulares, Dates, Maps, Sets y typed arrays, y elegí el enfoque correcto con una matriz de decisión.
RecipeAplanar y Reconstruir Objetos con Python, JS y Java
Cómo convertir objetos anidados en pares clave-valor planos y reconstruirlos, con soporte de notación por puntos, corchetes y separadores personalizados.
RecipeDeep Clone de Objetos en JavaScript: Mas alla de JSON.parse
Compara estrategias de deep clone incluyendo JSON.parse, structuredClone, recursion manual y librerias para copiar objetos anidados con referencias circulares