La lección anterior terminó con un span en rojo: inventario-1 no responde y pedidos ha abierto el circuito hacia él. La observabilidad ha hecho su trabajo, que era señalar; ahora la plataforma tiene que sobrevivir. Eso exige responder a preguntas que en el monolito ni se planteaban: cómo distinguir un nodo caído de uno lento, quién decide que la réplica pase a ser primaria y cómo se impide que el antiguo primario vuelva creyendo que aún manda, cómo se reanuda un proceso de dos horas que se cayó en el minuto 90, y cómo se recuperan los datos de km0_inventario cuando lo que ha fallado no es un proceso sino un DELETE sin WHERE un viernes por la tarde. Esta lección recorre el ciclo completo: detectar, tolerar (redundancia y failover con cercado), recuperar (checkpoints, backups, PITR, RPO/RTO, DR) y gestionar el incidente mientras todo eso ocurre. Los patrones de llamada entre servicios (timeouts, reintentos, circuit breaker) son la lección 07-04; la orquestación con Kubernetes, la 07-05.

Contenido

  1. Fallos en la práctica: del modelo teórico a lo que ocurre
  2. Detección: heartbeats, health checks y falsos positivos
  3. Tolerancia por redundancia
  4. Failover, split-brain y cercado
  5. Failover real: Patroni, Kafka y Cassandra
  6. Recuperación de estado: checkpoints de procesos largos
  7. Recuperación de datos: backups, PITR y pruebas de restauración
  8. RPO, RTO y plan de recuperación ante desastres
  9. Gestión de incidentes: runbooks, on-call y postmortems
  10. Errores comunes y consejos
  11. Ejercicios y soluciones
  12. Conclusión

  1. Fallos en la práctica: del modelo teórico a lo que ocurre

En 01-02 se presentaron los modelos de fallo: crash-stop, crash-recovery, omisión, temporización y bizantino. Sirven para razonar sobre algoritmos; en operación, lo que se ve es su encarnación concreta, y casi nunca es el caso limpio:

Lo que ocurre Modelo de 01-02 Ejemplo en Kilómetro Cero Por qué es traicionero
Caída de nodo Crash-stop / crash-recovery Se muere el contenedor inventario-1 por OOM Es el caso "fácil": se detecta bien y hay réplica
Partición de red Omisión El switch entre las zonas bcn y vlc pierde paquetes durante 40 s Cada lado cree que el otro ha muerto; ambos siguen vivos (03-02)
Disco lleno Crash con degradación previa El broker kafka-2 llena su disco con reparto.posiciones sin retención Antes de morir, escribe lento, timeouts en productores, lag en consumidores
Fallo gris (lento pero vivo) Temporización pedidos-db-primario responde, pero cada consulta tarda 4 s por un disco degradado Los health checks dicen "OK"; los clientes agotan sus pools esperando (07-04)
Fallos correlacionados Varios crash a la vez Un cambio en la CA de Vault deja sin certificado válido a las tres réplicas de pedidos en el mismo minuto La redundancia no protege: comparten la causa
Error humano Cualquiera DELETE FROM stock WHERE productor_id = ... sin transacción ni WHERE correcto Se replica perfectamente a inv-bcn e inv-vlc: la réplica no es un backup
Despliegue defectuoso Crash o gris, correlacionado Versión de pedidos con un bug que devuelve 500 en el 8 % de los casos (el sábado de 07-01) Rolling update lo lleva a todas las réplicas; hace falta rollback (07-05)

Dos lecciones se repiten: los fallos lentos son peores que los fallos limpios (un nodo muerto se sustituye; uno lento contagia), y la redundancia solo protege contra fallos independientes. Todo lo que sigue intenta convertir fallos grises en fallos limpios (detectarlos y apartar el nodo) y romper las correlaciones (zonas, versiones escalonadas, backups fuera del sistema).

  1. Detección: heartbeats, health checks y falsos positivos

Un sistema no puede tolerar un fallo que no detecta. Hay dos direcciones de detección:

  • Heartbeats: el nodo vigilado envía periódicamente "estoy vivo" (a un coordinador, a sus pares, a etcd renovando un lease). Si deja de llegar, se le da por muerto. Es lo que usan Raft (03-03), Cassandra (gossip), Kafka (los brokers con el controlador) y Patroni (con etcd).
  • Health checks: el vigilante pregunta "¿estás bien?" al vigilado. Es lo que hacen Kubernetes, los balanceadores y Kong con los upstreams.

Tres tipos de health check

La distinción, popularizada por Kubernetes (07-05), vale para cualquier plataforma:

Sonda Pregunta Si falla Qué debe comprobar
Liveness (/salud/vivo) "¿El proceso está atascado sin remedio?" Reiniciar el proceso Solo que el proceso responde: un bucle interno bloqueado, memoria agotada. Nunca dependencias externas
Readiness (/salud/listo) "¿Puedes atender tráfico ahora?" Sacarlo del balanceo sin reiniciarlo Dependencias imprescindibles: conexión a km0_pedidos, productor Kafka conectado, configuración cargada
Startup "¿Has terminado de arrancar?" Esperar (no aplicar liveness todavía) Migraciones, caches iniciales, carga de certificados de Vault

El error clásico es comprobar PostgreSQL en la liveness: si la base de datos cae, todas las réplicas de pedidos fallan la sonda, se reinician en bucle, y cuando PostgreSQL vuelve, no hay nadie para atenderlo. Con la readiness, en cambio, las réplicas se retiran del balanceo, siguen vivas, y vuelven solas en cuanto la dependencia responde.

# km0/servicios/pedidos/salud.py
"""Health checks de pedidos: liveness sin dependencias, readiness con ellas."""
import asyncio
import time

from fastapi import APIRouter, Response

router = APIRouter(prefix="/salud")
ARRANQUE = time.monotonic()


@router.get("/vivo")
async def vivo():
    # Liveness: si podemos ejecutar esta función, el proceso no está colgado.
    # Devolver 200 siempre; el orquestador lo reiniciará si ni siquiera responde.
    return {"estado": "vivo", "segundos_arriba": int(time.monotonic() - ARRANQUE)}


async def _comprobar_postgres(pool) -> tuple[bool, str]:
    try:
        async with pool.acquire(timeout=1.0) as conn:        # esperar como mucho 1 s por una conexión
            await asyncio.wait_for(conn.execute("SELECT 1"), timeout=1.0)
        return True, "ok"
    except Exception as e:                                    # timeout, conexión rechazada...
        return False, f"postgres: {type(e).__name__}"


async def _comprobar_kafka(productor) -> tuple[bool, str]:
    try:
        # list_topics con timeout: si el clúster no responde, lanza excepción
        await asyncio.get_running_loop().run_in_executor(
            None, lambda: productor.list_topics(topic="pedidos.eventos", timeout=1.0))
        return True, "ok"
    except Exception as e:
        return False, f"kafka: {type(e).__name__}"


