El problema fundamental que aborda esta propuesta es la complejidad inherente y la fragilidad de la orquestación de workflows distribuidos, especialmente en el contexto de cargas de trabajo de IA que requieren durabilidad y visibilidad. Las soluciones actuales, que a menudo se basan en coordinadores externos y sistemas de colas separados, introducen puntos de fallo adicionales, latencia y un alto costo operativo y de desarrollo. La tesis central es que los 'workflows son datos' y, por lo tanto, una base de datos existente, optimizada durante décadas para la gestión de datos transaccionales, es el lugar más adecuado y eficiente para almacenar y gestionar el estado de ejecución de los workflows.
Este enfoque busca simplificar la arquitectura de sistemas distribuidos, reducir la superficie de ataque y mejorar la fiabilidad al consolidar la lógica de orquestación dentro de la misma capa de persistencia de datos de la aplicación. La necesidad de durabilidad y ejecución 'exactly-once' se ha vuelto crítica con la creciente dependencia de la IA en procesos de negocio sensibles, donde los errores son costosos y la paciencia del usuario es limitada.
Arquitectura del Sistema
La arquitectura propuesta se centra en una librería 'in-process' que se integra directamente en el código de la aplicación, utilizando una base de datos relacional (como PostgreSQL) como el único componente de estado compartido para la orquestación. No hay servicios de orquestación externos, colas separadas o pools de workers dedicados.
Los workflows y sus pasos son funciones ordinarias en el código de la aplicación. La librería envuelve estas funciones para implementar la durabilidad. Cada workflow tiene un 'workflow ID' único (similar a una clave de idempotencia) y su estado se registra en una tabla workflow_status en la base de datos. Esta tabla almacena metadatos como el ID, nombre, estado (pendiente, completado, fallido), inputs y outputs del workflow. Cada paso individual dentro de un workflow también registra su output en una tabla step_outputs, con una clave foránea que referencia el workflow ID principal. Esto permite el 'checkpointing' del estado de ejecución.
En caso de fallo (ej. caída del servidor, OOM), la recuperación se logra listando los workflows pendientes de la base de datos. La reejecución de un workflow recuperado carga los inputs y el puntero a la función del workflow. Gracias a los checkpoints de los pasos, la librería puede 'saltar' los pasos ya completados y reanudar la ejecución desde el último paso exitoso, garantizando la semántica 'exactly-once'. Para la gestión de colas, la misma tabla workflow_status se extiende con columnas como queue_name, created_at y priority. Los workers realizan un 'polling' de esta tabla usando SELECT ... FOR UPDATE SKIP LOCKED para adquirir tareas de forma concurrente sin contención de bloqueos. La programación cron descentralizada se implementa haciendo que cada worker intente encolar el mismo workflow con un workflow ID basado en el tiempo programado y un 'jitter' aleatorio, aprovechando la restricción de clave primaria para asegurar una única inserción exitosa y evitando picos de carga en la base de datos.
Ejecución de Workflow con Checkpointing
- 1 Aplicación Invoca función de workflow envuelta
- 2 Librería Transact Genera Workflow ID, registra inputs en `workflow_status` (estado PENDING)
- 3 Aplicación Ejecuta funciones de paso envueltas
- 4 Librería Transact Para cada paso: verifica checkpoint, ejecuta, registra output en `step_outputs`
- 5 Aplicación Workflow finaliza
- 6 Librería Transact Registra output final en `workflow_status` (estado COMPLETED)
Recuperación y Reejecución de Workflow
- 1 Worker de Recuperación Lista workflows PENDING de `workflow_status`
- 2 Worker de Recuperación Carga inputs y definición de función del workflow
- 3 Librería Transact Reejecuta workflow con ID original e inputs
- 4 Librería Transact Para cada paso: si hay checkpoint en `step_outputs`, salta y usa resultado gu...
- 5 Librería Transact Si no hay checkpoint, ejecuta el paso
- 6 Librería Transact Continúa hasta el final o el último paso no completado
| Capa | Tecnología | Justificación |
|---|---|---|
| storage | PostgreSQL | Base de datos principal para almacenar el estado de los workflows, inputs, outputs de pasos y metadatos. Actúa como el único punto de verdad y coordinador de durabilidad. vs Otras bases de datos relacionales (MySQL, SQL Server), Sistemas de colas dedicados (RabbitMQ, Kafka), Orquestadores de workflows externos (AWS Step Functions, Temporal) Uso de `FOR UPDATE SKIP LOCKED` para gestión de colas concurrente. |
| compute | Python, TypeScript, Go, Java | Lenguajes de programación para implementar la lógica de negocio y la librería 'in-process' de orquestación. La librería se integra directamente en el runtime de la aplicación. |
Trade-offs
Ganancias
- ▲ Simplicidad arquitectónica
- ▲▲ Reducción de latencia
- ▲ Costo operacional
- ▲ Visibilidad y auditabilidad
- ▲ Experiencia de desarrollador
- ▲ Integración con ecosistema existente
Costes
- ▲ Acoplamiento a la base de datos
- △ Escalabilidad de escritura limitada por la DB
- ▲ Necesidad de implementar librerías por lenguaje
- ▲ Serialización/deserialización cross-lenguaje
SELECT workflow_id, inputs FROM workflow_status
WHERE status = 'ENQUEUED'
ORDER BY created_at ASC
FOR UPDATE SKIP LOCKED
LIMIT 1;function run_workflow_wrapper(workflow_func, inputs):
workflow_id = generate_unique_id()
store_workflow_status(workflow_id, PENDING, inputs)
try:
result = workflow_func(inputs)
store_workflow_status(workflow_id, COMPLETED, result)
except Exception as e:
store_workflow_status(workflow_id, FAILED, error=e)
return resultfunction run_step_wrapper(step_func, step_id, workflow_id, step_inputs):
if checkpoint_exists(workflow_id, step_id):
return load_checkpoint(workflow_id, step_id)
try:
result = step_func(step_inputs)
store_step_output(workflow_id, step_id, result)
return result
except Exception as e:
// Handle retry logic or mark step as failed
throw eFundamentos Teóricos
Este enfoque resuena con principios fundamentales de los sistemas distribuidos y la gestión de transacciones. La idea de 'workflows como datos' y el uso de la base de datos como un registro de eventos duradero se alinea con el concepto de Write-Ahead Log (WAL) y los sistemas de gestión de transacciones ACID, que han sido estudiados extensamente desde los años 70 y 80. La durabilidad y la atomicidad de las operaciones son pilares de la computación transaccional, y este sistema las extiende al nivel de la lógica de negocio.
La garantía de 'exactly-once execution' para los workflows, aunque dependiente de la idempotencia de los pasos individuales, se basa en la capacidad de la base de datos para mantener un estado consistente y recuperable. Esto se relaciona con los conceptos de 'determinismo' en la reejecución y 'idempotencia' en las operaciones, esenciales para la tolerancia a fallos en sistemas distribuidos. El uso de FOR UPDATE SKIP LOCKED para la gestión de colas es una aplicación directa de mecanismos de bloqueo de bases de datos para resolver problemas de concurrencia en sistemas de 'message passing' o 'task processing', un patrón bien conocido en la literatura de bases de datos y sistemas operativos.