← Volver al catálogo
WebReferenciaAvanzadoEn el pase

Orquestacion Saga para Transacciones Distribuidas

Guia de produccion para implementar el patron Saga en microservicios: coordinacion de transacciones distribuidas sin 2PC, acciones compensatorias, maquinas de estado y recuperacion ante fallos.

Comprobando acceso…

Incluida en el Pase · para Python, Kafka, RabbitMQ

// resultado_de_ejemplo

""" CultivaEdu — Saga de Matriculación con Patrón Orquestador

Implementación producción-grade del patrón Saga para coordinar la matriculación de alumnos en cursos de IA corporativa a través de 4 microservicios independientes sin bloqueos distribuidos (sin 2PC).

Stack: Python 3.12 · RabbitMQ 3.13 · PostgreSQL · Prometheus """

from future import annotations

import asyncio import logging import time import uuid from dataclasses import dataclass, field from datetime import datetime from enum import Enum from typing import Any, Dict, List, Optional

──────────────────────────────────────────────────────────────────────────────

Métricas Prometheus (stub compatible sin dependencia externa)

En producción: pip install prometheus-client y usar las clases reales

──────────────────────────────────────────────────────────────────────────────

class MetricStub: """Stub que reproduce la API de prometheus_client sin dependencia externa.""" def labels(self, **): return self def inc(self, amount=1): pass def observe(self, val): pass

SAGA_STARTED = _MetricStub() SAGA_COMPLETED = _MetricStub() SAGA_FAILED = _MetricStub() SAGA_DURATION = _MetricStub() SAGA_STUCK = _MetricStub() STEP_ERRORS = _MetricStub()

──────────────────────────────────────────────────────────────────────────────

Tipos base

──────────────────────────────────────────────────────────────────────────────

class SagaState(str, Enum): STARTED = "STARTED" # saga iniciada, primer paso en cola PENDING = "PENDING" # esperando respuesta de un participante COMPENSATING = "COMPENSATING" # un paso falló; revirtiendo pasos completados COMPLETED = "COMPLETED" # todos los pasos forward completados FAILED = "FAILED" # saga fallida y todas las compensaciones terminadas

@dataclass class SagaStep: name: str action: str # comando a enviar al participante (forward) compensation: str # comando para revertir si hay fallo posterior timeout_s: float # SLA máximo de este paso en segundos result: Optional[Dict] = None

@dataclass class SagaRecord: saga_id: str saga_type: str state: SagaState current_step: int data: Dict steps: List[SagaStep] created_at: datetime = field(default_factory=datetime.utcnow) updated_at: datetime = field(default_factory=datetime.utcnow) error: Optional[str] = None

──────────────────────────────────────────────────────────────────────────────

Repositorio de estado (mock en memoria; en producción → PostgreSQL)

──────────────────────────────────────────────────────────────────────────────

class SagaStore: """Simula la tabla saga_records en PostgreSQL."""

def __init__(self):
    self._store: Dict[str, SagaRecord] = {}

async def save(self, record: SagaRecord) -> None:
    record.updated_at = datetime.utcnow()
    self._store[record.saga_id] = record

async def load(self, saga_id: str) -> Optional[SagaRecord]:
    return self._store.get(saga_id)

async def find_stuck(self, older_than_s: float = 300) -> List[SagaRecord]:
    now = time.time()
    return [
        r for r in self._store.values()
        if r.state in (SagaState.PENDING, SagaState.COMPENSATING)
        and (now - r.updated_at.timestamp()) > older_than_s
    ]

──────────────────────────────────────────────────────────────────────────────

Bus de mensajes (mock; en producción → RabbitMQ exchange cultiva.saga)

──────────────────────────────────────────────────────────────────────────────

class EventBus: """Simula RabbitMQ con exchange cultiva.saga."""

def __init__(self):
    self._log: List[Dict] = []
    self._handlers: Dict[str, Any] = {}