@router.get("/listo")
async def listo(response: Response):
    from servicios.pedidos.app import pool_pg, productor_kafka   # recursos del servicio
    resultados = await asyncio.gather(_comprobar_postgres(pool_pg), _comprobar_kafka(productor_kafka))
    detalle = {"postgres": resultados[0][1], "kafka": resultados[1][1]}
    if all(ok for ok, _ in resultados):
        return {"estado": "listo", **detalle}
    response.status_code = 503        # el balanceador deja de enviarnos tráfico
    return {"estado": "no_listo", **detalle}

Los timeouts de 1 s en cada comprobación son parte del diseño: un health check sin timeout que se queda esperando a un PostgreSQL gris hace que la propia sonda sea lenta, y el orquestador la interpreta como fallo. Con timeouts, un PostgreSQL gris se convierte en "no listo" en un segundo, que es justo la conversión de fallo gris en fallo limpio que se buscaba.

Timeouts de detección y falsos positivos

Todo detector basado en tiempo tiene una tensión: un umbral corto detecta rápido pero declara muertos a nodos vivos que solo estaban lentos (falso positivo); uno largo tarda en reaccionar. En una red con partición, además, es imposible distinguir "muerto" de "inalcanzable" (03-02), así que la pregunta correcta no es "¿está muerto?" sino "¿cuánta confianza tengo en que lo está?". El detector phi-accrual (Hayashibara et al., usado por Cassandra y Akka) responde con un número: en lugar de un umbral fijo, aprende la distribución de los intervalos entre heartbeats y calcula la sospecha φ como "cuán improbable es este retraso dado lo que he visto". Una versión simplificada:

# km0/simulaciones/heartbeat_detector.py
"""Detector de fallos phi-accrual simplificado.

Cada nodo envía heartbeats; el detector aprende la media y la desviación de los
intervalos y expresa la sospecha como phi = -log10(P(retraso >= observado)).
phi = 1 -> 10 % de probabilidad de que sea un retraso normal; phi = 3 -> 0,1 %.
"""
import math
import statistics
import time
from collections import deque


class DetectorPhi:
    def __init__(self, umbral_phi: float = 8.0, ventana: int = 100, intervalo_inicial: float = 1.0):
        self.umbral = umbral_phi
        self.intervalos: dict[str, deque] = {}
        self.ultimo: dict[str, float] = {}
        self.ventana = ventana
        self.intervalo_inicial = intervalo_inicial

    def heartbeat(self, nodo: str, ahora: float | None = None) -> None:
        ahora = ahora or time.monotonic()
        if nodo in self.ultimo:
            self.intervalos.setdefault(nodo, deque(maxlen=self.ventana)).append(ahora - self.ultimo[nodo])
        self.ultimo[nodo] = ahora

    def phi(self, nodo: str, ahora: float | None = None) -> float:
        ahora = ahora or time.monotonic()
        if nodo not in self.ultimo:
            return 0.0
        muestras = self.intervalos.get(nodo)
        if not muestras or len(muestras) < 2:
            media, desv = self.intervalo_inicial, self.intervalo_inicial / 4
        else:
            media = statistics.mean(muestras)
            desv = max(statistics.pstdev(muestras), media * 0.05)   # evitar desviación 0
        retraso = ahora - self.ultimo[nodo]
        # Aproximación con distribución normal: P(X >= retraso)
        z = (retraso - media) / desv
        p = 0.5 * math.erfc(z / math.sqrt(2))
        return -math.log10(max(p, 1e-12))     # cota para no dividir por 0

    def sospechoso(self, nodo: str, ahora: float | None = None) -> bool:
        return self.phi(nodo, ahora) > self.umbral


if __name__ == "__main__":
    # Simulación: inv-bcn envía heartbeats cada 1 s con jitter; a los 30 s se calla.
    import random
    det = DetectorPhi(umbral_phi=8.0)
    t = 0.0
    for i in range(30):
        det.heartbeat("inv-bcn", ahora=t)
        t += 1.0 + random.gauss(0, 0.05)
    ultimo_hb = t
    for retraso in (0.5, 1.0, 1.2, 1.5, 2.0, 3.0, 5.0):
        ahora = ultimo_hb + retraso
        print(f"retraso {retraso:4.1f} s  phi = {det.phi('inv-bcn', ahora):6.2f}  "
              f"{'SOSPECHOSO' if det.sospechoso('inv-bcn', ahora) else 'vivo'}")

Salida típica:

retraso  0.5 s  phi =   0.00  vivo
retraso  1.0 s  phi =   0.30  vivo
retraso  1.2 s  phi =   3.90  vivo
retraso  1.5 s  phi =  12.00  SOSPECHOSO
retraso  2.0 s  phi =  12.00  SOSPECHOSO
...

Con una red muy estable (desviación de 50 ms), 1,5 s de silencio ya es altamente sospechoso; en una red con jitter de 300 ms, el mismo detector esperaría más antes de sospechar, sin cambiar ningún parámetro. Ese es el valor de un detector adaptativo: el umbral se ajusta al comportamiento observado, y el consumidor del detector (el que decide el failover) elige cuánta confianza exige. El gossip es la forma de propagar este estado entre muchos nodos sin un coordinador: cada nodo cuenta a unos pocos, al azar, lo que sabe de los demás (Cassandra lo usa para pertenencia y estado). Basta con saber que existe.

  1. Tolerancia por redundancia

Detectado el fallo, tolerarlo exige tener con qué sustituir lo que falla:

Esquema Qué es Coste Ejemplo en Kilómetro Cero
N+1 Capacidad para la carga con N unidades, más una de reserva 1/N extra pedidos con 3 réplicas cuando 2 bastan para la Semana de la Vendimia
Activo-activo Todas las unidades atienden tráfico a la vez Todo trabaja; hace falta que sean intercambiables (sin estado o con estado replicado) Réplicas de pedidos, catalogo, inventario tras Kong; nodos de Cassandra
Activo-pasivo Una unidad atiende; otra espera lista para relevar La pasiva está ociosa; hay un instante de conmutación pedidos-db-primario / pedidos-db-replica; el líder del relay outbox con lease en etcd
Zonas de disponibilidad Unidades repartidas por dominios de fallo independientes (edificio, alimentación, red) Latencia entre zonas; coste de duplicar inv-bcn en una zona, inv-vlc en otra; brokers de Kafka repartidos con rack awareness

Lo que hace útil la redundancia es la independencia: tres réplicas de pedidos en la misma máquina no toleran que caiga la máquina; tres nodos de Cassandra en el mismo rack no toleran que caiga el switch. Y, como mostró la tabla de fallos, tres réplicas con el mismo certificado caducado no toleran nada. La redundancia se diseña por dominio de fallo: máquina, rack, zona, versión de software, credencial.

  1. Failover, split-brain y cercado

