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
- Fallos en la práctica: del modelo teórico a lo que ocurre
- Detección: heartbeats, health checks y falsos positivos
- Tolerancia por redundancia
- Failover, split-brain y cercado
- Failover real: Patroni, Kafka y Cassandra
- Recuperación de estado: checkpoints de procesos largos
- Recuperación de datos: backups, PITR y pruebas de restauración
- RPO, RTO y plan de recuperación ante desastres
- Gestión de incidentes: runbooks, on-call y postmortems
- Errores comunes y consejos
- Ejercicios y soluciones
- Conclusión
- 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).
- 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.
- 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.
- 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:
- 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.
- 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
revisiondel 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.
- Token de cercado: cada líder recibe un número monótono creciente (la
- Timeouts generosos y
foren 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.
- 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.
- 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:
- El checkpoint se confirma en la misma transacción que el trabajo que representa. Si el proceso muere entre el
UPDATEy el checkpoint, la transacción no se confirma y ambos se pierden juntos; nunca queda un checkpoint que diga "hastaqueso-curado" conqueso-curadosin corregir, ni al revés. - El orden es determinista (
ORDER BY slug), de modo que "reanudar desde X" tiene significado. - 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.
- 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 minBackup 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 sobrescrituraPruebas 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.
- 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.
- 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
- Conceptos Básicos de Sistemas Distribuidos
- Modelos de Sistemas Distribuidos
- Ventajas y Desafíos de los Sistemas Distribuidos
- Las Falacias de la Computación Distribuida
- Tiempo, Relojes y Ordenación de Eventos
- Del Monolito a la Plataforma Distribuida: el Caso Kilómetro Cero
Módulo 2: Comunicación en Sistemas Distribuidos
- Protocolos de Comunicación
- RPC y RMI
- gRPC y Serialización de Datos
- Mensajería y Colas de Mensajes
- Patrones de Comunicación Asíncrona
Módulo 3: Consistencia y Replicación
- Modelos de Consistencia
- El Teorema CAP y PACELC
- Algoritmos de Consenso
- Replicación de Datos
- Transacciones Distribuidas y Sagas
Módulo 4: Almacenamiento Distribuido
- Particionado de Datos y Hashing Consistente
- Sistemas de Archivos Distribuidos
- Almacenamiento de Objetos
- Bases de Datos Distribuidas
- Cachés Distribuidos
Módulo 5: Computación Distribuida
- Modelos de Computación Distribuida
- MapReduce y Hadoop
- Spark y Computación en Memoria
- Procesamiento de Flujos de Datos
- Planificación de Trabajos y Pipelines de Datos
Módulo 6: Seguridad en Sistemas Distribuidos
- Autenticación y Autorización
- Cifrado y Protección de Datos
- Gestión de Identidades
- Seguridad entre Servicios: mTLS y Gestión de Secretos
- Puertas de Enlace, Limitación de Tasa y Auditoría
Módulo 7: Monitoreo y Mantenimiento
- Monitoreo de Sistemas Distribuidos
- Logs Centralizados y Trazabilidad Distribuida
- Gestión de Fallos y Recuperación
- Patrones de Resiliencia: Timeouts, Reintentos y Circuit Breaker
- Automatización y Orquestación
- Pruebas en Sistemas Distribuidos e Ingeniería del Caos