async def publish(self, routing_key: str, payload: Dict) -> None:
    self._log.append({"event": routing_key, "payload": payload, "ts": datetime.utcnow().isoformat()})
    logging.info(f"[RabbitMQ] → {routing_key}: saga_id={payload.get('saga_id','?')}")
    # Entrega inmediata al handler registrado (simula consumidor)
    handler = self._handlers.get(routing_key)
    if handler:
        await handler(payload)

def subscribe(self, routing_key: str, handler) -> None:
    self._handlers[routing_key] = handler

def event_log(self) -> List[Dict]:
    return self._log

──────────────────────────────────────────────────────────────────────────────

Orquestador abstracto base

──────────────────────────────────────────────────────────────────────────────

class SagaOrchestrator: """ Orquestador base para el patrón Saga.

Subclasificar e implementar:
  - saga_type  → identificador único de tipo
  - define_steps(data) → lista ordenada de SagaStep
"""

def __init__(self, store: SagaStore, bus: EventBus):
    self.store = store
    self.bus   = bus
    # Suscribirse a respuestas de participantes
    self.bus.subscribe("SagaStepCompleted",        self._on_step_completed)
    self.bus.subscribe("SagaStepFailed",           self._on_step_failed)
    self.bus.subscribe("SagaCompensationCompleted", self._on_compensation_completed)

@property
def saga_type(self) -> str:
    raise NotImplementedError

def define_steps(self, data: Dict) -> List[SagaStep]:
    raise NotImplementedError

# ── Ciclo de vida ────────────────────────────────────────────────────────

async def start(self, data: Dict) -> str:
    saga_id = str(uuid.uuid4())
    steps   = self.define_steps(data)
    record  = SagaRecord(
        saga_id=saga_id,
        saga_type=self.saga_type,
        state=SagaState.STARTED,
        current_step=0,
        data=data,
        steps=steps,
    )
    await self.store.save(record)
    SAGA_STARTED.labels(saga_type=self.saga_type).inc()
    logging.info(f"[SAGA] {self.saga_type} iniciada → saga_id={saga_id}")
    await self._dispatch_step(record)
    return saga_id

async def _dispatch_step(self, record: SagaRecord) -> None:
    step = record.steps[record.current_step]
    record.state = SagaState.PENDING
    await self.store.save(record)
    logging.info(f"[SAGA] Paso {record.current_step} '{step.name}' → {step.action}")
    await self.bus.publish(step.action, {
        "saga_id":   record.saga_id,
        "step_name": step.name,
        "timeout_s": step.timeout_s,
        **record.data,
    })

async def _on_step_completed(self, event: Dict) -> None:
    record = await self.store.load(event["saga_id"])
    if not record or record.state != SagaState.PENDING:
        return
    step = record.steps[record.current_step]
    if step.name != event["step_name"]:
        return
    # Guardar resultado del paso para la compensación posterior
    step.result = event.get("result", {})
    record.current_step += 1
    if record.current_step >= len(record.steps):
        record.state = SagaState.COMPLETED
        await self.store.save(record)
        SAGA_COMPLETED.labels(saga_type=record.saga_type).inc()
        duration = (datetime.utcnow() - record.created_at).total_seconds()
        SAGA_DURATION.labels(saga_type=record.saga_type).observe(duration)
        logging.info(f"[SAGA] ✓ COMPLETADA saga_id={record.saga_id} ({duration:.2f}s)")
    else:
        await self._dispatch_step(record)

async def _on_step_failed(self, event: Dict) -> None:
    record = await self.store.load(event["saga_id"])
    if not record:
        return
    step_name = event["step_name"]
    error     = event.get("error", "desconocido")
    logging.warning(f"[SAGA] Paso '{step_name}' fallido: {error} — iniciando compensación")
    STEP_ERRORS.labels(saga_type=record.saga_type, step=step_name).inc()
    record.state = SagaState.COMPENSATING
    record.error = error
    await self.store.save(record)
    await self._compensate(record)