El failover es la conmutación de un componente fallido a su sustituto. Para servicios sin estado es trivial: el balanceador deja de enviar a la réplica no lista (readiness) y punto. Para componentes con estado y un solo escritor (el primario de PostgreSQL, el líder de una partición de Kafka, el líder del relay outbox) es el problema más delicado de toda la lección, y 03-04 lo dejó anunciado: split-brain.

El escenario: inv-bcn es primario y inv-vlc su réplica en streaming. Una partición de red de 40 s separa a inv-bcn de todo lo demás. El detector, desde el otro lado, lo da por muerto y promociona a inv-vlc. Pero inv-bcn no está muerto: sigue aceptando escrituras de los clientes que aún lo alcanzan. Cuando la red vuelve, hay dos primarios con historias divergentes: dos reservas del último queso-curado en dos bases de datos que ya no se pueden reconciliar automáticamente.

Las defensas, en orden de fortaleza:

  1. Consenso para la decisión: no promociona "quien lo ve muerto", sino una mayoría a través de un almacén de consenso (etcd, 03-03). Como mucho un nodo puede tener el lease de "líder" a la vez.
  2. Cercado (fencing): garantizar que el antiguo primario no puede escribir antes de que el nuevo empiece. Hay varias formas, y se combinan:
    • Token de cercado: cada líder recibe un número monótono creciente (la revision del lease de etcd, como en 03-03). Cualquier recurso compartido (el almacenamiento, el consumidor del outbox) rechaza escrituras con un token menor del último que vio. El viejo líder, con el token 41, no puede escribir donde ya se ha visto el 42.
    • Autocercado: el líder se degrada solo si no consigue renovar su lease (por eso el relay outbox comprueba su lease antes de cada lote y para si no lo tiene). Depende de que el reloj del líder no esté muy desviado (01-05).
    • STONITH ("shoot the other node in the head"): apagar físicamente el viejo nodo (por IPMI, por la API de la nube) antes de promocionar. Brutal, pero definitivo.
  3. Timeouts generosos y for en la decisión: no promocionar por un parpadeo de 3 s.
sequenceDiagram
    participant C as Clientes (inventario)
    participant A as inv-bcn (primario, lease rev 41)
    participant E as etcd (3 nodos)
    participant B as inv-vlc (réplica)
    Note over A,E: Partición: inv-bcn no alcanza etcd
    A--xE: renovar lease (falla)
    Note over A: TTL agotado sin renovar:<br/>AUTOCERCADO -> modo solo lectura
    E->>B: lease de líder expirado
    B->>E: adquirir lease (rev 42)
    E-->>B: concedido
    Note over B: promote: pg_promote()
    B->>C: "soy primario, token 42"
    C->>B: escrituras con token 42
    Note over A,E: La red vuelve
    A->>C: escritura con token 41
    C-->>A: rechazada (41 < 42)
    A->>E: ¿quién es líder?
    E-->>A: inv-vlc (42)
    Note over A: se reincorpora como réplica<br/>(pg_rewind si divergió)

El failback (volver al nodo original cuando se recupera) rara vez merece la pena de forma automática: cada conmutación es un riesgo. Lo habitual es que el recuperado se incorpore como réplica y se quede así hasta un switchover planificado.

  1. Failover real: Patroni, Kafka y Cassandra

5.1 Patroni para PostgreSQL

Patroni es un agente que corre junto a cada PostgreSQL y hace exactamente lo descrito: usa etcd (o Consul, o la API de Kubernetes) como almacén de consenso, mantiene un lease de líder con TTL, promociona la réplica más avanzada cuando el lease expira, y degrada al viejo primario que no logra renovar. Para km0_inventario, el docker-compose.yml:

# km0/docker-compose.yml (fragmento): Patroni con 3 nodos para km0_inventario
  etcd-1:
    image: quay.io/coreos/etcd:v3.5.15
    command: ["etcd", "--name=etcd-1", "--initial-cluster=etcd-1=http://etcd-1:2380,etcd-2=http://etcd-2:2380,etcd-3=http://etcd-3:2380",
              "--listen-peer-urls=http://0.0.0.0:2380", "--listen-client-urls=http://0.0.0.0:2379",
              "--advertise-client-urls=http://etcd-1:2379", "--initial-advertise-peer-urls=http://etcd-1:2380"]
  # etcd-2 y etcd-3 iguales cambiando el nombre

  inv-bcn: &patroni
    image: ghcr.io/zalando/spilo-16:3.3-p1      # PostgreSQL 16 + Patroni
    environment: &patroni_env
      SCOPE: km0-inventario                      # nombre del clúster en etcd
      ETCD3_HOSTS: "etcd-1:2379,etcd-2:2379,etcd-3:2379"
      PGPASSWORD_SUPERUSER_FILE: /run/secrets/pg_super
      PATRONI_TTL: "30"                          # lease del líder: 30 s
      PATRONI_LOOP_WAIT: "10"                    # cada 10 s renueva
      PATRONI_RETRY_TIMEOUT: "10"
      PATRONI_MAXIMUM_LAG_ON_FAILOVER: "1048576" # no promocionar una réplica con > 1 MB de retraso
      PATRONI_SYNCHRONOUS_MODE: "true"           # al menos una réplica síncrona: RPO 0 en failover
    hostname: inv-bcn
    volumes: ["inv-bcn-datos:/home/postgres/pgdata"]
  inv-vlc:
    <<: *patroni
    hostname: inv-vlc
    volumes: ["inv-vlc-datos:/home/postgres/pgdata"]
  inv-gir:
    <<: *patroni
    hostname: inv-gir
    volumes: ["inv-gir-datos:/home/postgres/pgdata"]

  # Los clientes no conocen quién es primario: HAProxy pregunta a Patroni (/primary devuelve 200 solo en el líder)
  inventario-db:
    image: haproxy:2.9
    volumes: ["./observabilidad/haproxy-patroni.cfg:/usr/local/etc/haproxy/haproxy.cfg:ro"]
    ports: ["5432:5432", "5433:5433"]     # 5432 -> primario (escritura); 5433 -> réplicas (lectura)

Los parámetros que importan: TTL y LOOP_WAIT definen la ventana de detección (hasta 30 s sin renovar = lease perdido); MAXIMUM_LAG_ON_FAILOVER evita promocionar una réplica muy atrasada (perdería más datos de los aceptables); SYNCHRONOUS_MODE hace que cada commit espere a una réplica, con lo que un failover no pierde transacciones confirmadas (el precio es latencia de escritura, la misma tensión de 03-01). HAProxy consulta el endpoint REST de Patroni para saber a quién enviar, de modo que inventario se conecta siempre a inventario-db:5432 y nunca necesita saber quién es primario.

Operación del día a día:

