🌐 Sistemas distribuidos de élite
Consenso real (Raft en detalle), los papers fundacionales (Dynamo, Spanner, GFS, Bigtable, Kafka), y los sistemas reales (Kafka, Cassandra, Redis). Teoría, código y práctica.
Sistemas distribuidos de élite
Este es el conocimiento que separa al senior del staff: los papers y los sistemas reales que sostienen internet. Se estudia leyendo fuentes primarias y construyendo. La base conceptual (CAP, replicación, particionado, consistencia) ya está en Fundamentos: distribuidos — aquí bajamos al detalle de cómo funcionan de verdad los sistemas que usas.
Los sistemas distribuidos no se aprenden en tutoriales: se aprenden leyendo papers, el código de sistemas reales (Kafka, Cassandra) y haciendo los labs. Cada paper de esta lista cambió la industria: dedícales tiempo y toma notas.
Raft a fondo: el consenso que entiendes
Ya viste la idea en Fundamentos: distribuidos. Aquí, los detalles que importan para entender un sistema real (etcd, Consul, CockroachDB usan Raft).
Los tres subproblemas
- Elección de líder: decidir QUIÉN manda (con mayoría de votos).
- Replicación de log: el líder replica su log y los seguidores lo aplican.
- Seguridad: garantizar que solo haya un líder válido por término.
Términos y reglas
- Cada elección tiene un término (número de elección). Los votos valen para un término.
- Un nodo solo vota una vez por término, al primer candidato que le pida y cuyo log esté al día.
- Mayoría (quórum) = n/2 + 1: con 3 nodos necesitas 2 votos; con 5, 3.
3 nodos: el líder necesita 2 de 3 confirmaciones
5 nodos: el líder necesita 3 de 5
Por eso 2 nodos es peor que 1 (necesitas 2 de 2: un solo fallo paraliza todo)
y 3 nodos tolera 1 fallo, 5 tolera 2.
Replicación del log
[Líder] term=3
│ AppendEntries (log: [1:x] [2:y] [3:z])
├──▶ [Seguidor A] → confirma
├──▶ [Seguidor B] → confirma
└── mayoría confirmada → COMMIT → aplica al state machine
💡 La regla de seguridad clave: un nodo solo es elegible como líder si su log está al día (tiene todos los commits de términos anteriores). Eso evita que un líder con log viejo sobrescriba datos ya confirmados.
El clúster en la práctica
# Un clúster etcd (Raft) de 3 nodos
etcd --name n1 --initial-cluster n1=http://10.0.0.1:2380,n2=...,n3=...
etcdctl member list
# mata el líder → el clúster elige otro en segundos (tolerancia a 1 fallo)
⚠️ Split-brain y la mayoría: si la red se parte en 2 y 1, el lado con 2 nodos tiene mayoría y sigue; el lado con 1 pierde el quórum y deja de escribir (evita dos líderes). Eso es CP: prefieres no escribir a escribir dos veces.
Kafka: el log distribuido que mueve el mundo
Kafka es el broker de eventos dominante. Su modelo: un log distribuido, particionado y replicado.
Partitions y offset
Cada topic se divide en particiones (shards del log). Los mensajes de cada partición tienen un offset (posición). El orden se garantiza dentro de una partición, no global.
Topic "pedidos" (3 particiones):
P0: [pedido_1][pedido_4][pedido_7] → ordenados
P1: [pedido_2][pedido_5][pedido_8]
P2: [pedido_3][pedido_6][pedido_9]
| Concepto | Qué es |
|---|---|
| Topic | Categoría de eventos |
| Partition | Unidad de paralelismo y de orden |
| Offset | Posición del mensaje en la partición |
| Producer | El que publica |
| Consumer group | Grupo de consumidores que reparten las particiones |
| Replication factor | Cuántas copias de cada partición (3 = tolera 2 fallos) |
| ISR | Réplicas «in sync» (las que están al día) |
| Retention | Cuánto tiempo se guardan los mensajes (configurable) |
Consumer groups: repartir el trabajo
Un consumer group reparte las particiones entre sus miembros: cada consumidor procesa un subconjunto.
Topic con 4 particiones, grupo de 2 consumidores:
C1 → P0, P1
C2 → P2, P3 ← reparto automático (rebalanceo)
💡 Kafka garantiza que un mensaje se entrega a una sola instancia del group (procesan en paralelo sin duplicar el trabajo) y que dentro de una partición el orden se mantiene. Para orden global, usa UNA partición o un id de partición por clave (
key → partition).
Exactly-once vs at-least-once
| Garantía | Qué significa | Coste |
|---|---|---|
| At-most-once | El mensaje se pierde a veces (no se reenvía) | Bajo |
| At-least-once | El mensaje puede duplicarse, nunca perderse | Medio (necesita idempotencia) |
| Exactly-once | Cada mensaje se procesa exactamente una vez | Alto (transacciones en el stream) |
⚠️ En la práctica, el estándar es at-least-once + consumidores idempotentes. El exactly-once «de verdad» (transacciones de Kafka) es caro y complejo: la mayoría de sistemas lo evitan.
Los papers fundacionales
Leer los papers originales es el rito del nivel élite. El orden recomendado:
1. Dynamo (Amazon, 2007) — el origen del NoSQL AP
Diseñado para el carrito de compras de Amazon: disponibilidad máxima (AP). Las ideas:
- Consistent hashing: cada nodo guarda un rango del anillo de hashes; añadir/quitar nodos solo mueve una fracción de los datos.
- Quórum: escrituras en
Wnodos, lecturas enRnodos; siR + W > Nsiempre lees lo último. - Vector clocks: detectan escrituras concurrentes y sus conflictos.
- Gossip: los nodos se pasan el estado por rumor (sin coordinador central).
- Hinted handoff: si un nodo cae, otro guarda sus escrituras y se las pasa cuando vuelve.
Anillo de consistent hashing (N=3 réplicas):
Clave "carrito:123" → hash → posición → se replica en los 3 nodos siguientes
Añadir nodo X: solo se mueven las claves de los vecinos de X
💡 Legado: DynamoDB (AWS), Cassandra, Riak y Voldemort descienden de este paper. Es la respuesta a «¿cómo se construye un AP real?».
2. Spanner (Google, 2012) — la DB globalmente consistente
Google necesitaba una base de datos global (multi-región) con consistencia fuerte y transacciones. El truco revolucionario: TrueTime.
- TrueTime: un reloj global con sincronización por GPS y relojes atómicos, que devuelve un intervalo
[earliest, latest]en vez de un instante. - Con TrueTime, Spanner puede ordenar transacciones a nivel global (consistencia externa) usando Paxos para replicación.
- Resultado: CP global: datos siempre consistentes en todo el planeta, a costa de latencia.
Time: [earliest ──(TrueTime)── latest] → espera a que el intervalo pase
→ garantiza que la transacción está "en el pasado" global
💡 Legado: Cloud Spanner, CockroachDB (open source con ideas similares, usa HLC — hybrid logical clocks). Es la respuesta a «¿cómo se construye un CP global?».
3. GFS / Bigtable / MapReduce — la tríada de Google
- GFS (Google File System): un filesystem distribuido diseñado para fallos constantes y archivos enormes. Escribe con réplicas y leases de escritura.
- Bigtable: el almacén columnar que dio origen a HBase y Cassandra. Combina el log (GFS) con tablas ordenadas (LSM).
- MapReduce: el modelo de cómputo distribuido (map → shuffle → reduce) que hizo escalable el procesamiento (→ Hadoop, y luego Spark/Flink).
4. Kafka (LinkedIn, 2011)
El paper del log como backbone: ya lo vimos arriba. Cambió el event-driven para siempre.
Los sistemas reales: cómo operan
Cassandra: AP descentralizado
Cassandra desciende de Dynamo: sin líder, quórum, gossip, consistencia configurable.
-- Cassandra: consistencia configurable por query
SELECT * FROM usuarios WHERE id = 42 CONSISTENCY QUORUM;
INSERT INTO usuarios (id, nombre) VALUES (42, 'Ana') CONSISTENCY ONE;
| Fuerza | Uso |
|---|---|
| Escrituras masivas en cualquier nodo | Telemetría, IoT, eventos |
| Sin punto único de fallo | Multi-región |
| Tunable consistency | Eliges por query entre ONE / QUORUM / ALL |
Redis: el sistema in-memory
No es «solo una caché»: es un servidor de estructuras de datos en memoria (strings, lists, sets, hashes, sorted sets, streams, pub/sub).
SET usuario:42 '{"nombre":"Ana"}' EX 300 # con TTL
INCR contador # atómico
LPUSH cola "tarea" # listas como cola
ZADD ranking 100 "usuario" # sorted set
PUBLISH canal "mensaje" # pub/sub
| Patrón | Uso |
|---|---|
| Caché | Cache-aside con TTL |
| Sesiones | Guardar sesiones compartidas entre instancias |
| Rate limiting | INCR + TTL por ventana |
| Colas ligeras | Listas / streams |
| Pub/sub | Notificaciones en tiempo real |
| Semaforos / locks | SETNX (distributed lock) |
Kafka vs RabbitMQ: cuándo cada uno
| Kafka | RabbitMQ | |
|---|---|---|
| Modelo | Log distribuido (replay, múltiples consumidores) | Cola de mensajes (competing consumers) |
| Orden | Por partición | Por cola |
| Persistencia | Alta (retention) | Media (por config) |
| Uso | Event-driven, streams, pipelines | Task queues, request/reply, routing complejo |
💡 Regla rápida: ¿es un evento que debes poder reproducir y leer varios consumidores? → Kafka. ¿Es una tarea que un worker debe procesar una vez? → RabbitMQ/SQS.
MIT 6.824/6.5840: la formación hands-on
Los labs del curso de MIT son la forma definitiva de aprender:
| Lab | Qué construyes |
|---|---|
| Lab 1 | MapReduce (workers + coordinador) |
| Lab 2 | Raft en Go (elección, replicación, persistencia) |
| Lab 3 | KV service con Raft (transacciones, particiones) |
| Lab 4 | Sharded KV store con rebalanceo |
// La forma del lab de Raft (concepto): implementas la interfaz
type Raft interface {
RequestVote(args *RequestVoteArgs, reply *RequestVoteReply) error
AppendEntries(args *AppendEntriesArgs, reply *AppendEntriesReply) error
Start(command interface{}) (index, term, isleader)
}
💡 Escribir Raft a mano te hace entender cada regla (y por qué la mayoría no es opcional). Es el proyecto que más «hace click» en distribuidos.
Cheatsheet de conceptos
| Término | En una frase |
|---|---|
| Término (Raft) | Número de elección; votos y log valen para un término |
| Quórum / mayoría | n/2+1 nodos; la base del consenso |
| ISR | Réplicas de Kafka que están al día |
| Offset | Posición del mensaje en la partición |
| Consumer group | Consumidores que reparten las particiones |
| Consistent hashing | Hash en anillo; añadir nodos mueve poco |
| Vector clock | Reloj que detecta escrituras concurrentes |
| Gossip | Nodos que se pasan estado por rumor |
| Hinted handoff | Un nodo guarda escrituras de otro caído |
| TrueTime | Reloj global de Spanner (GPS + relojes atómicos) |
| LSM-tree | Estructura de escritura secuencial (Cassandra, RocksDB) |
| Tunable consistency | Consistencia configurable por query |
| At-least-once | Puede duplicar, nunca perder |
| Replay | Reproducir eventos desde el inicio |
Práctica propuesta
- Lab 2 de MIT 6.824 (Raft en Go): implementa elección y replicación de log. Pásale los tests.
- Levanta Kafka (con Docker) y un consumer group con 2 instancias. Observa el reparto de particiones y el rebalanceo al matar una.
- Lee el paper de Dynamo y dibuja el anillo de consistent hashing con 5 nodos. Añade 2 nodos: ¿qué claves se mueven?
- Monta etcd de 3 nodos, mata el líder y observa la elección. Mata 2 nodos: verifica que pierde el quórum.
- Prueba Cassandra CONSISTENCY ONE vs QUORUM con dos nodos y observa la diferencia de lectura.
- Compara Redis vs RabbitMQ vs Kafka resolviendo el mismo problema (una tarea que se procesa una vez). Anota cuándo elegirías cada uno.
Para profundizar
- MIT 6.5840 / 6.824: el curso graduado de MIT con los labs de Raft.
- Distributed Systems for Fun and Profit: el mapa mental antes de los papers.
- raft.github.io: el paper con animaciones interactivas.
- Dynamo paper, Spanner paper, GFS, Bigtable: las fuentes primarias.
- Apache Kafka docs: el sistema real en profundidad.
- The Log (Kreps): el log como backbone.
- Fundamentos: distribuidos: la base conceptual (CAP, replicación, particionado).
- Patrones avanzados: sagas, outbox, CQRS, event sourcing.
- Sigue con 💼 Entrevistas o 👑 Nivel Élite.