async def _compensate(self, record: SagaRecord) -> None:
    """Ejecuta compensaciones en orden inverso para los pasos ya completados."""
    for i in range(record.current_step - 1, -1, -1):
        step = record.steps[i]
        if step.result is not None:
            logging.info(f"[SAGA] ↩ Compensando paso {i} '{step.name}' → {step.compensation}")
            await self.bus.publish(step.compensation, {
                "saga_id":         record.saga_id,
                "step_name":       step.name,
                "original_result": step.result,
            })

async def _on_compensation_completed(self, event: Dict) -> None:
    record = await self.store.load(event["saga_id"])
    if not record or record.state != SagaState.COMPENSATING:
        return
    # Verificar si todas las compensaciones terminaron
    compensations_pending = any(
        s.result is not None
        for s in record.steps[:record.current_step]
        if s.compensation not in [
            e["payload"]["step_name"]
            for e in self.bus.event_log()
            if e["event"] == "SagaCompensationCompleted"
               and e["payload"]["saga_id"] == record.saga_id
        ]
    )
    if not compensations_pending:
        record.state = SagaState.FAILED
        await self.store.save(record)
        SAGA_FAILED.labels(saga_type=record.saga_type).inc()
        logging.error(f"[SAGA] ✗ FALLIDA saga_id={record.saga_id} — {record.error}")

──────────────────────────────────────────────────────────────────────────────

Saga concreta: MatriculacionSaga

──────────────────────────────────────────────────────────────────────────────

class MatriculacionSaga(SagaOrchestrator): """ Saga que coordina la matriculación de un alumno en CultivaEdu:

  1. CursosService.ReservarPlaza       (compensación: LiberarPlaza)
  2. PagosService.ProcesarPago         (compensación: ReembolsarPago)
  3. CertificacionService.CrearExpediente (compensación: EliminarExpediente)
  4. NotificacionService.EnviarBienvenida (compensación: EnviarCancelacion)
"""

@property
def saga_type(self) -> str:
    return "Matriculacion"

def define_steps(self, data: Dict) -> List[SagaStep]:
    return [
        SagaStep(
            name="reservar_plaza",
            action="CursosService.ReservarPlaza",
            compensation="CursosService.LiberarPlaza",
            timeout_s=5,
        ),
        SagaStep(
            name="procesar_pago",
            action="PagosService.ProcesarPago",
            compensation="PagosService.ReembolsarPago",
            timeout_s=30,
        ),
        SagaStep(
            name="crear_expediente",
            action="CertificacionService.CrearExpediente",
            compensation="CertificacionService.EliminarExpediente",
            timeout_s=10,
        ),
        SagaStep(
            name="enviar_bienvenida",
            action="NotificacionService.EnviarBienvenida",
            compensation="NotificacionService.EnviarCancelacion",
            timeout_s=5,
        ),
    ]

──────────────────────────────────────────────────────────────────────────────

Participantes (servicios consumidores, handlers de comandos)

──────────────────────────────────────────────────────────────────────────────

class CursosService: def init(self, bus: EventBus, fail_on: Optional[str] = None): self.bus = bus self.fail_on = fail_on self._reservas: Dict[str, str] = {} bus.subscribe("CursosService.ReservarPlaza", self.reservar_plaza) bus.subscribe("CursosService.LiberarPlaza", self.liberar_plaza)

async def reservar_plaza(self, cmd: Dict) -> None:
    if self.fail_on == "reservar_plaza":
        await self.bus.publish("SagaStepFailed", {
            "saga_id":   cmd["saga_id"],
            "step_name": "reservar_plaza",
            "error":     "Sin plazas disponibles en curso-ia-fundamentos-v3",
        })
        return
    reserva_id = f"res-{uuid.uuid4().hex[:8]}"
    self._reservas[cmd["saga_id"]] = reserva_id
    logging.info(f"  [CursosService] Plaza reservada → {reserva_id}")
    await self.bus.publish("SagaStepCompleted", {
        "saga_id":   cmd["saga_id"],
        "step_name": "reservar_plaza",
        "result":    {"reserva_id": reserva_id, "curso_id": cmd.get("curso_id")},
    })