$ patronictl -c /etc/patroni.yml list
+ Cluster: km0-inventario -----+---------+-----------+----+-----------+
| Member  | Host     | Role    | State     | TL | Lag in MB |
+---------+----------+---------+-----------+----+-----------+
| inv-bcn | inv-bcn  | Leader  | running   | 12 |           |
| inv-vlc | inv-vlc  | Sync Standby | streaming | 12 |     0 |
| inv-gir | inv-gir  | Replica | streaming | 12 |         0 |
+---------+----------+---------+-----------+----+-----------+

# Conmutación planificada (mantenimiento de inv-bcn): sin pérdida, con confirmación
$ patronictl -c /etc/patroni.yml switchover --leader inv-bcn --candidate inv-vlc --scheduled now
Are you sure you want to switchover cluster km0-inventario, demoting current leader inv-bcn? [y/N]: y
Successfully switched over to "inv-vlc"

# Tras un failover no planificado, el viejo líder se reincorpora; si divergió, Patroni ejecuta pg_rewind
$ patronictl -c /etc/patroni.yml list
| inv-bcn | inv-bcn  | Replica | streaming | 13 |         0 |
| inv-vlc | inv-vlc  | Leader  | running   | 13 |           |

La columna TL (timeline) sube con cada promoción: es el "token de cercado" de PostgreSQL. Un primario en la timeline 12 no puede enviar WAL a réplicas que ya están en la 13.

5.2 Kafka: réplicas de partición, ISR y elección de líder

Kafka no tiene "un primario": cada partición tiene un líder y N-1 réplicas seguidoras, repartidas por brokers. El conjunto de réplicas al día se llama ISR (in-sync replicas). Cuando cae el broker líder de una partición, el controlador del clúster elige nuevo líder entre las ISR; los productores y consumidores descubren el cambio en la siguiente petición de metadatos. Los parámetros que gobiernan qué se pierde:

Parámetro Valor en Kilómetro Cero Efecto
replication.factor 3 (pedidos.eventos, auditoria.eventos), 2 (reparto.posiciones) Cuántas copias de cada partición
min.insync.replicas 2 Un acks=all solo se confirma si al menos 2 réplicas lo tienen; si solo queda 1 ISR, el productor recibe NotEnoughReplicas (se prefiere no aceptar a perder)
acks (productor) all en pedidos, 1 en reparto Cuándo considera el productor que el mensaje está guardado
unclean.leader.election.enable false Nunca elegir líder a una réplica fuera de ISR: mejor partición no disponible que perder mensajes confirmados

La combinación replication.factor=3 + min.insync.replicas=2 + acks=all tolera la caída de un broker sin perder ni un evento de pedido confirmado; con dos brokers caídos, pedidos.eventos deja de aceptar escrituras (y el outbox de 02-05 las retiene hasta que vuelva). reparto.posiciones acepta perder posiciones a cambio de latencia.

5.3 Cassandra: sin líder que conmutar

km0_pedidos en Cassandra (04-04) no necesita failover porque no hay líder: cada fila vive en 3 nodos y cada lectura y escritura con LOCAL_QUORUM necesita 2. Si cae un nodo, las operaciones siguen con los otros dos; el detector phi-accrual y el gossip marcan el nodo caído; cuando vuelve, los hinted handoffs (escrituras guardadas para él por sus vecinos) y la reparación (nodetool repair) lo ponen al día. El precio ya se pagó en el diseño: consistencia eventual entre réplicas y un modelo de datos por consultas.

  1. Recuperación de estado: checkpoints de procesos largos

No todo es bases de datos. Cada noche, inventario ejecuta la reconciliación de stock: recorre los tres productores, compara el stock de km0_inventario con las reservas confirmadas en km0_pedidos y las entregas de reparto, y corrige las desviaciones. Tarda unas dos horas. Si el proceso muere a los 90 minutos (OOM, despliegue, nodo caído), empezar de cero significa otras dos horas, y puede que no acabe antes de que abran los mercados. La solución es la misma idea que los checkpoints de Flink (05-04): persistir el progreso periódicamente en un lugar que sobreviva al proceso, y reanudar desde el último checkpoint de forma idempotente.

# km0/servicios/inventario/reconciliacion_nocturna.py
"""Reconciliación nocturna de stock con checkpoint persistido y reanudación."""
import json
import time
from dataclasses import dataclass, asdict

import psycopg
from servicios.comun.logs import log

CHECKPOINT_CADA = 500          # productos procesados entre checkpoints


@dataclass
class Checkpoint:
    ejecucion_id: str          # p. ej. "2026-09-13"
    productor_actual: str      # slug del productor que se está procesando
    ultimo_producto: str       # último producto confirmado dentro de ese productor ("" = ninguno)
    procesados: int
    corregidos: int


def cargar_checkpoint(conn, ejecucion_id: str) -> Checkpoint | None:
    fila = conn.execute(
        "SELECT estado FROM reconciliacion_checkpoints WHERE ejecucion_id = %s", (ejecucion_id,)
    ).fetchone()
    return Checkpoint(**json.loads(fila[0])) if fila else None


def guardar_checkpoint(conn, cp: Checkpoint) -> None:
    # UPSERT: la fila del día se sobreescribe; se confirma en la misma transacción
    # que las correcciones del lote, así nunca hay correcciones sin checkpoint ni al revés.
    conn.execute(
        """INSERT INTO reconciliacion_checkpoints (ejecucion_id, estado, actualizado)
           VALUES (%s, %s, now())
           ON CONFLICT (ejecucion_id) DO UPDATE SET estado = EXCLUDED.estado, actualizado = now()""",
        (cp.ejecucion_id, json.dumps(asdict(cp))),
    )


def productos_desde(conn, productor: str, ultimo_producto: str):
    # Orden determinista por slug: reanudar es "seguir a partir del último confirmado"
    yield from conn.execute(
        "SELECT slug FROM productos WHERE productor = %s AND slug > %s ORDER BY slug",
        (productor, ultimo_producto),
    )


def reconciliar_producto(conn, productor: str, producto: str) -> bool:
    """Compara stock con reservas y entregas; corrige si hay desviación. Idempotente:
    ejecutarlo dos veces sobre el mismo producto deja el mismo resultado."""
    esperado = calcular_stock_esperado(conn, productor, producto)     # consulta a km0_pedidos y reparto
    actual = conn.execute("SELECT unidades FROM stock WHERE producto = %s FOR UPDATE", (producto,)).fetchone()[0]
    if actual != esperado:
        conn.execute("UPDATE stock SET unidades = %s WHERE producto = %s", (esperado, producto))
        log.warning("stock_corregido", productor=productor, producto=producto, antes=actual, despues=esperado)
        return True
    return False


