Arquitecturas de datos: cómo garantizar idempotencia cuando el pipeline reintenta llamadas a APIs externas
Arquitecturas de datos: cómo garantizar idempotencia cuando el pipeline reintenta llamadas a APIs externas
Un pipeline cae a mitad de ejecución. El scheduler lo relanza. La llamada a la API externa se repite. El registro llega dos veces al almacén. Nadie lo nota hasta que el análisis de tendencias devuelve volúmenes inflados un 40 % y el equipo de negocio toma decisiones sobre datos duplicados.
Este escenario no es un accidente raro. Es la consecuencia predecible de construir pipelines que consumen APIs externas sin diseñar idempotencia desde el primer sprint. La tolerancia a fallos y el backpressure tienen su espacio en la arquitectura, pero sin idempotencia todo lo demás es decoración: puedes reintentar perfectamente y aun así corromper el estado final.
El problema no está en el reintento. Está en que el reintento no sabe qué ya ocurrió.
Por qué las APIs externas amplifican el problema
Cuando el pipeline es completamente interno, la idempotencia es difícil pero controlable: tienes acceso al estado, a los logs de transacción, a los mecanismos de commit. Cuando dependes de una API externa, pierdes parte de ese control.
Las APIs de datos del universo público —señales de medios, menciones, análisis derivados— no están obligadas a garantizar que una respuesta ya procesada llegue marcada como tal. Devuelven lo que hay en ese momento. Si tu pipeline pregunta dos veces por el mismo rango temporal, recibirá dos respuestas potencialmente solapadas.
Tres patrones agravan esto:
- Reinicios automáticos sin checkpoint: el orquestador (Airflow, Prefect, Temporal) relanza el DAG desde el inicio porque no hay marcas de progreso persistentes.
- Ventanas de tiempo abiertas: el pipeline define rangos de consulta dinámicos (
now - 15m) en lugar de rangos fijos anclados a un ID de ejecución. - Almacenes sin constraint de unicidad: si el destino (BigQuery, S3, Elasticsearch) no rechaza duplicados, el problema se acumula en silencio.
El resultado es un dataset que parece correcto en volumen pero está contaminado en contenido.
La clave primaria como contrato arquitectónico
El primer mecanismo de idempotencia no es técnico: es conceptual. Antes de escribir una línea de código, tienes que responder: ¿qué hace que un registro sea único en mi dominio?
Para pipelines que consumen APIs de análisis de datos, la clave primaria suele ser una combinación de:
- Identificador del documento o señal en el proveedor (
doc_id,mention_id). - Timestamp de publicación normalizado (no de ingestión).
- Fuente o canal de origen.
Con esa clave definida, el almacén puede rechazar o ignorar inserciones duplicadas. En PostgreSQL basta con INSERT ... ON CONFLICT DO NOTHING. En BigQuery, una tabla particionada con MERGE sobre la clave compuesta. En Elasticsearch, usar el _id del proveedor como identificador del documento en lugar de dejar que Elasticsearch genere uno automático.
El error más común es usar el timestamp de ingestión como parte de la clave. Cada reintento genera un timestamp diferente y el conflicto nunca se detecta.
Checkpoints de progreso: el pipeline debe recordar dónde estaba
La idempotencia en el almacén es necesaria pero no suficiente. Si el pipeline reprocesa miles de llamadas a la API para llegar al mismo resultado, el coste de cómputo y de llamadas se dispara. En modelos pay-as-you-go, cada reintento innecesario tiene un coste directo.
La solución es externalizar el estado de progreso del pipeline. Opciones según el nivel de madurez:
Nivel 1 — Tabla de control en la base de datos: una tabla pipeline_runs con columnas (run_id, page_cursor, status, updated_at). Antes de cada llamada a la API, el worker persiste el cursor. Si cae y se reinicia, lee el cursor y continúa desde ahí.
Nivel 2 — Offset en el broker de mensajes: si el pipeline usa Kafka o Kinesis, el offset del consumer es el checkpoint. El problema es que la API externa no es el broker; el checkpoint cubre la cola interna pero no la posición en la API.
Nivel 3 — Cursores de API como estado inmutable: algunas APIs exponen cursores de paginación opacos que son estables entre sesiones. Si el proveedor los soporta, persiste el cursor en Redis o DynamoDB con TTL. Al reiniciar, el pipeline reanuda desde ese cursor exacto, sin repetir páginas anteriores.
FeedScale, por ejemplo, expone parámetros de paginación reproducibles que permiten anclar la consulta a un estado determinado del índice. Usarlos como checkpoint reduce el coste de reintento a cero páginas duplicadas.
Ventanas de tiempo fijas frente a ventanas deslizantes
Un antipatrón frecuente en pipelines de análisis de señales es definir la ventana de consulta de forma relativa: "dame los documentos de las últimas dos horas". Si el pipeline se ejecuta tres veces (porque falló dos), obtiene tres rangos ligeramente distintos con solapamiento no determinista.
La alternativa es anclar la ventana al ID de ejecución o al slot temporal del scheduler:
# Antipatrón: ventana relativa al momento de ejecución
end_time = datetime.utcnow()
start_time = end_time - timedelta(hours=2)
# Patrón correcto: ventana anclada al slot del scheduler
slot = execution_date # parámetro inyectado por Airflow/Prefect
start_time = slot
end_time = slot + timedelta(hours=2)
Con ventanas fijas, el reintento del slot 2026-08-23T10:00:00Z siempre consultará exactamente el mismo rango. El resultado puede diferir si el proveedor actualizó su índice, pero el rango de consulta es idéntico. Eso facilita la detección y resolución de duplicados en el almacén.
Validación post-ingestión: la red de seguridad que nadie construye
Incluso con claves primarias correctas, checkpoints y ventanas fijas, los duplicados pueden colarse. Proveedores que reindexan documentos con timestamps modificados, cambios de zona horaria sin aviso, cursores expirados que obligan a reiniciar desde el principio.
Un job de validación post-ingestión, ejecutado con baja frecuencia (diario o semanal), añade una capa de seguridad sin coste operativo elevado:
-- Detectar duplicados por clave compuesta
SELECT doc_id, source, DATE(published_at), COUNT(*) AS occurrences
FROM signals
GROUP BY 1, 2, 3
HAVING COUNT(*) > 1
ORDER BY occurrences DESC
LIMIT 100;
Este job no bloquea el pipeline principal. Corre en paralelo, alimenta un dashboard de calidad y dispara alertas solo cuando el porcentaje de duplicados supera un umbral definido (por ejemplo, 0,5 % del volumen diario).
La idempotencia no es un detalle de implementación
Equipos que tratan la idempotencia como un refinamiento posterior descubren, meses después, que su almacén histórico está contaminado y que reconstruirlo implica repetir semanas de procesamiento con coste real en llamadas a APIs externas y tiempo de cómputo.
Diseñar idempotencia desde el inicio —clave primaria clara, checkpoint persistente, ventanas fijas, validación periódica— no añade complejidad estructural. Añade disciplina. Y en pipelines que dependen de APIs externas, esa disciplina es lo que separa un sistema que escala de uno que colapsa cuando el scheduler decide reintentar en el peor momento posible.