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

  1. Qué es una base de datos distribuida y sus transparencias
  2. Fragmentación horizontal y vertical
  3. Camino (a): la base de datos relacional escalada y NewSQL
  4. Camino (b): NoSQL de estilo Dynamo con Cassandra
  5. Modelado orientado a consultas en Cassandra
  6. Camino (c): bases de datos documentales con MongoDB
  7. Tabla comparativa
  8. La decisión de Kilómetro Cero
  9. Práctica: Cassandra en docker-compose.yml y repositorio_cassandra.py
  10. Errores Comunes y Consejos
  11. Ejercicios
  12. Conclusión

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

  1. 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, pagos y reparto son fragmentos verticales del antiguo esquema, cada uno dueño de sus columnas, con pedido_id como 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_pedidos es un fragmento vertical que a su vez se particiona horizontalmente por cliente_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.

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

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

  1. 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_id ordenan 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') exige ALLOW FILTERING y 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, con PRIMARY KEY ((pedido_id)), y la aplicación escribe en ambas en un BATCH registrado (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.

  1. 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 nuevo W y R con 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 mongos enrutan 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.

  1. 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í (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 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

  1. 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:

  1. 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.
  2. 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).
  3. Disponibilidad: crear un pedido debe funcionar aunque un nodo o incluso un centro de datos falle; con LOCAL_QUORUM y 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.
  4. 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:

  1. Contadores CP: el stock es una restricción (CHECK (disponible >= 0)) que hay que hacer cumplir en cada reserva, con UPDATE … WHERE disponible >= 2 bloqueando 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 un UPDATE en PostgreSQL.
  2. 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.
  3. 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 un COMMIT.

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

  1. Práctica: Cassandra en docker-compose.yml y repositorio_cassandra.py

Tres 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 status
Datacenter: 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.
  • ExecutionProfile fija LOCAL_QUORUM por defecto, y cada sentencia puede sobrescribirlo con consistency_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_pedidos usa ONE en la ruta normal (la lista de pedidos tolera unos milisegundos de retraso) y LOCAL_QUORUM cuando la web la llama justo después de crear_pedido, que también fue LOCAL_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_estado muestra el coste de la desnormalización: para actualizar pedidos_por_cliente hace falta la clave completa (cliente_id, fecha, pedido_id), que se obtiene de pedidos_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-3

Con 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 pedidos normalizada con índices secundarios para cada consulta termina en ALLOW FILTERING y 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).
  • SimpleStrategy en producción. No entiende de racks ni de centros de datos: tres réplicas en el mismo rack. Siempre NetworkTopologyStrategy, incluso con un solo DC.
  • gc_grace_seconds menor 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 repair programado).
  • 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 counter no 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-driver es LOCAL_ONE. Fíjalo explícitamente en el ExecutionProfile y 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 resultado

Cada 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

Módulo 2: Comunicación en Sistemas Distribuidos

Módulo 3: Consistencia y Replicación

Módulo 4: Almacenamiento Distribuido

Módulo 5: Computación Distribuida

Módulo 6: Seguridad en Sistemas Distribuidos

Módulo 7: Monitoreo y Mantenimiento

Módulo 8: Casos de Estudio y Aplicaciones

© Copyright 2026. Todos los derechos reservados