def ejecutar(ejecucion_id: str, productores: list[str]) -> None:
    t0 = time.monotonic()
    with psycopg.connect(DSN_INVENTARIO) as conn:
        cp = cargar_checkpoint(conn, ejecucion_id) or Checkpoint(ejecucion_id, productores[0], "", 0, 0)
        if cp.procesados:
            log.info("reconciliacion_reanudada", desde_productor=cp.productor_actual,
                     desde_producto=cp.ultimo_producto, procesados=cp.procesados)
        # Saltar los productores ya completados
        for productor in productores[productores.index(cp.productor_actual):]:
            desde = cp.ultimo_producto if productor == cp.productor_actual else ""
            pendiente_en_lote = 0
            for (producto,) in productos_desde(conn, productor, desde):
                if reconciliar_producto(conn, productor, producto):
                    cp.corregidos += 1
                cp.procesados += 1
                cp.productor_actual, cp.ultimo_producto = productor, producto
                pendiente_en_lote += 1
                if pendiente_en_lote >= CHECKPOINT_CADA:
                    guardar_checkpoint(conn, cp)
                    conn.commit()            # correcciones + checkpoint, atómicamente
                    pendiente_en_lote = 0
            guardar_checkpoint(conn, cp)
            conn.commit()
        log.info("reconciliacion_terminada", procesados=cp.procesados, corregidos=cp.corregidos,
                 duracion_s=int(time.monotonic() - t0))


if __name__ == "__main__":
    ejecutar(time.strftime("%Y-%m-%d"), ["huerta-la-vega", "queseria-montblanc", "bodega-roble-alto"])

Las tres propiedades que hacen que esto funcione, y que valen para cualquier proceso largo:

  1. El checkpoint se confirma en la misma transacción que el trabajo que representa. Si el proceso muere entre el UPDATE y el checkpoint, la transacción no se confirma y ambos se pierden juntos; nunca queda un checkpoint que diga "hasta queso-curado" con queso-curado sin corregir, ni al revés.
  2. El orden es determinista (ORDER BY slug), de modo que "reanudar desde X" tiene significado.
  3. Cada unidad de trabajo es idempotente: si el checkpoint se guardó antes de un lote de 499 productos y el proceso murió, esos 499 se reprocesan al reanudar, y reprocesarlos no hace daño.

Airflow (05-05) es quien lanza este proceso y quien lo relanza si falla (retries): la reanudación desde checkpoint hace que el reintento cueste minutos y no horas.

  1. Recuperación de datos: backups, PITR y pruebas de restauración

La replicación protege contra la caída de un nodo; no protege contra un dato malo. El DELETE erróneo de la sección 1 llega a inv-vlc en milisegundos. Para eso están los backups, que son copias desacopladas en el tiempo del sistema.

Tipo Cómo Ventajas Inconvenientes En Kilómetro Cero
Lógico pg_dump: volcado SQL o formato propio Portable entre versiones; restaurar una tabla suelta Lento en bases grandes; restaurar es reejecutar; no permite PITR km0_analitica semanal (se puede regenerar del lago)
Físico pg_basebackup: copia de los ficheros de datos Rápido; base para PITR Misma versión mayor; todo o nada km0_inventario diario
PITR (point-in-time recovery) Backup físico + archivo continuo del WAL Restaurar a cualquier instante (las 16:59, un minuto antes del DELETE) Hay que guardar todo el WAL desde el último base backup km0_inventario con WAL en km0-backups
Snapshot Copia de los SSTables de Cassandra (nodetool snapshot), enlaces duros instantáneos Casi gratis de crear Hay que copiarlos fuera del nodo; por nodo km0_pedidos diario, a km0-backups
Versionado de objetos MinIO conserva versiones anteriores de cada objeto Un borrado o sobrescritura se deshace Coste de almacenamiento km0-facturas, km0-auditoria (con object lock, 06-05)

PITR en PostgreSQL paso a paso

PostgreSQL escribe cada cambio primero en el WAL (write-ahead log), en segmentos de 16 MB. Si se guardan todos los segmentos desde un backup base, se puede reproducir la historia hasta el instante deseado.

flowchart LR
    subgraph normal["Operación normal"]
        PG[(inv-bcn<br/>primario)] -- "archive_command<br/>cada segmento WAL" --> M[(MinIO<br/>km0-backups/inventario/wal/)]
        PG -- "pg_basebackup<br/>diario 02:00" --> MB[(km0-backups/inventario/base/2026-09-12/)]
    end
    subgraph recuperacion["Recuperación a las 16:59"]
        MB --> R[Nodo de restauración]
        M -- "restore_command<br/>replay hasta recovery_target_time" --> R
        R --> V{¿datos correctos?}
        V -- sí --> P[pg_promote → nuevo primario<br/>Patroni reinicializa réplicas]
    end

Configuración del archivado en el primario (Patroni la aplica en postgresql.parameters):

# km0/sql/backup/archivar_wal.sh — invocado por PostgreSQL con %p (ruta) y %f (nombre del segmento)
#!/usr/bin/env bash
set -euo pipefail
# mc: cliente de MinIO; el alias 'km0' se configuró con credenciales de Vault (06-04)
mc cp --quiet "$1" "km0/km0-backups/inventario/wal/$2"
# Parámetros de PostgreSQL gestionados por Patroni (fragmento de patroni.yml)
postgresql:
  parameters:
    wal_level: replica
    archive_mode: "on"
    archive_command: "/opt/km0/sql/backup/archivar_wal.sh %p %f"
    archive_timeout: 60        # forzar un segmento cada minuto aunque no esté lleno: RPO <= 1 min

Backup base diario:

# km0/sql/backup/base_backup.sh
#!/usr/bin/env bash
set -euo pipefail
FECHA=$(date -u +%F)
DESTINO=/backups/base/$FECHA
mkdir -p "$DESTINO"
# -Ft: tar; -z: comprimido; -X stream: incluye el WAL necesario para que el backup sea consistente por sí mismo
pg_basebackup -h inventario-db -p 5432 -U replicador -D "$DESTINO" -Ft -z -X stream --checkpoint=fast
mc cp --recursive "$DESTINO" "km0/km0-backups/inventario/base/$FECHA/"
# Conservar 14 días de backups base; el WAL anterior al más antiguo ya no sirve
mc rm --recursive --force --older-than 14d km0/km0-backups/inventario/base/

Y la restauración a las 16:59 del viernes, un minuto antes del DELETE de las 17:00:

# 1. Nodo limpio (inv-rest): descargar el último backup base ANTERIOR al instante objetivo
mc cp --recursive km0/km0-backups/inventario/base/2026-09-11/ /restauracion/base/
mkdir -p /restauracion/pgdata && cd /restauracion/pgdata
tar -xzf /restauracion/base/base.tar.gz
tar -xzf /restauracion/base/pg_wal.tar.gz -C pg_wal/