async def liberar_plaza(self, cmd: Dict) -> None:
    reserva_id = cmd["original_result"].get("reserva_id", "?")
    logging.info(f"  [CursosService] ↩ Plaza liberada → {reserva_id}")
    await self.bus.publish("SagaCompensationCompleted", {
        "saga_id":   cmd["saga_id"],
        "step_name": "reservar_plaza",
    })

class PagosService: def init(self, bus: EventBus, fail_on: Optional[str] = None): self.bus = bus self.fail_on = fail_on bus.subscribe("PagosService.ProcesarPago", self.procesar_pago) bus.subscribe("PagosService.ReembolsarPago", self.reembolsar_pago)

async def procesar_pago(self, cmd: Dict) -> None:
    if self.fail_on == "procesar_pago":
        await self.bus.publish("SagaStepFailed", {
            "saga_id":   cmd["saga_id"],
            "step_name": "procesar_pago",
            "error":     "Tarjeta rechazada (stripe: insufficient_funds)",
        })
        return
    pago_id = f"pi-{uuid.uuid4().hex[:12]}"
    logging.info(f"  [PagosService] Pago procesado → {pago_id} ({cmd.get('importe_eur')}€)")
    await self.bus.publish("SagaStepCompleted", {
        "saga_id":   cmd["saga_id"],
        "step_name": "procesar_pago",
        "result":    {"pago_id": pago_id, "importe_eur": cmd.get("importe_eur")},
    })

async def reembolsar_pago(self, cmd: Dict) -> None:
    pago_id = cmd["original_result"].get("pago_id", "?")
    logging.info(f"  [PagosService] ↩ Reembolso emitido → {pago_id}")
    await self.bus.publish("SagaCompensationCompleted", {
        "saga_id":   cmd["saga_id"],
        "step_name": "procesar_pago",
    })

class CertificacionService: def init(self, bus: EventBus, fail_on: Optional[str] = None): self.bus = bus self.fail_on = fail_on bus.subscribe("CertificacionService.CrearExpediente", self.crear_expediente) bus.subscribe("CertificacionService.EliminarExpediente", self.eliminar_expediente)

async def crear_expediente(self, cmd: Dict) -> None:
    if self.fail_on == "crear_expediente":
        await self.bus.publish("SagaStepFailed", {
            "saga_id":   cmd["saga_id"],
            "step_name": "crear_expediente",
            "error":     "Alumno ya tiene expediente activo en este curso",
        })
        return
    expediente_id = f"exp-{uuid.uuid4().hex[:10]}"
    logging.info(f"  [CertificacionService] Expediente creado → {expediente_id}")
    await self.bus.publish("SagaStepCompleted", {
        "saga_id":   cmd["saga_id"],
        "step_name": "crear_expediente",
        "result":    {"expediente_id": expediente_id},
    })

async def eliminar_expediente(self, cmd: Dict) -> None:
    exp_id = cmd["original_result"].get("expediente_id", "?")
    logging.info(f"  [CertificacionService] ↩ Expediente eliminado → {exp_id}")
    await self.bus.publish("SagaCompensationCompleted", {
        "saga_id":   cmd["saga_id"],
        "step_name": "crear_expediente",
    })

class NotificacionService: def init(self, bus: EventBus, fail_on: Optional[str] = None): self.bus = bus self.fail_on = fail_on bus.subscribe("NotificacionService.EnviarBienvenida", self.enviar_bienvenida) bus.subscribe("NotificacionService.EnviarCancelacion", self.enviar_cancelacion)

async def enviar_bienvenida(self, cmd: Dict) -> None:
    logging.info(
        f"  [NotificacionService] Email bienvenida → {cmd.get('alumno_email')} "
        f"+ HR {cmd.get('gestor_hr_email')}"
    )
    await self.bus.publish("SagaStepCompleted", {
        "saga_id":   cmd["saga_id"],
        "step_name": "enviar_bienvenida",
        "result":    {"emails_enviados": [cmd.get("alumno_email"), cmd.get("gestor_hr_email")]},
    })

