🔎 Buscar

🌐 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.

System Design📖 Contenido

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.

📌 Mentalidad

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

  1. Elección de líder: decidir QUIÉN manda (con mayoría de votos).
  2. Replicación de log: el líder replica su log y los seguidores lo aplican.
  3. 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 W nodos, lecturas en R nodos; si R + W > N siempre 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

  1. Lab 2 de MIT 6.824 (Raft en Go): implementa elección y replicación de log. Pásale los tests.
  2. Levanta Kafka (con Docker) y un consumer group con 2 instancias. Observa el reparto de particiones y el rebalanceo al matar una.
  3. Lee el paper de Dynamo y dibuja el anillo de consistent hashing con 5 nodos. Añade 2 nodos: ¿qué claves se mueven?
  4. Monta etcd de 3 nodos, mata el líder y observa la elección. Mata 2 nodos: verifica que pierde el quórum.
  5. Prueba Cassandra CONSISTENCY ONE vs QUORUM con dos nodos y observa la diferencia de lectura.
  6. 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

Estudio · Recursos de todo el mundo (inglés, chino, japonés, español, francés, ruso…) curados y traducidos al español.