# 2. Decirle a PostgreSQL hasta dónde reproducir y de dónde sacar el WAL
cat >> postgresql.auto.conf <<'EOF'
restore_command = 'mc cp --quiet km0/km0-backups/inventario/wal/%f %p'
recovery_target_time = '2026-09-11 16:59:00+02'
recovery_target_action = 'promote'      # al llegar, salir del modo recuperación
EOF
touch recovery.signal

# 3. Arrancar: reproduce el WAL del día hasta las 16:59 y se promociona
pg_ctl -D /restauracion/pgdata start
# LOG:  starting point-in-time recovery to 2026-09-11 16:59:00+02
# LOG:  restored log file "000000010000004A000000F3" from archive
# ...
# LOG:  recovery stopping before commit of transaction 8812345, time 2026-09-11 16:59:12.8
# LOG:  database system is ready to accept connections

# 4. Verificar antes de tocar producción
psql -d km0_inventario -c "SELECT count(*), sum(unidades) FROM stock WHERE productor = 'queseria-montblanc';"

# 5. Decidir: (a) extraer las filas borradas y reinsertarlas en producción con un script (lo habitual,
#    si producción ha seguido recibiendo pedidos después de las 17:00), o
#    (b) convertir inv-rest en el nuevo primario del clúster Patroni y reinicializar las réplicas
#    (si el daño es tan grande que se prefiere perder lo posterior a las 16:59).

El paso 5 es el que se olvida en los ejercicios de manual: entre las 17:00 y el momento de la restauración, km0_inventario ha seguido recibiendo reservas. Volver entera a las 16:59 las perdería. Casi siempre se opta por restaurar aparte y reinyectar lo que faltaba.

Cassandra y MinIO