async def enviar_cancelacion(self, cmd: Dict) -> None:
    logging.info(f"  [NotificacionService] ↩ Email cancelación enviado")
    await self.bus.publish("SagaCompensationCompleted", {
        "saga_id":   cmd["saga_id"],
        "step_name": "enviar_bienvenida",
    })

──────────────────────────────────────────────────────────────────────────────

Escenarios de demostración

──────────────────────────────────────────────────────────────────────────────

INPUT_MATRICULACION = { "empresa_id": "emp-0042", "alumno_id": "usr-8891", "alumno_nombre": "María García", "alumno_email": "maria.garcia@techcorp.es", "curso_id": "curso-ia-fundamentos-v3", "metodo_pago": "stripe_pm_abc123", "importe_eur": 299.00, "gestor_hr_email": "rrhh@techcorp.es", }

logging.basicConfig(level=logging.INFO, format="%(message)s")

async def escenario_exito() -> None: """Saga que completa los 4 pasos sin errores.""" print("\n" + "="*60) print("ESCENARIO 1: Matriculación exitosa") print("="*60) bus = EventBus() store = SagaStore() CursosService(bus) PagosService(bus) CertificacionService(bus) NotificacionService(bus) saga = MatriculacionSaga(store, bus) saga_id = await saga.start(INPUT_MATRICULACION) record = await store.load(saga_id) print(f"\nResultado final: {record.state.value}")

async def escenario_fallo_pago() -> None: """Pago rechazado — debe compensar (liberar plaza).""" print("\n" + "="*60) print("ESCENARIO 2: Fallo en pago → compensación automática") print("="*60) bus = EventBus() store = SagaStore() CursosService(bus) # reserva exitosa PagosService(bus, fail_on="procesar_pago") # pago falla CertificacionService(bus) NotificacionService(bus) saga = MatriculacionSaga(store, bus) saga_id = await saga.start(INPUT_MATRICULACION) record = await store.load(saga_id) print(f"\nResultado final: {record.state.value} — {record.error}")

async def escenario_fallo_expediente() -> None: """Error al crear expediente — compensa pago Y plaza.""" print("\n" + "="*60) print("ESCENARIO 3: Fallo en expediente → compensación en cadena") print("="*60) bus = EventBus() store = SagaStore() CursosService(bus) PagosService(bus) CertificacionService(bus, fail_on="crear_expediente") # falla NotificacionService(bus) saga = MatriculacionSaga(store, bus) saga_id = await saga.start(INPUT_MATRICULACION) record = await store.load(saga_id) print(f"\nResultado final: {record.state.value} — {record.error}")

async def main() -> None: await escenario_exito() await escenario_fallo_pago() await escenario_fallo_expediente() print("\n" + "="*60) print("Saga Pattern implementado correctamente en CultivaEdu.") print("="*60)

if name == "main": asyncio.run(main())

// qué_hace

Implementa el patron Saga para gestionar transacciones distribuidas de larga duracion a traves de multiples microservicios sin bloqueos distribuidos.

// cómo_lo_hace

Define pasos ordenados con comandos de accion y compensacion, configura timeouts por paso, gestiona maquinas de estado y provee monitoring con Prometheus y recuperacion via colas de mensajes muertos (DLQ).

// ejemplo_de_uso

Úsala cuando un proceso de negocio toca varios microservicios y necesitas garantizar consistencia sin bloqueos distribuidos. Ej.: implementar el flujo de compra (reserva de stock, cobro, envío) como una saga con pasos de compensación automáticos si el cobro falla.

// plataformas

PythonKafkaRabbitMQSQSPrometheus
Categoría
Web
Tipo
Referencia
Nivel
Avanzado
Licencia
MIT
Seguridad
seguro · riesgo bajo
Versión
1.0.0

// opiniones_de_la_comunidad

Opiniones

Cargando opiniones…