HDFS guarda los eventos masivos y MinIO sirve las fotos, pero ninguno de los dos puede responder en pocos milisegundos a "los últimos pedidos de Ana" ni descontar dos unidades de queso-curado sin vender la misma pieza dos veces. Para eso están las bases de datos, y cuando una sola máquina ya no basta, las bases de datos distribuidas. Esta lección reúne todo lo que el Módulo 3 y las lecciones anteriores han preparado: el particionado y el hashing consistente de 04-01, la replicación y los quórums de 03-04, los modelos de consistencia y la tabla CAP de 03-02, y las transacciones de 03-05. Veremos los tres caminos que existen para distribuir una base de datos (relacional escalada, NoSQL de estilo Dynamo y documental), los compararemos, y tomaremos la decisión más importante del módulo para Kilómetro Cero: pedidos migra a Cassandra e inventario se queda en PostgreSQL. Después lo haremos real: Cassandra en docker-compose.yml, el keyspace km0_pedidos, las tablas pedidos_por_cliente y pedidos_por_id, y un repositorio Python que elige el nivel de consistencia en cada operación.
Contenido
- Qué es una base de datos distribuida y sus transparencias
- Fragmentación horizontal y vertical
- Camino (a): la base de datos relacional escalada y NewSQL
- Camino (b): NoSQL de estilo Dynamo con Cassandra
- Modelado orientado a consultas en Cassandra
- Camino (c): bases de datos documentales con MongoDB
- Tabla comparativa
- La decisión de Kilómetro Cero
- Práctica: Cassandra en
docker-compose.ymlyrepositorio_cassandra.py - Errores Comunes y Consejos
- Ejercicios
- Conclusión
- Qué es una base de datos distribuida y sus transparencias
Una base de datos distribuida es una colección de datos lógicamente relacionados, físicamente repartidos entre varios nodos conectados por red, que se presenta a las aplicaciones como una sola base de datos. La definición tiene dos mitades: el reparto físico (que ya sabemos hacer con particiones y réplicas) y la ilusión de unidad, que es de nuevo la lista de transparencias de 01-01 aplicada a los datos:
| Transparencia | Qué significa | Cuánto la ofrece cada camino |
|---|---|---|
| De fragmentación | La aplicación consulta la tabla pedidos, no la partición 7 |
Total en Cassandra y Spanner; parcial en sharding por aplicación (la aplicación elige el shard) |
| De replicación | La aplicación no sabe cuántas copias hay ni cuál lee | Total, aunque el nivel de consistencia lo puede elegir la aplicación |
| De ubicación | Ninguna consulta menciona un nodo | Total con un router o un driver inteligente (04-01, apartado 9) |
| De fallo | La caída de un nodo no interrumpe el servicio | Depende del factor de replicación y del nivel de consistencia |
| De transacción | Una operación sobre varios nodos es atómica | Total en NewSQL (con coste); limitada a una partición en Cassandra y MongoDB (o con transacciones multi-documento más lentas) |
La última fila es la que separa los caminos. Mantener la transparencia de transacción entre nodos exige 2PC o consenso (03-03, 03-05) en cada escritura que cruce particiones, y cada camino decide cuánto de eso se permite.
- Fragmentación horizontal y vertical
En el vocabulario clásico de las bases de datos distribuidas, particionar una tabla se llama fragmentar:
- Fragmentación horizontal: cada fragmento tiene un subconjunto de las filas, con todas las columnas. Es exactamente el particionado de 04-01 (por rango, hash o compuesto) aplicado a una tabla, y en la jerga NoSQL se llama sharding. Los pedidos de Ana en un nodo, los de Marc en otro.
- Fragmentación vertical: cada fragmento tiene un subconjunto de las columnas (con la clave primaria repetida). Las columnas de facturación de un pedido en un nodo, las de reparto en otro. Es poco frecuente entre nodos de una misma base de datos, pero es precisamente lo que hizo Kilómetro Cero al dividir el monolito:
pedidos,inventario,pagosyrepartoson fragmentos verticales del antiguo esquema, cada uno dueño de sus columnas, conpedido_idcomo clave compartida y sin joins entre servicios (01-06). - Fragmentación mixta: vertical entre servicios, horizontal dentro de cada uno. Es la situación real de la plataforma:
km0_pedidoses un fragmento vertical que a su vez se particiona horizontalmente porcliente_id.
Tres reglas de una buena fragmentación (Özsu y Valduriez): completitud (toda fila o columna está en algún fragmento), reconstrucción (la tabla original se puede recomponer con uniones) y disyunción (una fila está en un solo fragmento, salvo la clave en la vertical). La replicación relaja la tercera a propósito.
- Camino (a): la base de datos relacional escalada y NewSQL
El primer camino conserva SQL, el modelo relacional y las transacciones ACID, y añade distribución por capas:
Réplicas de lectura. Lo montamos en 03-04: pedidos-db-primario recibe todas las escrituras y pedidos-db-replica (y las que se añadan) sirven lecturas con replicación en streaming. Escala las lecturas, no las escrituras ni el tamaño, y trae consigo la consistencia eventual en las réplicas (leer tu propia escritura exige ir al primario o esperar el LSN).
Sharding por aplicación. Cuando las escrituras o el volumen desbordan un primario, la aplicación reparte las filas entre varias instancias PostgreSQL independientes con la lógica de 04-01: shard = anillo.nodo_para(cliente_id), una conexión por shard, y la aplicación sabe en qué shard está cada dato. Es lo que hicieron durante años Instagram, Notion o Shopify. Funciona, pero la transparencia de fragmentación desaparece: los joins entre shards se hacen en la aplicación, las transacciones entre shards no existen (o son sagas, 03-05), reequilibrar shards es un proyecto en sí mismo, y cada nuevo caso de uso tiene que respetar la clave de partición.
Citus. Extensión de PostgreSQL que automatiza ese sharding: un nodo coordinador recibe las consultas SQL normales, y las tablas "distribuidas" (SELECT create_distributed_table('pedidos', 'cliente_id')) se reparten por hash en nodos trabajadores. El coordinador reescribe cada consulta en subconsultas por shard y combina resultados; las transacciones que tocan un solo shard son locales, y las que tocan varios usan 2PC entre trabajadores (03-05, con sus latencias). Es el punto medio: SQL y PostgreSQL con reparto automático, ideal para cargas multiinquilino donde casi todas las consultas incluyen la clave de distribución.
NewSQL. Bases de datos diseñadas desde cero para ser distribuidas sin perder SQL ni ACID:
| Sistema | Idea central | Consistencia | Coste |
|---|---|---|---|
| Google Spanner | Particiones replicadas con Paxos; transacciones globales con 2PC sobre Paxos; TrueTime (relojes atómicos y GPS con incertidumbre acotada) para ordenar commits globalmente (retoma 01-05) | Serializabilidad externa (linealizable y serializable) | Latencia de commit = quórum entre regiones (decenas de ms); hardware especial |
| CockroachDB | Rangos de claves replicados con Raft (03-03); transacciones distribuidas con 2PC optimizado; relojes híbridos (HLC) sin hardware especial | Serializable | Escrituras con latencia de consenso; las lecturas pueden requerir esperas por incertidumbre de reloj |
| YugabyteDB | Tablets replicados con Raft; capa SQL compatible con PostgreSQL (reutiliza su parser) | Snapshot isolation o serializable | Similar a CockroachDB |
NewSQL es CP con transacciones (03-02: elige consistencia en la partición y consistencia sobre latencia en operación normal). Lo que se paga es latencia por escritura (un consenso por rango, más 2PC si hay varios) y complejidad operativa. Es la respuesta correcta cuando se necesitan transacciones entre particiones y SQL completo a escala; no es necesaria para la mayoría de servicios de Kilómetro Cero, que ya han renunciado a las transacciones entre dominios con las sagas.
- Camino (b): NoSQL de estilo Dynamo con Cassandra
En 2007 Amazon publicó el diseño de Dynamo, su almacén clave-valor para el carrito de la compra, que optimizaba por disponibilidad y latencia de escritura: sin líder, hashing consistente, quórums configurables, vectores de versión y reparación en lectura. Cassandra (Facebook, 2008; hoy Apache) combinó la arquitectura de Dynamo con el modelo de datos de familias de columnas de Bigtable. Riak, Voldemort y ScyllaDB pertenecen a la misma familia. Todo lo que sigue es la puesta en práctica de 04-01 y 03-04:
Anillo con hashing consistente. Cada nodo posee rangos de tokens del anillo de Murmur3 (−2⁶³ a 2⁶³−1); con num_tokens: 16 (Cassandra 4) cada nodo tiene 16 vnodes. La clave de partición se hashea y el token resultante determina el nodo coordinador de esa partición y sus réplicas: las siguientes N−1 posiciones distintas del anillo (saltando vnodes del mismo nodo y, con NetworkTopologyStrategy, repartiendo entre racks y centros de datos).
flowchart TB
subgraph anillo["Anillo km0_pedidos · RF=3 · 6 nodos"]
direction LR
n1(("nodo1<br/>rack1"))
n2(("nodo2<br/>rack2"))
n3(("nodo3<br/>rack1"))
n4(("nodo4<br/>rack2"))
n5(("nodo5<br/>rack1"))
n6(("nodo6<br/>rack2"))
n1 --> n2 --> n3 --> n4 --> n5 --> n6 --> n1
end
k["partición cliente_id = 'ana'<br/>token = 0x3F… → cae entre nodo2 y nodo3"] -.-> n3
n3 -. "réplica 1" .- r1[" "]
n4 -. "réplica 2 (otro rack)" .- r2[" "]
n5 -. "réplica 3" .- r3[" "]
style r1 fill:none,stroke:none
style r2 fill:none,stroke:none
style r3 fill:none,stroke:none
Sin líder. Cualquier nodo acepta cualquier petición (opción "cualquier nodo" de 04-01): el nodo que la recibe actúa como coordinador, reenvía la escritura a las N réplicas de la partición y espera las confirmaciones que exija el nivel de consistencia. No hay elección de líder ni failover: si una réplica está caída, las demás siguen aceptando escrituras y el coordinador guarda un hinted handoff (03-04) para entregárselo cuando vuelva.
Factor de replicación y estrategia. Se fijan por keyspace: SimpleStrategy (réplicas en los siguientes nodos del anillo, solo para pruebas) o NetworkTopologyStrategy con un factor por centro de datos ({'dc-bcn': 3, 'dc-vlc': 3}), que además reparte las réplicas entre racks. Cambiar el factor exige nodetool repair para que las nuevas réplicas se llenen.
Niveles de consistencia por operación. Es la aportación más elegante de Cassandra: cada lectura y cada escritura elige cuántas réplicas deben responder. Con N = 3:
| Nivel | Escritura: confirma cuando… | Lectura: responde con… | W+R>N con RF=3 |
|---|---|---|---|
ONE |
1 réplica ha escrito (en commit log y memtable) | La primera réplica que responde | ONE + ONE = 2 ≤ 3: puede leer datos viejos |
QUORUM |
⌊N/2⌋+1 = 2 réplicas | 2 réplicas, se devuelve la más reciente (por marca de tiempo) | QUORUM + QUORUM = 4 > 3: lectura fuerte |
ALL |
Las 3 | Las 3 | Máxima consistencia, disponibilidad mínima: un nodo caído bloquea |
LOCAL_QUORUM |
Quórum dentro del centro de datos local | Ídem | Fuerte dentro del DC, sin esperar a la otra región |
ANY |
Cualquier nodo, incluso solo un hint | — | Máxima disponibilidad de escritura, sin garantía de lectura |
SERIAL / LOCAL_SERIAL |
Transacciones ligeras (IF NOT EXISTS) con Paxos |
Lee el estado tras el último Paxos | Linealizable por partición, mucho más lento |
Es exactamente el W + R > N de 03-04 con quorum_wr.py, pero decidido en cada operación: la posición de furgoneta-3 se escribe con ONE (rápido, AP), la confirmación de un pedido con QUORUM (consistente, CP), y la lectura que muestra a Ana su pedido recién confirmado con QUORUM para que vea lo que acaba de escribir. Las lecturas hacen read repair cuando las réplicas discrepan, y nodetool repair (anti-entropía con árboles de Merkle) reconcilia en segundo plano.
Escritura en disco. Una escritura va al commit log (secuencial, duradero) y a la memtable en memoria; cuando la memtable se llena, se vuelca a un SSTable inmutable en disco. Las SSTables se compactan periódicamente fusionando versiones. Este diseño (LSM-tree) hace que las escrituras sean secuenciales y baratísimas, que es la razón por la que Cassandra es la elección clásica para cargas de escritura intensiva; las lecturas pueden tener que consultar varias SSTables (mitigado con filtros de Bloom y cachés).
Tombstones. Como las SSTables son inmutables, borrar no borra: escribe una lápida (tombstone) con marca de tiempo que oculta el valor. La lápida se elimina en la compactación tras gc_grace_seconds (10 días por defecto), tiempo que debe superar la reparación de cualquier réplica que estuviera caída; si no, la réplica caída "resucitaría" el dato borrado al reconciliar. Muchas lápidas en una partición (borrados masivos, colas implementadas en Cassandra) degradan gravemente las lecturas: es el antipatrón más famoso de Cassandra.
- Modelado orientado a consultas en Cassandra
En SQL se modela el dominio (normalizado) y después se escriben las consultas que hagan falta, con joins e índices. En Cassandra no hay joins ni consultas ad hoc eficientes, así que se invierte el orden: primero las consultas, luego una tabla por consulta, desnormalizando lo que haga falta. La clave primaria tiene dos partes:
CREATE TABLE pedidos_por_cliente (
cliente_id text,
fecha timestamp,
pedido_id text,
estado text,
total decimal,
lineas list<frozen<linea>>,
PRIMARY KEY ((cliente_id), fecha, pedido_id)
) WITH CLUSTERING ORDER BY (fecha DESC, pedido_id ASC);- La clave de partición
(cliente_id)(los paréntesis interiores) decide el token y por tanto el nodo: todos los pedidos de Ana están juntos, en el mismo nodo (y sus réplicas). Es la decisión de 04-01, apartado 10, hecha esquema. - Las clustering columns
fecha, pedido_idordenan las filas dentro de la partición, en disco.CLUSTERING ORDER BY (fecha DESC)guarda los más recientes primero, así que "los últimos 20 pedidos de Ana" es leer las primeras 20 filas de la partición: una sola búsqueda, un solo nodo. - Las consultas eficientes son las que fijan la clave de partición completa y, opcionalmente, un rango sobre las clustering columns en orden:
WHERE cliente_id = 'ana',WHERE cliente_id = 'ana' AND fecha > '2026-09-01'. Una consulta sin clave de partición (WHERE estado = 'PENDIENTE') exigeALLOW FILTERINGy recorre todo el clúster: es el scatter/gather de 04-01 en su peor forma, y se prohíbe en producción. - Para "pedido por id" (la confirmación, el enlace del correo) se necesita otra tabla,
pedidos_por_id, conPRIMARY KEY ((pedido_id)), y la aplicación escribe en ambas en unBATCHregistrado (atómico entre las dos tablas, aunque no aislado). Es el precio de la desnormalización, y es barato porque las escrituras lo son. - Las particiones deben mantenerse acotadas (recomendación: < 100 MB y < 100 000 filas): un cliente restaurante con decenas de miles de pedidos rompería la regla, y por eso 04-01 propuso el cubo mensual
(cliente_id, año_mes)como clave de partición compuesta. Lo aplicaremos como ejercicio.
- Camino (c): bases de datos documentales con MongoDB
Las bases de datos documentales guardan documentos JSON/BSON con esquema flexible y consultas ricas sobre cualquier campo, con índices secundarios. MongoDB las distribuye en dos niveles:
- Réplica set: un primario y varios secundarios con replicación asíncrona por oplog y elección automática de primario (un protocolo derivado de Raft, 03-03). La aplicación elige el write concern (
w: 1,w: "majority") y el read preference (primary,secondaryPreferred), que son de nuevoWyRcon otros nombres. Es líder-seguidor con failover (03-04). - Sharding: las colecciones se reparten en chunks por rango o hash de una shard key, cada chunk vive en un réplica set, los routers
mongosenrutan las consultas (opción "capa de enrutamiento" de 04-01), y los config servers (un réplica set) guardan el mapa de chunks, que un balanceador mueve para reequilibrar. Las consultas con la shard key van a un shard; sin ella, a todos. - Transacciones: ACID en un documento siempre (un pedido con sus líneas embebidas se escribe atómicamente); desde la versión 4.0/4.2, transacciones multi-documento y multi-shard con 2PC interno, más lentas y con límites.
Encaja cuando el dominio se expresa bien como documentos autocontenidos con consultas variadas (catálogos, perfiles, contenido). Lo tratamos brevemente porque, para Kilómetro Cero, el catálogo (documento por producto con variantes y fotos) sería un candidato natural, pero el equipo decidió mantenerlo en PostgreSQL con columnas jsonb (que cubren el 90 % del caso documental) más la caché de Redis (04-05): una base de datos menos que operar.
- Tabla comparativa
| PostgreSQL + réplicas | Citus / sharding | NewSQL (CockroachDB, Spanner) | Cassandra | MongoDB | |
|---|---|---|---|---|---|
| Modelo de datos | Relacional | Relacional | Relacional | Familia de columnas (tablas anchas, particiones) | Documentos |
| Consistencia (03-02) | CP en el primario; réplicas eventuales | CP por shard | CP, serializable global | Ajustable por operación (ONE…ALL); AP por defecto, CP con QUORUM | CP en el primario; ajustable con write/read concern |
| Escalado de escrituras | No (un primario) | Sí, por shard | Sí, por rango | Sí, lineal (sin líder) | Sí, por shard |
| Escalado de lecturas | Sí (réplicas) | Sí | Sí | Sí | Sí (secundarios) |
| Consultas | SQL completo, joins | SQL; joins eficientes solo colocalizados | SQL completo | Solo por clave de partición; sin joins; una tabla por consulta | Ricas sobre cualquier campo; agregaciones |
| Transacciones | ACID completo | ACID en un shard; 2PC entre shards | ACID distribuido | Atomicidad por partición; batch entre tablas; LWT con Paxos | ACID por documento; multi-documento con coste |
| Índices secundarios | Sí | Sí | Sí | Locales (scatter/gather) o vistas materializadas | Sí, incluidos globales por shard |
| Latencia de escritura | Baja (1 nodo) | Baja/media | Media (consenso) | Muy baja (LSM, sin líder) | Baja |
| Punto fuerte | Todo lo que cabe en una máquina grande; contadores y restricciones | Multiinquilino con SQL | Transacciones globales | Escritura masiva, series temporales, disponibilidad | Flexibilidad de esquema |
| Punto débil | Escalado de escritura | Consultas sin la clave de distribución | Latencia y complejidad | Modelado rígido, tombstones, sin agregaciones | Memoria, sharding difícil de cambiar |
- La decisión de Kilómetro Cero
Con la tabla delante y la tabla de decisiones CAP de 03-02, el equipo decide por servicio:
pedidos migra a Cassandra. Razones:
- Volumen y escritura intensiva: 40 000 pedidos/día en campaña con 6 eventos por pedido, más el historial completo (la ley obliga a conservarlo años) y la tabla
sagas(03-05) con una transición por paso. Es una carga de escritura secuencial y de lectura por clave conocida, el punto fuerte del LSM-tree y del anillo sin líder. - Consultas conocidas y estables: "mis pedidos" (por cliente, ordenados por fecha), "pedido por id", y los índices por mercado y productor materializados desde
pedidos.eventos. Ninguna necesita joins ni agregaciones en línea (esas van al lago de datos y al Módulo 5). - Disponibilidad: crear un pedido debe funcionar aunque un nodo o incluso un centro de datos falle; con
LOCAL_QUORUMy RF=3 por DC, la escritura sigue con un nodo caído. El pedido ya no necesita transacciones entre servicios (sagas) ni entre tablas fuera de un batch, así que la atomicidad por partición es suficiente. - Escala lineal: añadir nodos añade capacidad proporcional, y 04-01 mostró que el anillo con vnodes mueve solo lo imprescindible.
inventario se queda en PostgreSQL (primario + réplica de 03-04, particionado por producto_slug cuando haga falta). Razones:
- Contadores CP: el stock es una restricción (
CHECK (disponible >= 0)) que hay que hacer cumplir en cada reserva, conUPDATE … WHERE disponible >= 2bloqueando la fila. Cassandra no tiene restricciones, sus contadores no son idempotentes ni transaccionales, y las transacciones ligeras (Paxos por operación) serían mucho más lentas que unUPDATEen PostgreSQL. - Volumen modesto: 4 200 productos por 4 mercados son unas 17 000 filas de stock; el problema no es el tamaño sino la corrección y la contención en las claves calientes, que se atacará con la caché y con colas (04-05), no con más nodos.
- Transacciones locales: reservar stock y registrar la reserva idempotente (
mensajes_procesados, 02-05) en la misma transacción, y la tabla outbox (02-05) en la misma base: PostgreSQL lo hace en unCOMMIT.
pagos también se queda en PostgreSQL (CP, bajo volumen, auditoría); catalogo en PostgreSQL con jsonb y Redis delante; reparto (posiciones de furgoneta-3, 2,4 millones al día, AP) es candidato a Cassandra con particiones diarias por repartidor, como planteó el ejercicio 2 de 04-01; analitica en el lago (04-02) y en un almacén analítico que el Módulo 5 elegirá.
- Práctica: Cassandra en
docker-compose.yml y repositorio_cassandra.py
docker-compose.yml y repositorio_cassandra.pyTres nodos Cassandra en un solo centro de datos (dc-bcn), con dos racks simulados, para poder usar NetworkTopologyStrategy y ver el reparto de réplicas:
# km0/docker-compose.yml (fragmento)
x-cassandra-env: &cassandra-env
CASSANDRA_CLUSTER_NAME: km0
CASSANDRA_SEEDS: cassandra-1
CASSANDRA_DC: dc-bcn
CASSANDRA_ENDPOINT_SNITCH: GossipingPropertyFileSnitch # necesario para NetworkTopologyStrategy
CASSANDRA_NUM_TOKENS: "16"
MAX_HEAP_SIZE: 1G
HEAP_NEWSIZE: 200M
services:
cassandra-1:
image: cassandra:5.0
hostname: cassandra-1
environment:
<<: *cassandra-env
CASSANDRA_RACK: rack1
ports: ["9042:9042"]
volumes: [cassandra-1-data:/var/lib/cassandra]
healthcheck:
test: ["CMD-SHELL", "nodetool status | grep -q '^UN'"]
interval: 15s
retries: 20
cassandra-2:
image: cassandra:5.0
hostname: cassandra-2
environment:
<<: *cassandra-env
CASSANDRA_RACK: rack2
volumes: [cassandra-2-data:/var/lib/cassandra]
depends_on:
cassandra-1: { condition: service_healthy }
cassandra-3:
image: cassandra:5.0
hostname: cassandra-3
environment:
<<: *cassandra-env
CASSANDRA_RACK: rack1
volumes: [cassandra-3-data:/var/lib/cassandra]
depends_on:
cassandra-2: { condition: service_started }
volumes:
cassandra-1-data:
cassandra-2-data:
cassandra-3-data:Los nodos deben arrancar de uno en uno (por eso las dependencias): dos nodos uniéndose al anillo a la vez con el mismo seed es una de las causas clásicas de un clúster mal formado. Arrancar y comprobar el anillo:
docker compose up -d cassandra-1 cassandra-2 cassandra-3 # tarda 2-3 minutos
docker compose exec cassandra-1 nodetool statusDatacenter: dc-bcn ================== Status=Up/Down |/ State=Normal/Leaving/Joining/Moving -- Address Load Tokens Owns (effective) Host ID Rack UN 172.21.0.3 104.51 KiB 16 100.0% 5f2c… rack1 UN 172.21.0.4 109.87 KiB 16 100.0% 8a1d… rack2 UN 172.21.0.5 98.22 KiB 16 100.0% c73e… rack1
UN = Up, Normal. Owns (effective) es 100 % en cada nodo porque, con RF=3 y 3 nodos, cada nodo tiene una réplica de todo; con 6 nodos veríamos ~50 %. nodetool ring muestra los 48 tokens (16 por nodo) del anillo, y nodetool getendpoints km0_pedidos pedidos_por_cliente ana dirá qué tres nodos guardan la partición de Ana.
El esquema, con cqlsh:
-- docker compose exec cassandra-1 cqlsh
CREATE KEYSPACE IF NOT EXISTS km0_pedidos
WITH replication = {'class': 'NetworkTopologyStrategy', 'dc-bcn': 3};
USE km0_pedidos;
CREATE TYPE IF NOT EXISTS linea (
producto text,
productor text,
cantidad int,
precio decimal
);
CREATE TABLE IF NOT EXISTS pedidos_por_cliente (
cliente_id text,
fecha timestamp,
pedido_id text,
mercado text,
estado text,
total decimal,
lineas list<frozen<linea>>,
PRIMARY KEY ((cliente_id), fecha, pedido_id)
) WITH CLUSTERING ORDER BY (fecha DESC, pedido_id ASC)
AND comment = 'Consulta: mis pedidos, más recientes primero';
CREATE TABLE IF NOT EXISTS pedidos_por_id (
pedido_id text PRIMARY KEY,
cliente_id text,
fecha timestamp,
mercado text,
estado text,
total decimal,
lineas list<frozen<linea>>
) WITH comment = 'Consulta: pedido por id (confirmación, correo, saga)';El tipo linea es un UDT (user-defined type) y frozen indica que la lista se guarda como un solo valor serializado (no se pueden modificar elementos sueltos, pero se lee y escribe de una vez, que es lo que queremos). pedidos_por_cliente y pedidos_por_id contienen los mismos datos con distinta clave: es la desnormalización del apartado 5.
El repositorio Python con cassandra-driver (el driver oficial de DataStax; pip install cassandra-driver):
# km0/servicios/pedidos/repositorio_cassandra.py
"""Acceso a km0_pedidos en Cassandra con nivel de consistencia por operación."""
from datetime import datetime, timedelta, timezone
from decimal import Decimal
from cassandra.cluster import Cluster, ExecutionProfile, EXEC_PROFILE_DEFAULT
from cassandra.policies import DCAwareRoundRobinPolicy, TokenAwarePolicy
from cassandra.query import BatchStatement, BatchType, ConsistencyLevel
class RepositorioPedidos:
def __init__(self, contactos=("localhost",), dc_local="dc-bcn"):
perfil = ExecutionProfile(
load_balancing_policy=TokenAwarePolicy(DCAwareRoundRobinPolicy(local_dc=dc_local)),
consistency_level=ConsistencyLevel.LOCAL_QUORUM, # valor por defecto: fuerte dentro del DC
request_timeout=5.0,
)
self.cluster = Cluster(contact_points=list(contactos),
execution_profiles={EXEC_PROFILE_DEFAULT: perfil})
self.session = self.cluster.connect("km0_pedidos")
self.session.cluster.register_user_type("km0_pedidos", "linea", Linea)
# Sentencias preparadas: se analizan una vez en el servidor y se reutilizan (más rápido y seguro)
self._ins_cliente = self.session.prepare(
"INSERT INTO pedidos_por_cliente (cliente_id, fecha, pedido_id, mercado, estado, total, lineas) "
"VALUES (?, ?, ?, ?, ?, ?, ?)")
self._ins_id = self.session.prepare(
"INSERT INTO pedidos_por_id (pedido_id, cliente_id, fecha, mercado, estado, total, lineas) "
"VALUES (?, ?, ?, ?, ?, ?, ?)")
self._sel_cliente = self.session.prepare(
"SELECT pedido_id, fecha, mercado, estado, total FROM pedidos_por_cliente "
"WHERE cliente_id = ? LIMIT ?")
self._sel_id = self.session.prepare("SELECT * FROM pedidos_por_id WHERE pedido_id = ?")
self._upd_estado_cliente = self.session.prepare(
"UPDATE pedidos_por_cliente SET estado = ? WHERE cliente_id = ? AND fecha = ? AND pedido_id = ?")
self._upd_estado_id = self.session.prepare(
"UPDATE pedidos_por_id SET estado = ? WHERE pedido_id = ?")
# --- escrituras -------------------------------------------------------------------------
def crear_pedido(self, pedido_id: str, cliente_id: str, mercado: str, lineas: list, total: Decimal,
fecha: datetime | None = None) -> None:
"""Escribe en las dos tablas dentro de un batch registrado: o entran las dos filas o ninguna."""
fecha = fecha or datetime.now(timezone.utc)
lote = BatchStatement(batch_type=BatchType.LOGGED, consistency_level=ConsistencyLevel.LOCAL_QUORUM)
lote.add(self._ins_cliente, (cliente_id, fecha, pedido_id, mercado, "CREADO", total, lineas))
lote.add(self._ins_id, (pedido_id, cliente_id, fecha, mercado, "CREADO", total, lineas))
self.session.execute(lote)
def cambiar_estado(self, pedido_id: str, nuevo_estado: str) -> None:
"""Necesita la clave completa de pedidos_por_cliente: la obtiene de pedidos_por_id."""
fila = self.session.execute(self._sel_id, (pedido_id,)).one()
if fila is None:
raise KeyError(pedido_id)
lote = BatchStatement(batch_type=BatchType.LOGGED, consistency_level=ConsistencyLevel.LOCAL_QUORUM)
lote.add(self._upd_estado_cliente, (nuevo_estado, fila.cliente_id, fila.fecha, pedido_id))
lote.add(self._upd_estado_id, (nuevo_estado, pedido_id))
self.session.execute(lote)
# --- lecturas ---------------------------------------------------------------------------
def ultimos_pedidos(self, cliente_id: str, limite: int = 20, fuerte: bool = False) -> list:
"""Listado de 'mis pedidos'. Por defecto ONE (rápido); `fuerte=True` usa LOCAL_QUORUM
justo después de crear un pedido, para que el cliente vea su propia escritura."""
consulta = self._sel_cliente.bind((cliente_id, limite))
consulta.consistency_level = ConsistencyLevel.LOCAL_QUORUM if fuerte else ConsistencyLevel.ONE
return list(self.session.execute(consulta))
def pedido(self, pedido_id: str, nivel=ConsistencyLevel.LOCAL_QUORUM):
consulta = self._sel_id.bind((pedido_id,))
consulta.consistency_level = nivel
return self.session.execute(consulta).one()
def cerrar(self):
self.cluster.shutdown()
class Linea:
def __init__(self, producto, productor, cantidad, precio):
self.producto, self.productor, self.cantidad, self.precio = producto, productor, cantidad, precio
if __name__ == "__main__":
repo = RepositorioPedidos()
repo.crear_pedido("P-2026-000123", "ana", "Girona",
[Linea("queso-curado", "queseria-montblanc", 2, Decimal("14.50")),
Linea("tomate-rosa", "huerta-la-vega", 3, Decimal("3.20"))], Decimal("38.60"))
repo.crear_pedido("P-2026-000124", "ana", "Girona",
[Linea("vino-crianza", "bodega-roble-alto", 1, Decimal("15.90"))], Decimal("15.90"))
repo.crear_pedido("P-2026-000125", "marc", "Valencia",
[Linea("queso-fresco", "queseria-montblanc", 1, Decimal("6.10"))], Decimal("6.10"))
print("Pedidos de Ana (lectura fuerte tras escribir):")
for p in repo.ultimos_pedidos("ana", fuerte=True):
print(" ", p.pedido_id, p.fecha, p.mercado, p.estado, p.total)
repo.cambiar_estado("P-2026-000123", "PAGADO")
print("P-2026-000123:", repo.pedido("P-2026-000123").estado)
repo.cerrar()Lo que conviene entender del código:
TokenAwarePolicy(DCAwareRoundRobinPolicy): el driver descarga el mapa de tokens del anillo y envía cada petición directamente a una réplica de la partición (opción "cliente informado" de 04-01), evitando el salto extra del coordinador; y prefiere los nodos del centro de datos local.ExecutionProfilefijaLOCAL_QUORUMpor defecto, y cada sentencia puede sobrescribirlo conconsistency_level. Este es el punto central de la práctica: la consistencia es una propiedad de la operación, no de la base de datos.ultimos_pedidosusaONEen la ruta normal (la lista de pedidos tolera unos milisegundos de retraso) yLOCAL_QUORUMcuando la web la llama justo después decrear_pedido, que también fueLOCAL_QUORUM: 2 + 2 > 3, así que la lectura ve la escritura. Es el modelo read-your-writes de 03-01 conseguido con quórums.BatchStatement(LOGGED): Cassandra escribe primero el lote en un batchlog replicado y garantiza que las dos inserciones acaben aplicándose (atomicidad eventual, sin aislamiento: un lector puede ver una tabla actualizada y la otra no durante milisegundos). Los batches sirven para mantener tablas desnormalizadas coherentes, no para "ir más rápido": un batch con cientos de particiones distintas es un antipatrón.- Sentencias preparadas con
?: se envían al servidor una vez, se reutilizan con parámetros, y el driver sabe qué parámetro es la clave de partición para el enrutamiento por token. cambiar_estadomuestra el coste de la desnormalización: para actualizarpedidos_por_clientehace falta la clave completa (cliente_id,fecha,pedido_id), que se obtiene depedidos_por_id. En la saga de 03-05, el orquestador ya conoce esos datos y se ahorra la lectura.
Ejecuta el script, y luego prueba los niveles con el clúster degradado:
python -m servicios.pedidos.repositorio_cassandra
docker compose stop cassandra-3
docker compose exec cassandra-1 nodetool status # cassandra-3 aparece como DN (Down, Normal)
python - <<'EOF'
from cassandra.query import ConsistencyLevel
from servicios.pedidos.repositorio_cassandra import RepositorioPedidos
repo = RepositorioPedidos()
print("QUORUM con 2 de 3:", repo.pedido("P-2026-000123", ConsistencyLevel.QUORUM).estado) # funciona
try:
repo.pedido("P-2026-000123", ConsistencyLevel.ALL) # falla
except Exception as e:
print("ALL con 2 de 3:", type(e).__name__) # Unavailable: no hay 3 réplicas vivas
EOF
docker compose start cassandra-3Con un nodo caído, QUORUM sigue leyendo y escribiendo (2 réplicas vivas de 3) y ALL responde Unavailable inmediatamente: la disponibilidad y la consistencia elegidas operación a operación, como prometía 03-02. Al volver cassandra-3, los hints acumulados por cassandra-1 y cassandra-2 se le entregan y nodetool repair km0_pedidos reconcilia cualquier diferencia restante.
Errores Comunes y Consejos
- Modelar Cassandra como SQL. Una tabla
pedidosnormalizada con índices secundarios para cada consulta termina enALLOW FILTERINGy scatter/gather. Lista las consultas, una tabla por consulta, desnormaliza con batches. - Particiones sin límite. Un cliente, un repartidor o un producto con millones de filas en la misma partición degrada compactaciones y lecturas. Añade un cubo temporal a la clave de partición.
- Usar Cassandra para colas o para borrados masivos. Las lápidas se acumulan y las lecturas mueren. Si necesitas una cola, ya tienes Kafka (02-04).
SimpleStrategyen producción. No entiende de racks ni de centros de datos: tres réplicas en el mismo rack. SiempreNetworkTopologyStrategy, incluso con un solo DC.gc_grace_secondsmenor que el tiempo máximo de reparación. Un nodo que vuelve después resucita datos borrados. Repara cada nodo dentro de esa ventana (nodetool repairprogramado).- Elegir NewSQL o Cassandra "porque escala" con 20 GB de datos. Un PostgreSQL bien indexado en una máquina grande sirve decenas de miles de transacciones por segundo. Distribuye cuando el volumen, las escrituras o la disponibilidad lo exijan, y pagarás el modelado rígido o la latencia de consenso solo entonces.
- Contadores de negocio en Cassandra. Los
counterno son idempotentes (un reintento duplica el incremento) y no admiten condiciones. El stock, los saldos y todo lo que tenga una restricción van a PostgreSQL o a NewSQL. - Ignorar el nivel de consistencia por defecto del driver. En
cassandra-driveresLOCAL_ONE. Fíjalo explícitamente en elExecutionProfiley decide por operación.
Ejercicios
Ejercicio 1. Aplica la recomendación de particiones acotadas: redefine pedidos_por_cliente con clave de partición compuesta ((cliente_id, anyo_mes), fecha, pedido_id) donde anyo_mes es un texto como '2026-09'. Escribe el CREATE TABLE, adapta crear_pedido para calcular anyo_mes a partir de fecha, y reescribe ultimos_pedidos(cliente_id, limite) para que devuelva los últimos 20 pedidos aunque estén repartidos en varios meses (pista: recorrer meses hacia atrás hasta completar el límite o llegar a un tope). ¿Qué ocurre con un cliente que no ha comprado en los últimos 12 meses?
Ejercicio 2. Para cada operación de Kilómetro Cero, elige el nivel de consistencia de Cassandra (con RF=3 en dc-bcn y RF=3 en dc-vlc) y justifícalo con W + R > N y con la tabla CAP de 03-02: (a) escribir una posición de furgoneta-3; (b) leer la última posición para la web del cliente; (c) confirmar el pedido P-2026-000126 al recibir pago.confirmado; (d) leer ese pedido desde la página de confirmación 200 ms después; (e) el informe nocturno de analitica que lee todos los pedidos del día; (f) reservar el nombre de usuario de un nuevo productor, que debe ser único.
Ejercicio 3. Un compañero propone migrar también inventario a Cassandra "para tener una sola base de datos", implementando el stock como UPDATE stock SET disponible = disponible - 2 WHERE producto = 'queso-curado' con un tipo counter, y comprobando después con un SELECT que no ha quedado negativo. Explica con un escenario concreto (dos reservas concurrentes de Ana y Marc sobre las 3 últimas piezas) por qué falla, qué alternativa ofrece Cassandra (transacciones ligeras con IF) y por qué, aun así, la decisión del apartado 8 se mantiene.
Soluciones
Solución 1:
CREATE TABLE pedidos_por_cliente_mes (
cliente_id text, anyo_mes text, fecha timestamp, pedido_id text,
mercado text, estado text, total decimal, lineas list<frozen<linea>>,
PRIMARY KEY ((cliente_id, anyo_mes), fecha, pedido_id)
) WITH CLUSTERING ORDER BY (fecha DESC, pedido_id ASC);def _anyo_mes(fecha: datetime) -> str:
return fecha.strftime("%Y-%m")
# en crear_pedido: lote.add(self._ins_cliente, (cliente_id, _anyo_mes(fecha), fecha, pedido_id, ...))
# en __init__: self._sel_cliente_mes = session.prepare(
# "SELECT * FROM pedidos_por_cliente_mes WHERE cliente_id = ? AND anyo_mes = ? LIMIT ?")
def ultimos_pedidos(self, cliente_id: str, limite: int = 20, meses_max: int = 12) -> list:
resultado, cursor = [], datetime.now(timezone.utc).replace(day=1)
for _ in range(meses_max):
filas = self.session.execute(self._sel_cliente_mes, (cliente_id, _anyo_mes(cursor), limite - len(resultado)))
resultado.extend(filas)
if len(resultado) >= limite:
break
cursor = (cursor - timedelta(days=1)).replace(day=1) # mes anterior
return resultadoCada iteración es una lectura de una partición distinta (una partición por mes), así que "los últimos 20" cuesta entre 1 y meses_max lecturas, casi siempre 1 o 2 para un cliente activo. Un cliente sin compras en 12 meses recibe una lista vacía tras 12 lecturas rápidas (particiones inexistentes se resuelven con filtros de Bloom sin tocar disco): el tope evita recorrer años hacia atrás, y si el negocio necesita el historial completo se añade una consulta explícita por rango de meses o una tabla resumen meses_con_pedidos_por_cliente.
Solución 2:
| Operación | Nivel | Justificación |
|---|---|---|
(a) Escribir posición de furgoneta-3 |
ONE (o ANY) |
AP en 03-02: 2,4 M escrituras/día; perder una posición es irrelevante; latencia mínima. W=1 |
| (b) Leer última posición | ONE |
Un dato de hace 3 s vale igual; W+R = 2 ≤ 3, aceptamos leer viejo |
(c) Confirmar P-2026-000126 |
LOCAL_QUORUM (W=2 en el DC local) |
CP para el pedido; no esperar al otro DC (latencia entre BCN y VLC); un nodo caído no bloquea |
| (d) Leer la confirmación 200 ms después | LOCAL_QUORUM (R=2) |
W+R = 4 > 3 dentro del DC: read-your-writes garantizado si la lectura va al mismo DC (el driver lo asegura con DCAwareRoundRobinPolicy); si la web pudiera leer desde el otro DC, haría falta EACH_QUORUM en la escritura |
| (e) Informe nocturno | ONE o LOCAL_ONE |
Lectura masiva, sin prisa ni sensibilidad a milisegundos de retraso; y mejor aún desde el lago de datos (04-02) que desde Cassandra |
| (f) Nombre de usuario único | LOCAL_SERIAL con INSERT … IF NOT EXISTS |
Es una decisión que exige linealizabilidad (03-01): solo una transacción ligera (Paxos por partición) garantiza que dos productores no obtengan el mismo nombre; se acepta su latencia porque ocurre una vez por productor |
Solución 3:
Escenario: quedan 3 piezas. Ana reserva 2 y Marc reserva 2 casi a la vez. Con counter, ambos UPDATE se aplican sin condición: el contador pasa a 3 − 2 − 2 = −1. Las comprobaciones posteriores con SELECT ven −1 las dos, y cada una intenta "deshacer" sumando 2: el contador queda en 3, o en 1, o en −1 según el entrelazado, y ninguno sabe si su reserva vale. Peor: si un UPDATE de contador se reintenta por un timeout (03-05, sagas reintentables), se aplica dos veces, porque los contadores no son idempotentes. La alternativa correcta en Cassandra es una transacción ligera: UPDATE stock SET disponible = 1 WHERE producto = 'queso-curado' IF disponible = 3, que ejecuta Paxos entre las réplicas de la partición y devuelve [applied] = false al que llega segundo, que debe releer y reintentar (compare-and-set). Funciona, pero cada reserva cuesta cuatro viajes de red entre réplicas (unas decenas de milisegundos) frente a un UPDATE … WHERE disponible >= 2 en PostgreSQL con bloqueo de fila (menos de un milisegundo), y sin restricciones declarativas ni transacción con la tabla mensajes_procesados y la outbox. Con queso-curado como clave caliente en la Semana del Queso Artesano, la contención en Paxos sería el cuello de botella. Por eso inventario se queda en PostgreSQL: el problema es la corrección de un contador disputado, no el volumen, y ese es el punto fuerte de una base de datos relacional CP.
Conclusión
Una base de datos distribuida reparte físicamente los datos y ofrece las transparencias de fragmentación, replicación, ubicación, fallo y, con distinto alcance, de transacción. La fragmentación horizontal es el particionado de 04-01 sobre filas, y la vertical es lo que Kilómetro Cero hizo al repartir el esquema del monolito entre servicios. Hay tres caminos: la base relacional escalada con réplicas de lectura, sharding por aplicación o Citus, y en su versión más ambiciosa NewSQL (Spanner con TrueTime, CockroachDB y YugabyteDB con Raft) que conserva ACID global a costa de latencia; el NoSQL de estilo Dynamo que Cassandra encarna con un anillo de hashing consistente, sin líder, con factor de replicación por keyspace y con el nivel de consistencia elegido en cada operación (ONE, QUORUM, ALL) como aplicación literal de W + R > N, con un modelado orientado a consultas (una tabla por consulta, clave de partición más clustering columns) y con las lápidas como su peaje; y las bases documentales como MongoDB con réplica sets y sharding. Kilómetro Cero ha decidido que pedidos migre a Cassandra por su volumen, su escritura intensiva y sus consultas estables, y que inventario siga en PostgreSQL porque un contador con restricción es un problema CP de corrección, no de escala. Lo hemos montado con tres nodos en docker-compose.yml, el keyspace km0_pedidos con NetworkTopologyStrategy, las tablas pedidos_por_cliente y pedidos_por_id, y repositorio_cassandra.py escribiendo en batch y leyendo con ONE o LOCAL_QUORUM según lo que cada operación necesita, comprobando con nodetool status y un nodo parado que QUORUM sobrevive y ALL no.
Los datos de Kilómetro Cero tienen ya casi todos su sitio. Pero el catálogo se consulta cientos de veces por cada vez que cambia, y cada consulta que llega a PostgreSQL durante la Semana del Queso Artesano es una consulta que la base de datos podría no tener que responder. La última pieza del módulo está entre la aplicación y las bases de datos: los cachés distribuidos, con Redis, sus patrones, sus problemas clásicos y un Redis Cluster cuyos 16 384 slots son la última reencarnación del particionado con el que empezó el módulo.
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