# Snapshot en cada nodo de km0_pedidos (enlaces duros a los SSTables actuales: instantáneo)
nodetool snapshot -t diario-2026-09-12 km0_pedidos
# Copiar fuera del nodo (por keyspace/tabla) y borrar el snapshot local
mc cp --recursive /var/lib/cassandra/data/km0_pedidos/*/snapshots/diario-2026-09-12/ \
      km0/km0-backups/cassandra/$(hostname)/2026-09-12/
nodetool clearsnapshot -t diario-2026-09-12
# Restaurar: copiar los SSTables al directorio de la tabla y ejecutar `nodetool refresh km0_pedidos pedidos`
# (o cargarlos con sstableloader en un clúster distinto)

# MinIO: activar versionado en los buckets cuyo borrado accidental importe
mc version enable km0/km0-facturas
mc undo km0/km0-facturas/2026/09/F-2026-004411.pdf   # deshacer el último borrado o sobrescritura

Pruebas de restauración

Un backup que nunca se ha restaurado no es un backup: es una esperanza. Los fallos habituales son de lo más prosaicos: el archive_command lleva tres semanas fallando en silencio (alerta: pg_stat_archiver.failed_count sube, y es la métrica pg_stat_archiver_failed_count de postgres_exporter); el backup base está corrupto; la restauración necesita un parámetro que nadie recuerda; tarda seis horas y el RTO era una. Por eso Kilómetro Cero tiene un DAG de Airflow, km0_prueba_restauracion, que cada domingo restaura el último backup de km0_inventario en un contenedor efímero, ejecuta unas consultas de verificación (número de productos por productor, suma de stock) contra lo que se registró en el momento del backup, mide el tiempo, y publica km0_backup_restauracion_ok{bd="inventario"} y km0_backup_restauracion_segundos a Prometheus. Si falla, es una alerta de ticket con la misma prioridad que un fallo de producción.

  1. RPO, RTO y plan de recuperación ante desastres

Dos números resumen lo que un sistema promete ante un desastre:

  • RPO (Recovery Point Objective): cuántos datos, medidos en tiempo, se aceptan perder. Un RPO de 1 minuto significa que la última copia utilizable tiene como mucho un minuto de antigüedad.
  • RTO (Recovery Time Objective): cuánto tiempo puede estar el servicio caído hasta recuperarse.

Ambos se fijan por negocio y se pagan en arquitectura: RPO 0 exige replicación síncrona; RTO de minutos exige failover automático y probado. Para Kilómetro Cero:

Componente RPO RTO Cómo se consigue Qué se pierde si se supera
km0_inventario (PostgreSQL) 0 en failover (réplica síncrona); 1 min en PITR (archive_timeout) 1 min (Patroni); 1 h (PITR) Patroni con 3 nodos; WAL a km0-backups; prueba semanal Reservas de stock: sobreventa a los productores
km0_pedidos (Cassandra, 3 nodos) 0 con un nodo caído; 24 h ante pérdida total (snapshot diario) 0 con un nodo caído; 4 h ante pérdida total RF 3 + LOCAL_QUORUM; snapshots a MinIO; los eventos de pedidos.eventos permiten reconstruir Pedidos: el daño más directo a clientes
pedidos-db-* (PostgreSQL, streaming) segundos (réplica asíncrona) 5 min (promoción manual documentada) Streaming + runbook Estado de sagas: compensaciones pendientes
Kafka (pedidos.eventos, auditoria.eventos) 0 (RF 3, min.insync.replicas 2, acks=all) 30 s (elección de líder) Configuración de tópicos Eventos: analítica y auditoría incompletas
Kafka (reparto.posiciones) minutos (RF 2, acks=1) 30 s Se acepta: las posiciones son efímeras Nada relevante
Redis (caché) ∞ (se regenera) 1 min Sin backup; el catálogo se recalienta (04-05) Un pico de carga en PostgreSQL
MinIO km0-facturas, km0-auditoria 0 (versionado, object lock, réplica de bucket a otra región) 1 h Replicación de bucket Obligaciones legales
km0_analitica 24 h 8 h Volcado semanal + reejecución del DAG desde el lago Informes: se recalculan

La última línea de defensa es el plan de recuperación ante desastres (DR): qué hacer si desaparece un centro de datos entero (incendio, corte eléctrico prolongado, error del proveedor). Las estrategias, de más barata a más cara:

Estrategia Qué hay en la región secundaria RPO / RTO típicos Coste
Backup y restaurar Solo los backups (km0-backups replicado) Horas / horas-días Mínimo
Pilot light Datos replicados en caliente (réplicas de PostgreSQL, mirror de Kafka); servicios apagados, listos para arrancar Minutos / decenas de minutos Bajo
Warm standby Todo desplegado a escala reducida y recibiendo réplica; se escala al conmutar Segundos-minutos / minutos Medio
Multi-región activa Ambas regiones atienden tráfico; datos replicados en ambos sentidos ~0 / ~0 Alto, y complejo (conflictos, 03-04)

Kilómetro Cero opta por pilot light: réplica de bucket de MinIO a la región secundaria, una réplica de PostgreSQL de cada clúster en la otra región (asíncrona, fuera del quórum de Patroni) y un mirror de los tópicos críticos de Kafka. Dónde vive esa región secundaria y cómo se despliega con servicios gestionados es materia de 08-03. Lo que sí es de esta lección: el plan se ensaya (un game day al semestre, 07-06), tiene un runbook, y el criterio para activarlo está escrito de antemano, porque en mitad de un desastre nadie razona bien.

  1. Gestión de incidentes: runbooks, on-call y postmortems

Todo lo anterior son mecanismos; los incidentes los gestionan personas bajo presión. Tres prácticas breves:

Runbooks. Cada alerta de página (07-01) enlaza a un documento con: qué significa, cómo confirmar, qué mirar (paneles y consultas concretas), acciones ordenadas de menos a más invasivas, y cuándo escalar. El runbook de PedidosBurnRateRapido empieza con "¿ha habido un despliegue en la última hora? Si sí, rollback primero, investigar después". Se escriben en frío y se corrigen después de cada uso.

On-call. Alguien (Jordi o Marta esta semana) recibe las páginas, con un secundario de respaldo, rotación semanal, y compensación. La carga se mide: más de dos páginas por noche es un problema del sistema, no de la persona. Un incidente tiene un coordinador (que comunica y decide) separado de quien investiga, y un canal con un registro de tiempos, que será la materia prima del postmortem.

Postmortems sin culpa. Después de cada incidente relevante: cronología, impacto (en términos de SLO y presupuesto de error), causas contribuyentes (en plural: casi nunca es una), qué funcionó, qué no, y acciones con dueño y fecha. "Sin culpa" no es cortesía: si señalar a quien ejecutó el DELETE fuera el resultado, la próxima vez nadie contará qué pasó, y el sistema que permitió ejecutar un DELETE sin transacción en producción seguirá igual. La acción correcta es "las sesiones de operador en km0_inventario arrancan con SET default_transaction_read_only = on" y "PITR probado semanalmente", no "más cuidado".

Errores Comunes y Consejos

  • Comprobar dependencias en la liveness. Reinicia en bucle todas las réplicas cuando cae PostgreSQL. Liveness: solo el proceso. Readiness: las dependencias, con timeouts.
  • Health checks sin timeout. Una dependencia gris hace lenta la sonda y el orquestador la interpreta como fallo, pero tarde y mal. Cada comprobación con su propio límite de 1 s.
  • Promocionar sin cercado. "Lo vemos caído, promocionamos" es la receta del split-brain. Consenso para decidir (etcd/Patroni), token o timeline para rechazar al viejo, autocercado por lease.
  • Creer que la réplica es un backup. Replica los errores igual de bien que los aciertos. Backup físico + WAL archivado, fuera del sistema, con retención.
  • Backups que nunca se restauran. Son una esperanza. Prueba de restauración automatizada, medida y alertada.
  • Restaurar entera a un instante pasado sin pensar en lo posterior. Se pierden las transacciones legítimas desde entonces. Restaurar aparte y reinyectar.
  • RPO/RTO sin escribir. Si no están escritos, el sistema tiene los que salgan. Tabla por componente, acordada con negocio, y arquitectura que la cumpla.
  • Redundancia sin independencia. Tres réplicas en la misma máquina, con el mismo certificado, desplegadas a la vez. Dominios de fallo distintos para cada eje.
  • Postmortem que termina en "más cuidado". No es una acción. Cada causa contribuyente tiene un cambio en el sistema con dueño y fecha.

Ejercicios

Ejercicio 1. Durante la Semana de la Vendimia, el disco de inv-bcn empieza a degradarse: sigue respondiendo, pero cada UPDATE tarda 3-4 s. Patroni renueva el lease sin problemas (etcd responde), /salud/listo de inventario sigue devolviendo 200 porque el SELECT 1 tarda 20 ms, y las alertas de latencia de pedidos saltan. (a) ¿Qué tipo de fallo es y por qué ninguno de los dos detectores lo ve? (b) Propón un cambio en la readiness de inventario y otro en la configuración de Patroni o en la operativa que conviertan este fallo gris en una conmutación limpia, y comenta el riesgo de falsos positivos de cada uno. (c) Tras el switchover a inv-vlc, ¿qué pasa con las transacciones que inv-bcn tenía a medias, y qué garantiza que ninguna reserva confirmada a pedidos se pierda?

Ejercicio 2. El viernes a las 17:00 un operador ejecuta en km0_inventario un UPDATE stock SET unidades = 0 WHERE productor = 'bodega-roble-alto' creyendo estar en el entorno de pruebas. Se detecta a las 17:25 por la alerta de negocio "reservas rechazadas por stock" (07-01). El último backup base es de las 02:00 y el WAL se archiva cada minuto. (a) Describe el procedimiento completo, con los comandos de la sección 7, para recuperar el stock de Bodega Roble Alto sin perder las reservas que otros productores recibieron entre 17:00 y 17:25. (b) ¿Cuánto WAL hay que reproducir aproximadamente, y de qué depende el RTO real? (c) Escribe dos acciones de postmortem que no sean "más cuidado".

Ejercicio 3. La reconciliación nocturna se cae a las 03:40 tras 95 minutos, con el checkpoint {"productor_actual": "queseria-montblanc", "ultimo_producto": "queso-fresco", "procesados": 1150}. Airflow la relanza a las 03:42. (a) ¿Desde dónde continúa exactamente y qué productos se reprocesan? (b) Un compañero propone guardar el checkpoint en Redis "porque es más rápido". ¿Qué se pierde? (c) Otro propone que reconciliar_producto envíe un evento a inventario.alertas cada vez que corrige stock; ¿qué problema aparece al reanudar, y cómo lo resolverías con lo visto en 02-05?

Soluciones

Ejercicio 1.

(a) Es un fallo gris (modelo de temporización): el nodo está vivo y responde a las comprobaciones ligeras (lease de Patroni, SELECT 1), pero es inservible para el trabajo real. Los detectores miden "responde", no "responde a lo que importa". (b) Readiness de inventario: sustituir SELECT 1 por una comprobación representativa con el mismo timeout de 1 s (por ejemplo, UPDATE sonda SET ts = now() WHERE id = 1, una fila propia que ejercita el WAL y el disco); si tarda más de 1 s, 503 y las réplicas de inventario salen del balanceo, lo que al menos deja de contagiar a pedidos (el circuito de 07-04 hará el resto); riesgo: un pico puntual de E/S saca a todas las réplicas a la vez, así que hace falta for/umbral de fallos consecutivos (3 seguidos) antes de retirar. En Patroni: no hay detección de latencia de disco integrada, así que la vía es una alerta sobre pg_stat_*/latencia de E/S de node_exporter (USE del disco, 07-01) con un runbook cuyo primer paso es patronictl switchover --leader inv-bcn --candidate inv-vlc; o, si se quiere automático, un watchdog propio que ejecute ese switchover cuando la latencia de escritura supere N segundos durante M minutos, con el riesgo de conmutar por una carga alta legítima; conmutar es más barato que sobrevender, así que umbrales conservadores pero automatizados. (c) El switchover de Patroni hace un checkpoint, espera a que la réplica síncrona esté al día y degrada a inv-bcn; las transacciones a medias (no confirmadas) se abortan y el cliente recibe un error, que pedidos traduce a UNAVAILABLE y reintenta contra el nuevo primario (07-04, siempre que la reserva sea idempotente, 02-05); ninguna transacción confirmada se pierde porque SYNCHRONOUS_MODE garantizaba que el commit ya estaba en inv-vlc antes de responder a pedidos.

