Autores: Esteban Glochon, Matías Hernández, Sebastián Figueroa Retamal, Víctor Duarte Arce — Laboratorio de Diseño de Sistemas Intensivos en Datos, USACH.
Prueba de concepto distribuida que materializa, mediante contenedores, una parte relevante de la arquitectura de telemetría diseñada en el Laboratorio 1 (Diseño de Sistemas Intensivos en Datos, USACH), y la somete a fallos controlados. El alcance de este repositorio corresponde específicamente a lo exigido en el enunciado del Laboratorio 2. Es deliberadamente más acotado que el diseño conceptual de la Entrega 1 (que contemplaba múltiples regiones, líder múltiple y Protobuf), ya que aquí el objetivo es observar el comportamiento real de un sistema pequeño y reproducible, no escalarlo.
Los cuatro resultados de las pruebas de caos, con evidencia cruda completa en evidence/
(informe formal en evidence/informe-evidencia.md):
-
Failover automático de Shard 1 funciona. Al detener el primario con
docker stop,repmgrdpromovió un standby en ~3 segundos y la aplicación reanudó las escrituras 651 ms después de que terminó la promoción (02:40:46.233 → 02:40:46.884), sin ningún reinicio manual. Ventana total de indisponibilidad de escritura: ~27 segundos. Detalle en Replicación yevidence/repmgr-failover/. -
Split-brain reproducido bajo partición de red. Aislando al primario en minoría (1 contra 2), este siguió reportando
pg_is_in_recovery() = falsey aceptó unINSERTde prueba mientras el lado mayoritario promovía su propio primario: dos escritores simultáneos, confirmado con evidencia directa.REPMGR_PRIMARY_VISIBILITY_CONSENSUS=trueevita que un standby se promueva cuando el primario sigue siendo visible para otros nodos, pero no fuerza al primario aislado a dejar de aceptar escrituras — es una garantía distinta de la que ese flag ofrece. Detalle en Replicación yevidence/repmgr-failover/conclusion.md. -
Aislamiento de fallos entre shards, corregido. La versión inicial llamaba a la persistencia síncrona (psycopg2) directamente dentro del único agente Faust: un shard caído colgaba el consumo de todos los shards. Verificado apagando Shard 2 con tráfico real: la latencia de Shard 1 (sano) llegó a 424 segundos. Con una cola y un hilo de escritura independientes por shard (
processor/stream/app.py), el mismo escenario dejó a Shard 1 con ~0.5 segundos de atraso durante toda la caída. Detalle en Limitaciones conocidas. -
Nodo lento: el P99 no se movió. Inyectar 500 ms de latencia en una réplica de Shard 1 no degradó ninguna métrica del cliente (P99: 46.08 ms → 37.39 ms, dentro del ruido de un benchmark de ~20 s). Confirma que, con replicación asíncrona, la réplica queda fuera del camino crítico de escritura: el costo de esa decisión es durabilidad ante un failover, no latencia. Detalle en Escenario 1: nodo lento.
- Hallazgos principales
- Qué hace
- Arquitectura ejecutada
- Tecnologías y versiones
- Requisitos de hardware aproximados
- Cómo levantar y detener el stack
- Estructura del repositorio
- Particionamiento
- Replicación
- Modelo de consistencia y nivel de aislamiento
- Ruta modular seleccionada: procesamiento en stream
- Operación transaccional
- Escenarios de caos
- Limitaciones conocidas
- Herramienta de depuración del equipo
Un simulador genera telemetría de una flota de vehículos autónomos (posición, velocidad, batería) y la publica en Kafka. Un procesador en modo stream (Faust) la consume, detecta condiciones de riesgo (frenado de emergencia, exceso de velocidad sostenido) dentro de ventanas temporales, y persiste el resultado en dos shards independientes de PostgreSQL. Uno de esos shards está replicado. El repositorio incluye además los scripts para inyectar dos tipos de fallo (nodo lento, partición de red) sobre esa replicación y la evidencia recolectada al ejecutarlos.
Diagrama generado a partir de diagrams/arquitectura.puml. Para regenerarlo tras un cambio:
java -jar plantuml.jar -tpng diagrams/arquitectura.puml (requiere Java; el jar de PlantUML no
se incluye en este repositorio).
Cada evento se enruta por MD5(vehicle_id) % 2, calculado de forma independiente (pero idéntica)
tanto en el simulador como en el stream processor, de modo que ambos lados siempre coinciden en a
qué shard pertenece cada vehículo.
| Componente | Imagen / paquete | Versión |
|---|---|---|
| Broker de eventos | apache/kafka |
4.0.0 (KRaft, sin ZooKeeper) |
| Shard 1 (cluster repmgr, 3 nodos) | bitnamilegacy/postgresql-repmgr |
17.6.0 |
| Shard 2 | postgres |
fijada por digest (sha256:cd78ca5...) |
| Simulador | Python 3.9-slim |
kafka-python-ng >=2.0,<3.0, psycopg2-binary 2.9.12 |
| Stream processor | Python 3.11-slim |
faust-streaming >=0.13,<1.0, aiokafka >=0.10,<1.0, psycopg2-binary 2.9.12 |
| Fencing witness | Python 3.11-slim |
psycopg2-binary 2.9.12 |
| Orquestación | Docker Compose | v2 |
Las imágenes de Shard 2 y de Kafka están fijadas por digest o tag exacto. La imagen de Shard 1
(bitnamilegacy/postgresql-repmgr) quedó fijada por tag (17.6.0), no por digest, como pendiente
de seguimiento: Bitnami movió sus imágenes gratuitas de bitnami/ a bitnamilegacy/ en agosto de
2025 (bitnami/postgresql-repmgr quedó sin tags publicados), y al momento de escribir esto no se
verificó si bitnamilegacy/ republica esas imágenes con la misma estabilidad de versiones que
bitnami/; pineala por digest antes de depender de este repo a largo plazo.
Con los 5 vehículos simulados y carga de prueba normal, el stack completo (8 contenedores: kafka, 3 nodos de Shard 1, Shard 2, simulador, processor, fencing witness) consume del orden de 1.3 GB de RAM en régimen estable. Kafka concentra la mayor parte (~1 GB por overhead de JVM); el resto de los contenedores usa entre 10 y 75 MB cada uno (el fencing witness es el más liviano, ~9 MB). El uso de CPU es marginal (todos los contenedores bajo 15% en un núcleo). Se recomienda asignarle a Docker Desktop al menos 4 GB de RAM y 2 CPUs para tener margen durante los escenarios de caos y el benchmark de 1000 escrituras, más ~2 GB de disco para las imágenes.
# clonar y entrar al repo, luego:
docker compose up -d --build
# esperar a que todo quede healthy (~15-20s)
docker compose ps
# ver los logs en vivo
docker compose logs -f simulador
docker compose logs -f stream_processor
# detener sin borrar datos
docker compose down
# detener y borrar también los volúmenes (reinicio limpio)
docker compose down -vNo es necesario crear un archivo .env para el uso normal: todas las variables tienen valores por
defecto en docker-compose.yml (ver .env.example para personalizarlas, por ejemplo si algún
puerto local ya está ocupado).
telemetry-lab-2/
├── docker-compose.yml # kafka, Shard 1 (3 nodos repmgr), Shard 2, simulador, processor
├── .env.example # variables de entorno con sus defaults
├── database/schema.sql # esquema aplicado en los 4 nodos de Postgres
├── shared/
│ └── primary_discovery.py # descubre cual nodo de Shard 1 es primario ahora mismo
├── fencing-witness/ # vigila Shard 1 y fuerza solo-lectura si detecta 2 primarios
├── simulator/ # genera telemetría y la publica en Kafka
├── processor/
│ ├── transaction.py # la operación transaccional (invariante de la Sección 4)
│ ├── test_transaction.py # prueba de rollback/idempotencia/evento atrasado
│ └── stream/ # el consumidor Faust (ruta modular elegida)
├── chaos/ # scripts de inyección/recuperación de fallos + benchmark
│ └── kill-shard1-primary.sh # identifica y detiene al primario actual de Shard 1
├── evidence/
│ ├── informe-evidencia.md # informe de evidencia del enunciado §8.3 (ambos escenarios)
│ ├── baseline/ # captura de estado general del clúster
│ ├── transactions/ # log de la prueba transaccional
│ ├── slow-node/ # escenario 6.1 (histórico, topología de 2 nodos)
│ ├── network-partition/ # escenario 6.2 (histórico, topología de 2 nodos)
│ ├── repmgr-failover/ # failover automático y split brain, topología actual
│ └── fencing/ # intento de cerrar la brecha de split brain (parcial)
├── queries/ # SQL de inspección (sharding, replicación)
├── diagrams/ # diagramas de arquitectura y de la transacción (PlantUML)
└── debug-dashboard/ # herramienta de depuración del equipo, no es parte de la entrega
La clave de partición es MD5(vehicle_id) % 2: hash par va a Shard 1, hash impar va a Shard 2. Se
eligió vehicle_id (no región ni modelo) porque es el campo de mayor cardinalidad disponible y
porque todo el estado que necesita mantenerse por vehículo (ventanas de frenado/velocidad, último
estado conocido) vive naturalmente en un único shard, evitando consultas entre nodos para el camino
de escritura. El costo de esta decisión es que las consultas agregadas por región sí requieren
combinar resultados de ambos shards.
Con los 5 vehículos de la configuración de prueba (simulator/config.py), la distribución
resultante es 2 vehículos en Shard 1 (veh-001, veh-004) y 3 en Shard 2 (veh-002, veh-003,
veh-005). Es verificable con queries/count_by_shard.sql o evidence/baseline/collect.sh. Con
una muestra tan chica, el desbalance de 2 contra 3 es esperable por azar del hash y no se considera
evidencia de un problema de la función de partición, que con más vehículos tiende a uniformarse (es
la misma lógica de hashing ya justificada, con mayor detalle estadístico, en la Entrega 1).
Shard 1 tiene factor de replicación 3 (db_shard_1_node1/2/3), gestionado por
repmgr (imagen bitnamilegacy/postgresql-repmgr) con failover
automático (REPMGR_FAILOVER=automatic, el valor por defecto). Shard 2 no tiene réplica; se dejó
así intencionalmente para cumplir el mínimo del enunciado ("al menos un componente crítico
replicado") sin duplicar infraestructura que no aporta información adicional al experimento.
Por qué 3 nodos y no 2. Con failover automático, un cluster de 2 nodos corre el riesgo de que
un nodo aislado por una partición de red se autopromueva sin saber que quedó en minoría,
exactamente el split brain que se evitaba en el diseño anterior (líder-seguidor sin failover). Con
3 nodos y REPMGR_PRIMARY_VISIBILITY_CONSENSUS=true, la intención es que un nodo aislado reconozca
que no tiene visibilidad de la mayoría del cluster y dejar de aceptar escrituras. Esto se probó
explícitamente y el resultado fue negativo: ver más abajo y evidence/repmgr-failover/.
Quién es el primario ya no es un dato fijo. A diferencia del diseño anterior
(db_shard_1_leader como hostname fijo), el nodo primario de Shard 1 puede cambiar en cualquier
momento tras un failover. La aplicación descubre cuál nodo es escribible en tiempo de ejecución
(shared/primary_discovery.py): prueba cada candidato con SELECT pg_is_in_recovery();, cachea el
resultado por unos segundos, y si una escritura falla invalida el caché de inmediato y reintenta
contra un primario recién descubierto, sin esperar a que venza el caché.
La replicación de Shard 1 sigue siendo asíncrona (repmgr no cambia esto: se verificó en
evidence/slow-node/, histórico pero con conclusión todavía vigente, que la latencia de escritura
no se ve afectada por una réplica lenta). Esto implica que, si el primario falla, se pueden perder
las transacciones que haya confirmado pero que ningún standby haya replicado todavía.
Se probaron dos escenarios reales sobre el cluster de 3 nodos (evidence/repmgr-failover/):
- Caída limpia del primario (
docker stop): funciona como se espera.repmgrdpromovió un standby a primario en ~27 segundos, y el stream processor detectó el fallo, redescubrió al nuevo primario y retomó la escritura sin ningún reinicio manual, en el mismo segundo en que terminó la promoción. - Partición de red del primario (queda aislado, en minoría 1 contra 2): no funcionó como
se esperaba. El primario aislado siguió reportando
pg_is_in_recovery() = falsey aceptó una escritura de prueba mientras estaba aislado, al mismo tiempo que el lado mayoritario promovía su propio primario nuevo: dos primarios simultáneos, split brain real, no solo teórico. Al reconectar la red, la reconciliación tampoco fue automática ni inmediata: el nodo divergente necesitó un reinicio/reclonado para volver a un estado consistente como standby, descartando el dato escrito durante el split brain.
En síntesis: REPMGR_PRIMARY_VISIBILITY_CONSENSUS=true, tal como se configuró en este PoC, no
fue suficiente para impedir el split brain durante una partición activa: ese flag evita que un
standby se promueva cuando el primario sigue siendo visible para otros nodos, pero no fuerza al
primario aislado a dejar de aceptar escrituras (confirmado leyendo la documentación de repmgr, no
solo por la prueba).
Se agregó un servicio adicional (fencing_witness, ver fencing-witness/witness.py) que vigila
los 3 nodos y, si detecta más de uno declarándose primario a la vez, fuerza a los perdedores a modo
solo lectura (ALTER SYSTEM SET default_transaction_read_only = on + reload). El mecanismo en sí
se probó y funciona: aplicado directamente contra un primario real, bloquea escrituras nuevas con
el mismo error que produce un standby genuino, sin necesidad de reiniciar PostgreSQL.
Pero no cerró la brecha para el caso general: el testigo corre como un contenedor más dentro de
la misma red Docker que los nodos de Shard 1, así que una partición de red que aísla a un nodo lo
aísla también del testigo, el mismo punto ciego que tendría la aplicación real. Repitiendo la
prueba con el testigo desplegado, el análisis honesto quedó así (detalle completo en
evidence/fencing/conclusion.md):
- El camino de escritura que sí logró colarse en
evidence/repmgr-failover/fue específicamente vía el puerto publicado al host, un camino que la aplicación real (dentro de otro contenedor en la misma red interna) nunca usa. El riesgo confirmado es más acotado de lo que sonaba al principio. - Se descubrió, sin buscarlo, que
repmgrdsí tiene su propia salvaguarda de auto-terminación ("degraded monitoring timeout") que apaga al nodo aislado por su cuenta, aunque con timing variable e impredecible en las pruebas (6 segundos en una corrida, más de un minuto en otra). No es un fencing instantáneo garantizado.
Conclusión honesta: migrar a failover automático cambió el trade-off del diseño anterior, no lo eliminó. Se ganó disponibilidad automática tras perder el primario; se perdió la garantía estructural (no negociable, a nivel de protocolo de PostgreSQL) que sí tenía el diseño líder-seguidor sin failover. Cerrar esa brecha para el caso general (un cliente con ruta de red independiente hacia el nodo aislado) requeriría fencing a nivel de red real o un testigo desplegado fuera de la red interna del cluster, trabajo de infraestructura fuera del alcance razonable de este laboratorio.
Estas son dos propiedades distintas y se declaran por separado:
- Consistencia entre réplicas de Shard 1: eventual, consecuencia directa de la replicación asíncrona descrita arriba. Una lectura contra el follower puede devolver una versión desactualizada respecto del líder; qué tan desactualizada depende de cuánto WAL le falte por aplicar.
- Nivel de aislamiento transaccional:
READ COMMITTED, el valor por defecto de PostgreSQL. No se modificó explícitamente para ninguna de las operaciones de este proyecto. Esto es lo que rige qué cambios concurrentes puede ver una transacción dentro de un mismo nodo, y es independiente de si esa transacción, una vez confirmada, ya llegó o no a las réplicas.
Sobre split brain: con el diseño anterior (líder-seguidor de 2 nodos, sin failover automático),
un nodo aislado por una partición de red no podía aceptar escrituras en absoluto: PostgreSQL lo
restringía a solo lectura por estar en modo standby, así que el split brain era estructuralmente
imposible. Con la migración a repmgr (Shard 1 = 3 nodos, failover automático), esa garantía
ya no aplica: se probó explícitamente y un primario aislado por partición de red sí llegó a
aceptar una escritura mientras el lado mayoritario promovía su propio primario nuevo. Ver
"Failover automático" arriba y evidence/repmgr-failover/conclusion.md para el detalle completo.
Se implementó la Opción A del enunciado, con Faust sobre Kafka.
- Partición del stream: por
vehicle_id(misma clave que el sharding de base de datos), lo que permite mantener el estado de ventana (última velocidad, contador de exceso de velocidad sostenido) coherente por vehículo. - Reglas implementadas:
- Frenado de emergencia: caída de velocidad ≥ 40 km/h dentro de una ventana de 3 segundos.
- Exceso de velocidad sostenido: velocidad por sobre el límite de la vía durante 3 o más eventos consecutivos.
event_timevs.processing_time: cada evento lleva su propioevent_timegenerado por el simulador; el processor calcula la latencia contra el momento real de procesamiento y marca como tardío (late_events) todo evento cuya latencia supere 30 segundos. Un evento tardío igual se procesa (no se descarta), pero no reemplaza un estado de vehículo más reciente ya persistido (ver invariante transaccional más abajo).- Persistencia: cada evento se escribe transaccionalmente en el shard correspondiente
(
processed_events,vehicle_state,telemetry_eventsy, si corresponde,alerts), y las alertas además se publican en el tópico Kafkaalerts. - Métricas: expuestas en
http://localhost:6066/metrics, incluyendo eventos procesados, alertas generadas, eventos tardíos, latencia de procesamiento (avg/P95/P99), latencia de escritura a BD (avg/P95/P99) y consumer lag por partición.
Definida en processor/transaction.py. Dentro de una misma transacción se inserta el event_id en
processed_events y se actualiza vehicle_state; si falla cualquier paso intermedio, PostgreSQL
revierte ambos cambios y no queda un estado parcial. La cláusula ON CONFLICT DO NOTHING sobre
processed_events hace que reintentar un evento ya procesado sea un no-op, y la condición
vehicle_state.event_time < EXCLUDED.event_time evita que un evento atrasado sobrescriba un
estado más reciente.
python processor/test_transaction.pyCorre la prueba exigida en la Sección 4.2 (fallo simulado antes de actualizar vehicle_state,
rollback, reintento, verificación de efecto único) y deja el resultado en
evidence/transactions/test_transaction.log.
Ambos escenarios obligatorios del enunciado están implementados en chaos/ y ejecutados con
evidencia completa en evidence/. El informe formal de evidencia exigido por la sección 8.3 del
enunciado (hipótesis, configuración, comandos, línea temporal, métricas antes/durante/después,
logs, consultas de validación y conclusión causal de cada escenario) está en
evidence/informe-evidencia.md.
# nodo lento: agrega latencia de red dentro de un contenedor
./chaos/slow-node.sh <container> [delay_ms] # ej: ./chaos/slow-node.sh db_shard_1_node2 500
# partición de red: aísla un contenedor de la red del proyecto
./chaos/network-partition.sh <container> [network] # ej: ./chaos/network-partition.sh db_shard_1_node2
# revierte cualquiera de los dos fallos anteriores sobre el contenedor indicado
./chaos/recover.sh <container> [network]
# failover: identifica y detiene al primario actual de Shard 1
./chaos/kill-shard1-primary.sh
# benchmark de la operación representativa (escritura transaccional real)
python chaos/bench_write.py --host localhost --port 5432 --count 1000 --out resultado.json
# recorrido guiado con pausas, para mostrar los escenarios en vivo
./chaos/demo.shCriterio de recuperación. Para slow-node.sh/network-partition.sh, "recuperado" significa
que tc qdisc show dev eth0 en el contenedor vuelve a mostrar la disciplina por defecto
(noqueue, sin netem) y que el contenedor figura de nuevo entre los conectados a la red del
proyecto (docker network inspect telemetry-lab-2_telemetry_network), sin intervención manual
adicional más allá de ejecutar recover.sh. Para kill-shard1-primary.sh, "recuperado" significa
que repmgr.nodes (consultable vía psql -U repmgr -d repmgr) vuelve a mostrar 3 nodos activos
con exactamente un primario, lo cual, según el escenario probado (ver abajo), puede requerir
reiniciar el nodo detenido en vez de resolverse solo.
Protocolo ejecutado según el enunciado (hipótesis → línea base → inyección → observación →
recuperación → validación), con evidencia completa en evidence/slow-node/: 1.000 escrituras
transaccionales de línea base contra Shard 1, inyección de 500ms de latencia en un nodo
(./chaos/slow-node.sh db_shard_1_node2 500), repetición de la misma carga y restauración con
recover.sh. La medición se realizó sobre la topología que Shard 1 tenía en ese momento
(líder-seguidor de 2 nodos); la conclusión sigue vigente para el cluster repmgr actual porque la
replicación sigue siendo asíncrona (ver Replicación).
Métricas (detalle en 01_baseline_write.json y 03_injected_write.json):
| Línea base | Con 500ms en la réplica | |
|---|---|---|
| Mediana | 19.49 ms | 16.73 ms |
| P95 | 33.06 ms | 31.93 ms |
| P99 | 46.08 ms | 37.39 ms |
| Errores / timeouts | 0 / 0 | 0 / 0 |
| Throughput | 47.8 ops/s | 51.9 ops/s |
El lag líder-follower, medido por comparación directa de LSN antes, durante y después de la inyección, fue idéntico en los tres puntos: no se observó crecimiento de lag con esta carga.
Preguntas de análisis del enunciado:
¿El P99 aumentó aunque la mediana se mantuviera estable?
No: ni la mediana ni el P99 aumentaron; ambas cifras bajaron levemente, variación atribuida al ruido de un benchmark de ~20 segundos en un host compartido, no a una mejora real. Ese patrón (mediana estable con P99 disparado) es justo lo que se habría observado si la réplica lenta hubiera estado en el camino crítico; su ausencia es el hallazgo relevante.
¿La réplica lenta quedó fuera del camino crítico?
Sí. El cliente escribe únicamente contra el primario; el streaming del WAL hacia el standby es un
flujo desacoplado que no participa en la confirmación del COMMIT.
¿Qué configuración determinó ese comportamiento?
La replicación asíncrona: no existe synchronous_standby_names configurado y repmgr no cambia eso.
Con replicación síncrona sobre esa réplica, cada COMMIT habría debido esperar su confirmación y
los 500ms inyectados se habrían reflejado directamente en la mediana y el P99.
¿La política elegida protege mejor la latencia o la consistencia?
La latencia. El costo asumido es durabilidad: si el primario falla, se pierden las transacciones confirmadas que ningún standby haya llegado a replicar. Es el trade-off deliberado para telemetría (eventos inmutables), ya adoptado en la Entrega 1.
¿Qué ocurriría si dos réplicas se volvieran lentas?
Para el cliente, nada: mientras ninguna figure en synchronous_standby_names, el primario no
espera a ninguna réplica, sin importar cuántas estén lentas. Lo que crecería es el lag acumulado de
cada una y, con ello, la ventana de pérdida ante un failover. En el cluster actual de 3 nodos, dos
standbys muy retrasados dejarían la durabilidad efectiva en un único nodo.
¿Cómo distinguir un nodo lento de uno completamente caído?
Un nodo caído termina su proceso walreceiver y desaparece de pg_stat_replication (además de
rechazar conexiones); un nodo lento permanece conectado, con lag creciente visible en write_lag,
flush_lag y replay_lag. Limitación conocida de esta implementación: telemetry_user no tiene
el privilegio pg_monitor, así que esas columnas no son visibles para la aplicación, que usa
comparación directa de LSN como método alternativo (funcional, aunque menos directo).
En síntesis: la política asíncrona saca a la réplica lenta del camino crítico (ni el líder ni un quórum esperan por ella), el P99 medido no cambió, y el costo de la decisión se manifiesta como ventana de pérdida potencial ante failover, no como latencia adicional.
Escenario ejecutado sobre la topología actual (cluster repmgr de 3 nodos), con evidencia completa
en evidence/repmgr-failover/. Línea base: node-2 era el primario y node-1/node-3 sus
standbys (03_scenario_b_baseline.txt). Se aisló node-2 con
./chaos/network-partition.sh db_shard_1_node2, quedando en minoría (1 contra 2), con el simulador
manteniendo tráfico durante toda la ventana (~40 segundos). Como contraste, la caída limpia del
primario (docker stop) sí se recupera sola en ~27s sin intervención manual; lo específico de la
partición es lo que sigue.
Resultados por lado, con el formato exigido por el enunciado:
| Lado | ¿Mayoría? | Operación | Resultado | Interpretación |
|---|---|---|---|---|
A: node-1 + node-3 |
Sí | Escritura | Aceptada: el lado mayoritario promovió a node-1 a primario y siguió aceptando escrituras |
Disponibilidad de escritura preservada por el lado mayoritario |
A: node-1 + node-3 |
Sí | Lectura reciente | Reciente: node-1, ya promovido, responde como primario con estado actualizado |
Garantía de lectura intacta del lado mayoritario |
B: node-2 (aislado) |
No | Escritura | Aceptada: pg_is_in_recovery() = f e INSERT 0 1 durante el aislamiento (05_isolated_primary_write_attempt.txt) |
Divergencia materializada: dos nodos aceptando escrituras simultáneamente |
B: node-2 (aislado) |
No | Lectura local | Reciente respecto de su propio estado (incluye el INSERT divergente); divergente respecto del clúster real | Costo de operar sin fencing alcanzable desde la minoría |
¿Hubo split brain?
Sí, y esta vez el término aplica literalmente. Una partición de red no es automáticamente split
brain: solo lo es cuando dos subconjuntos aceptan simultáneamente operaciones incompatibles como si
ambos fueran autoritativos. Eso fue exactamente lo observado: node-2 aislado seguía siendo
primario a nivel PostgreSQL (el aislamiento de red no cambia por sí solo pg_is_in_recovery()),
aceptó un INSERT directo, y al mismo tiempo el lado mayoritario promovía su propio primario. Dos
escritores simultáneos, confirmado con evidencia directa y no solo en teoría. En la topología
anterior de 2 nodos esto era estructuralmente imposible (un standby aislado rechazaba escrituras
por estar en recuperación; allí correspondía hablar de partición sin split brain, ver
evidence/network-partition/): la diferencia la determina la configuración de liderazgo, no la
partición en sí.
¿Por qué no lo evitó la configuración? REPMGR_PRIMARY_VISIBILITY_CONSENSUS=true evita que un
standby se promueva cuando el primario sigue visible para otros nodos, pero no fuerza al
primario aislado a dejar de aceptar escrituras (verificado también contra la documentación de
repmgr). Se intentó cerrar la brecha con un testigo de fencing (fencing-witness/, detalle en
evidence/fencing/conclusion.md): el mecanismo funciona cuando es alcanzable, pero comparte la red
Docker con los nodos que vigila, así que la partición que aísla a un nodo lo aísla también del
testigo. El riesgo confirmado quedó acotado al camino del puerto publicado al host, camino que la
aplicación real (dentro de la red interna) nunca usa. Adicionalmente, repmgrd tiene una
auto-terminación por "degraded monitoring" que actúa como mitigación parcial, pero con timing
variable e impredecible entre corridas (6 segundos en una, más de un minuto en otra): no es un
fencing instantáneo ni garantizado.
Análisis PACELC (referido a la configuración ejecutada, no a una etiqueta genérica del producto):
- Durante la partición (P): el sistema priorizó disponibilidad (A) en ambos lados al costo de consistencia (C): los dos lados siguieron aceptando escrituras sin coordinación, que es precisamente la condición del split brain observado. Un diseño con quórum real de escritura (W > n/2) habría rechazado al lado minoritario.
- Fuera de la partición (E): prioriza latencia (L) sobre consistencia fuerte (EL), consecuencia directa de la replicación asíncrona medida en el escenario 1: escribir es rápido aunque las réplicas vayan atrás.
Recuperación y convergencia. Revertida la partición (recover.sh), la reconciliación no fue
automática: node-2 siguió reportándose primario hasta que su propio repmgrd se autoterminó
(después de la reconexión, no durante el aislamiento activo), y el contenedor necesitó
reinicio/reclonado: el arranque detectó "This node was acting as a primary before restart!" y se
reclonó desde node-1, descartando su estado divergente. Se verificó que el INSERT hecho durante
el split brain no sobrevivió en el nodo ganador. La convergencia final se alcanzó, pero
requirió intervención manual y a costa de perder la escritura divergente.
-
Aislamiento de fallos entre shards, resuelto con colas de escritura por shard. El stream processor consumía el tópico con un único agente Faust que procesaba un evento a la vez, llamando a
persist_transactional()(psycopg2, síncrona) directamente dentro de ese agente. Sin timeout de conexión, un shard caído podía colgar indefinidamente el consumo de todos los shards. Conconnect_timeout=3(fix anterior) el sistema dejaba de colgarse para siempre, pero seguía degradándose sin límite: verificado en vivo apagandodb_shard_2,consumer_lagllegó a 783 eventos y la latencia deveh-001(Shard 1, el shard sano) a 424s, porque el throughput degradado (<1 evento/s) quedaba por debajo de la tasa de producción real (~5 eventos/s) y la cola de trabajo pendiente crecía indefinidamente para todos los vehículos, no solo los del shard caído.Se corrigió de raíz: cada shard tiene ahora su propia
asyncio.Queuey su propia corrutina de escritura (_shard1_writer/_shard2_writerenprocessor/stream/app.py), cada una con su propio hilo (ThreadPoolExecutor) para la llamada bloqueante a psycopg2. El agente principal ya no espera a que termine la escritura antes de consumir el siguiente evento de Kafka: solo encola (queue.put_nowait) y sigue. Verificado en vivo repitiendo el mismo escenario (apagardb_shard_2con tráfico real corriendo): la cola de Shard 2 creció (143 eventos en 30s, esperado) mientras la de Shard 1 se mantuvo en 0 todo el tiempo, yveh-001/veh-004(Shard 1) se mantuvieron con ~0.5 segundos de atraso durante toda la caída de Shard 2, sin degradación medible. El endpoint/metricsexponequeue_depthpor shard para observar esto en vivo.Trade-off aceptado: la garantía de entrega se debilita levemente: si el proceso muere con eventos todavía en una cola en memoria, esos eventos se pierden (hay un intento de drenado en el shutdown,
on_before_shutdown, con 10s de margen, pero un crash duro no pasa por ahí). No es una categoría de riesgo nueva: ya existía la misma postura para las Tablas de Faust (store="memory://"). Un evento que agota los 2 reintentos depersist_transactional()ya no se pierde con solo un log: se escribe a un archivo.jsonlde dead-letter por shard (DEAD_LETTER_DIR, por defecto/app/dead_letter/, ver_write_dead_letter()enprocessor/stream/app.py) para poder inspeccionarlo o reprocesarlo después. Las colas están acotadas (WRITE_QUEUE_MAXSIZE = 10_000, más de 30 minutos de outage al ritmo de carga actual) para no crecer sin límite en memoria ante una caída muy prolongada; al llenarse (asyncio.QueueFull, un caso distinto al de reintentos agotados), el evento todavía se descarta con solo un log, algo que solo se alcanza en outages excepcionalmente largos.Nota aparte: las cifras de latencia promedio/P95/P99 en
/metricsquedan infladas un buen rato después de un incidente así porque son acumulativas y no se resetean solas; los indicadores confiables de salud en tiempo real sonconsumer_lagyqueue_depth, no esas cifras. -
El modelo de movimiento del simulador limita la velocidad al 95% del límite de la vía y suaviza la aceleración, por lo que nunca genera de forma natural un frenado de emergencia ni un exceso de velocidad sostenido: las alertas requieren inyectar un evento de prueba a mano para poder observarse.
-
El stream processor usa
store="memory://"en Faust: el estado de ventana (última velocidad conocida, contador de exceso de velocidad) se pierde si el processor se reinicia. El estado persistido en PostgreSQL (vehicle_state) no se ve afectado, pero la detección de alertas puede perder continuidad justo después de un reinicio. -
El rol
telemetry_userno tiene el privilegiopg_monitor, por lo que las columnas de estado y lag depg_stat_replication/pg_stat_wal_receiverno son visibles para la aplicación. Como alternativa se usa la comparación directa de LSN entre líder y follower. -
Shard 2 no tiene réplica: solo Shard 1 cumple el requisito de replicación del enunciado.
-
El failover automático de Shard 1 no impide split brain durante una partición de red, ni siquiera con el testigo de fencing agregado. Probado explícitamente (
evidence/repmgr-failover/,evidence/fencing/): un primario aislado por partición siguió aceptando escrituras mientras el lado mayoritario promovía su propio primario, conREPMGR_PRIMARY_VISIBILITY_CONSENSUS=trueconfigurado. Se agregófencing_witness(mecanismo de fencing verificado, funciona cuando es alcanzable), pero comparte la red Docker con los nodos que vigila, así que una partición que aísla a un nodo lo aísla también del testigo, el mismo punto ciego que tendría la aplicación real. El riesgo confirmado quedó acotado al camino del puerto publicado al host (no al camino que usa la aplicación). Cerrar la brecha para el caso general requeriría fencing a nivel de red real o un testigo con ruta de red independiente, fuera del alcance razonable de este laboratorio. -
La reconciliación tras un split brain no es automática. El nodo que quedó con datos divergentes no volvió a ser standby por sí solo al reconectar la red; necesitó un reinicio/reclonado del contenedor. Con volúmenes de datos persistentes (no usados en este PoC) el procedimiento sería más delicado que simplemente recrear el contenedor.
-
La imagen
bitnamilegacy/postgresql-repmgrestá fijada por tag, no por digest (ver "Tecnologías y versiones"), pendiente de endurecer si este repositorio se usa a más largo plazo.
debug-dashboard/ es un cliente de solo lectura contra los puertos que Docker ya publica al
host, pensado para que el equipo vea en vivo el estado de la flota, los shards y el cluster de
Shard 1 mientras se prueba el stack o se ejecutan los escenarios de chaos/. Instrucciones de uso en debug-dashboard/README.md.