Ejercicio 2.

(a) 1) Confirmar el alcance: en producción, SELECT count(*) FROM stock WHERE productor='bodega-roble-alto' AND unidades=0 y el log de auditoría del operador (06-05) para fijar la hora exacta (17:00:12). 2) En un nodo aparte (inv-rest), restaurar el base backup de las 02:00 y reproducir WAL con recovery_target_time = '2026-09-11 17:00:00+02' (justo antes del UPDATE), recovery_target_action = 'promote'. 3) Verificar en inv-rest: SELECT producto, unidades FROM stock WHERE productor='bodega-roble-alto'. 4) Calcular el stock correcto actual: el de las 17:00 en inv-rest menos las reservas confirmadas entre 17:00 y 17:25 en producción (que sí son válidas: algunas fallaron por stock 0, pero las que entraron antes de las 17:00:12 o de otros productores no se tocan); en la práctica, para Bodega Roble Alto: unidades_17:00 - reservas_confirmadas_desde_17:00(producto), obtenidas de km0_pedidos/pedidos.eventos. 5) Aplicar en producción con un UPDATE ... FROM (VALUES ...) dentro de una transacción, comparar antes de COMMIT, y registrar la corrección. 6) Relanzar la reconciliación de la sección 6 para ese productor como verificación. Nunca convertir inv-rest en primario: se perderían 25 minutos de reservas de Huerta La Vega y Quesería Montblanc. (b) 15 horas de WAL (02:00 → 17:00); en un día normal de km0_inventario puede ser de decenas de GB durante la Vendimia; el RTO real depende de la velocidad de descarga desde MinIO y de replay (mono-hilo en PostgreSQL), de si el runbook está probado (el DAG dominical da la cifra: si dice 40 min, el RTO de 1 h es creíble) y del tiempo de calcular y aplicar la corrección. (c) "Las conexiones de operadores a producción se abren con default_transaction_read_only = on y un role distinto que exige SET ROLE escritor explícito con auditoría"; "el prompt de psql y el nombre del host de producción llevan PROD y color distinto, y los entornos de prueba no comparten credenciales de Vault con producción"; "alerta de negocio reservas rechazadas por stock con for: 2m en lugar de 15, que habría detectado a las 17:03".

Ejercicio 3.

(a) Empieza en queseria-montblanc, con productos_desde(conn, "queseria-montblanc", "queso-fresco"), es decir, el primer producto cuyo slug sea mayor que queso-fresco en orden alfabético; Huerta La Vega no se toca (ya completada). Como el checkpoint solo se confirma cada 500 productos o al acabar un productor, los productos procesados entre el último checkpoint y la caída (hasta 499) se reprocesan; es correcto porque reconciliar_producto es idempotente. (b) Se pierde la atomicidad entre checkpoint y correcciones: Redis y PostgreSQL no comparten transacción, así que puede quedar el checkpoint avanzado con las correcciones sin confirmar (se saltarían productos sin reconciliar) o al revés (solo se reprocesa, lo cual es inocuo); además, Redis sin persistencia puede perder el checkpoint en un reinicio. La rapidez es irrelevante: se escribe una fila cada 500 productos. (c) Al reanudar, los productos reprocesados que ya se corrigieron en la transacción confirmada no volverán a corregirse (actual ya es igual a esperado), pero los del lote no confirmado sí se corrigen otra vez... y en realidad el problema es el contrario: si el evento se envía a Kafka directamente desde reconciliar_producto, se publica antes del commit, y si el proceso muere, hay un evento en inventario.alertas de una corrección que nunca ocurrió; o, si se reintenta, un evento duplicado. La solución es el patrón Outbox de 02-05: la corrección y la fila de outbox se escriben en la misma transacción que el checkpoint, y el relay las publica después; los consumidores desduplican por clave (producto + ejecucion_id).

Conclusión

Detectar, tolerar, recuperar. La plataforma detecta con heartbeats y health checks separados por intención (liveness sin dependencias, readiness con ellas y con timeouts, startup para arrancar), y con detectores adaptativos que expresan sospecha en lugar de certeza, porque en una red no se puede distinguir muerto de inalcanzable. Tolera con redundancia que solo vale si es independiente por dominio de fallo, y con failover que solo es seguro si una mayoría decide y el antiguo primario queda cercado: Patroni sobre etcd para km0_inventario (lease, timeline, réplica síncrona, switchover), ISR y min.insync.replicas para Kafka, y quórum sin líder para Cassandra. Recupera estado con checkpoints confirmados en la misma transacción que el trabajo, y datos con backups físicos y WAL archivado en km0-backups que permiten volver al minuto anterior al error, snapshots de Cassandra y versionado de MinIO, todo ello inútil si no se restaura de verdad cada semana. RPO y RTO ponen números por componente a lo que se promete, el plan de DR en pilot light cubre la pérdida de una región, y runbooks, on-call y postmortems sin culpa convierten cada incidente en un cambio del sistema y no en un reproche.

Queda un hueco entre la detección y el failover: esos 30 segundos en que inv-bcn no responde y Patroni aún no ha promocionado a inv-vlc, o los cuatro segundos de cada UPDATE en el fallo gris del ejercicio 1. Durante ese intervalo, pedidos sigue llamando a inventario, cada llamada ocupa un hilo esperando, los hilos se agotan, Kong empieza a encolar, y un problema de un disco se convierte en una caída de toda la plataforma. Tolerar fallos no es solo sustituir lo que falla; es que quien llama a lo que falla no se hunda con ello. Eso son los patrones de resiliencia de la siguiente lección: timeouts con presupuesto, reintentos que no empeoran las cosas, el circuit breaker que 01-04 dejó pendiente, bulkheads, degradación controlada y load shedding.

Curso de Arquitecturas Distribuidas

Módulo 1: Introducción a los Sistemas Distribuidos

Módulo 2: Comunicación en Sistemas Distribuidos

Módulo 3: Consistencia y Replicación

Módulo 4: Almacenamiento Distribuido

Módulo 5: Computación Distribuida

Módulo 6: Seguridad en Sistemas Distribuidos

Módulo 7: Monitoreo y Mantenimiento

Módulo 8: Casos de Estudio y Aplicaciones

© Copyright 2026. Todos los derechos reservados