Replicación, particionamiento y sharding | Nicolás Garzón
Replicar significa mantener copias de los mismos datos. Particionar significa dividirlos en subconjuntos. Sharding distribuye esas particiones entre nodos. Cada técnica resuelve un problema diferente y añade nuevos modos de fallo.
Cuando los datos viven en un solo nodo, muchas garantías parecen simples: una escritura tiene una autoridad clara, una consulta ve una única copia y una transacción puede coordinar recursos localmente.
Al distribuirlos aparecen preguntas nuevas:
¿Qué copia acepta escrituras?
¿Cuándo se considera confirmada una operación?
¿Qué puede observar un lector retrasado?
¿Cómo se detecta un nodo perdido?
¿Quién decide el nuevo líder?
¿Cómo se enruta una consulta al shard correcto?
¿Qué ocurre durante un rebalanceo?
Distribuir datos puede aumentar capacidad, disponibilidad o cercanía geográfica, pero también introduce coordinación, staleness, fallos parciales y recuperación más compleja.
La replicación mantiene varias copias lógicas del mismo conjunto de datos.
Tolerar fallos de nodo.
Aumentar capacidad de lectura.
Reducir latencia geográfica.
Facilitar mantenimiento.
Crear réplicas para reporting o backups.
No aumenta automáticamente la capacidad de escritura. En muchos modelos, todas las escrituras siguen pasando por una autoridad.
Un líder acepta escrituras y propaga cambios a seguidores.
Texto
Copiar cliente escribe
→ leader confirma
→ log de cambios
→ follower A aplica
→ follower B aplica
Orden de escritura centralizado.
Failover comprensible.
Réplicas de lectura.
Buen encaje para muchas bases relacionales.
El líder puede ser cuello de botella.
Los seguidores pueden estar retrasados.
Failover incorrecto puede perder cambios.
Promover un nodo desactualizado rompe expectativas.
Un usuario crea un pedido y luego es enviado a una réplica retrasada. Puede parecer que el pedido desapareció.
Leer temporalmente desde líder.
Mantener afinidad de sesión.
Esperar una posición de replicación.
Mostrar estado pendiente.
Usar garantías de sesión del motor.
El líder espera confirmación de una o más réplicas antes de responder.
Texto
Copiar write
→ leader persiste
→ follower confirma
→ respuesta al clienteReduce pérdida potencial si el líder falla, pero añade latencia y depende de disponibilidad de réplicas.
El líder responde antes de que seguidores confirmen.
Texto
Copiar write
→ leader persiste
→ respuesta
→ follower aplica despuésMejora latencia y disponibilidad de escritura, pero existe una ventana en la que el cambio puede perderse durante failover.
La decisión debe expresarse mediante RPO, latencia y tolerancia a indisponibilidad, no con “más seguro” o “más rápido”.
Varios líderes aceptan escrituras, normalmente por región o entorno desconectado.
Texto
Copiar región A leader ⇄ región B leaderAporta escritura local y tolerancia regional, pero permite conflictos.
Texto
Copiar A cambia email a x@example.com
B cambia email a y@example.com
Last-write-wins.
Merge de campos.
Version vectors.
Resolución manual.
Reglas específicas del dominio.
Para datos financieros o inventario, resolver por timestamp suele ser insuficiente porque pierde intención.
En un modelo leaderless, el cliente o coordinador escribe en varias réplicas.
Texto
Copiar N = 3 réplicas
W = 2 confirmaciones de escritura
R = 2 respuestas de lecturaLa regla R + W > N aumenta la posibilidad de intersectar una copia reciente, pero no garantiza por sí sola linearizability. Importan:
Versionado.
Sloppy quorums.
Réplicas caídas.
Read repair.
Anti-entropy.
Relojes.
Conflictos concurrentes.
El lag es la distancia temporal o de posición entre líder y réplica.
Red lenta.
Escritura intensa.
Réplica con CPU o I/O saturado.
Transacción grande.
Índices costosos.
Mantenimiento.
Esquema incompatible.
Segundos de retraso.
Bytes o posiciones pendientes.
Tiempo desde último replay.
Tasa de aplicación.
Error y estado de conexión.
Un valor global puede ocultar que una réplica específica está detenida.
Failover sustituye un nodo no disponible.
Detectar que el líder no responde.
Evitar falsos positivos por una pausa temporal.
Elegir una réplica suficientemente actualizada.
Obtener quorum o autoridad para promoverla.
Redirigir clientes.
Reintegrar el nodo anterior sin crear doble líder.
El problema difícil no es iniciar otro proceso. Es distinguir fallo real de partición de red y evitar split brain.
Dos nodos creen ser líderes y aceptan escrituras divergentes.
Quorum.
Fencing tokens.
Leases.
STONITH o aislamiento del nodo anterior.
Coordinación externa.
Reglas de promoción estrictas.
Un lock o lease expirado no basta si el nodo antiguo puede seguir modificando recursos. El receptor debe rechazar tokens de fencing antiguos.
Una réplica protege frente a fallo físico, pero reproduce:
Borrado accidental.
Corrupción lógica.
Migración destructiva.
Ransomware con acceso compartido.
Se necesitan backups aislados, retención, point-in-time recovery y pruebas de restore.
Particionar divide un conjunto de datos.
Texto
Copiar tabla orders
├── partición 2026-01
├── partición 2026-02
└── partición 2026-03Puede ocurrir dentro del mismo nodo o distribuirse entre varios.
Reducir datos examinados.
Facilitar retención.
Aislar workloads.
Distribuir capacidad.
Mantener locality.
Divide valores contiguos: fecha, ID o región.
Consultas de rango eficientes.
Retención sencilla por partición.
Hotspot en el rango actual.
Distribución desigual.
Aplica una función a la clave.
Distribución más uniforme.
Consultas de rango difíciles.
Rebalanceo según implementación.
Un catálogo indica dónde vive cada clave o tenant.
Directorio crítico.
Cache invalidation y disponibilidad.
Todos los datos de un tenant permanecen juntos.
Locality.
Exportación y aislamiento.
Routing simple.
Tenants desiguales.
Un tenant enorme crea hotspot.
Consultas globales requieren fan-out.
Útil para logs, métricas y datos con retención.
Riesgo principal: todas las escrituras actuales pueden concentrarse en una partición.
Sharding distribuye particiones entre nodos independientes.
Texto
Copiar router
├── shard 1: tenants A–F
├── shard 2: tenants G–M
└── shard 3: tenants N–ZEl router necesita conocer la ubicación o calcularla.
Balance.
Localidad.
Consultas.
Transacciones.
Rebalanceo.
Hotspots.
Crecimiento futuro.
Una buena clave suele tener:
Cardinalidad suficiente.
Distribución razonable.
Presencia en consultas dominantes.
Estabilidad.
Capacidad de mantener operaciones relacionadas juntas.
No existe una clave perfecta.
Favorece operaciones locales por negocio, pero un tenant grande puede saturar un shard.
Distribuye pedidos, pero consultas por tenant pueden tocar muchos shards si la clave no incluye tenant.
Favorece rangos, pero concentra escrituras recientes.
tenant_id + bucket puede distribuir tenants grandes conservando algo de locality. Aumenta complejidad de routing y consultas.
Un hotspot aparece cuando una partición recibe carga desproporcionada.
Clave monotónica.
Tenant grande.
Producto viral.
Contador global.
Lock o fila central.
Índice secundario con mala distribución.
Bucketing.
Salting controlado.
Particiones dedicadas.
Agregación por lotes.
Caching.
Separar lectura y escritura.
Rediseñar invariante global.
Cada mitigación complica consultas o consistencia.
Si la clave de routing no está presente, el sistema puede enviar la consulta a todos los shards.
Texto
Copiar query global
→ shard 1
→ shard 2
→ shard 3
→ merge y ordenar
Latencia depende del shard más lento.
Más tráfico.
Paginación difícil.
Fallos parciales.
Costos de ordenar y agregar.
Índice global.
Proyección analítica.
Denormalización.
Restricción funcional.
Servicio de búsqueda.
Cada shard indexa sus datos. Son simples de mantener, pero una búsqueda global requiere fan-out.
Un índice separado localiza registros entre shards.
Facilita consultas, pero introduce:
Doble escritura o eventos.
Consistencia eventual.
Ownership adicional.
Reconstrucción.
Nuevo punto de fallo.
Una transacción distribuida requiere coordinación como two-phase commit u otro protocolo.
Más latencia.
Bloqueos prolongados.
Coordinador.
Recuperación de estados inciertos.
Menor disponibilidad durante particiones.
Por eso se diseñan límites para mantener invariantes dentro de un shard cuando sea posible.
No se debe romper una invariante solo para evitar coordinación; quizá la elección de shard key sea incorrecta.
Al crecer, los datos deben moverse.
Split de rangos.
Consistent hashing.
Virtual nodes.
Directorio reasignable.
Copia gradual y cambio de routing.
Un flujo seguro puede ser:
Texto
Copiar crear destino
→ copiar snapshot
→ capturar cambios incrementales
→ verificar
→ cambiar lecturas
→ cambiar escrituras
→ retirar origenDurante la migración deben controlarse escrituras concurrentes y rollback.
Escribir en origen y destino puede parecer sencillo, pero una de las dos operaciones puede fallar.
Alternativas más seguras:
Change data capture.
Outbox.
Log de migración.
Escritura primaria con replicación.
Reconciliación y checksums.
La migración necesita métricas de progreso, divergencia y lag.
DomiSys podría comenzar con una base compartida indexada por business_id.
Si algunos negocios crecen:
Texto
Copiar tenants pequeños → pool compartido
enterprise → shard o base dedicadaEl directorio de tenants decide routing. La aplicación debe evitar asumir que todos viven en la misma conexión.
Una estrategia híbrida conserva simplicidad inicial y permite aislamiento selectivo.
Puede perderse una escritura ya confirmada al cliente si no estaba replicada. Define RPO y política de promoción.
El usuario no ve su cambio. Aplica read-your-writes o UI pendiente.
Ambos aceptan cambios. Se necesita conflicto semántico o bloquear una región.
Una consulta global puede devolver resultados parciales. Decide si fallar todo, marcar parcial o usar datos antiguos.
Jobs, cachés y eventos deben utilizar routing actualizado. Un ID físico incrustado en mensajes antiguos puede romper procesamiento.
La copia compite por I/O y red. Usa throttling, ventanas, prioridad y observación.
Las réplicas de lectura pueden estar retrasadas. La aplicación debe conocer la garantía.
Replica errores lógicos. Mantén backups aislados y restaura de prueba.
Se paga complejidad sin necesidad. Primero optimiza consultas, índices, hardware, retención y caché.
Puede destruir locality y exigir fan-out para todas las consultas importantes.
La distribución por tenant parece uniforme en cantidad, pero no en carga o tamaño.
Mover datos sin checksums y reconciliación puede perder o duplicar registros.
Lag por réplica.
Estado del stream.
Errores de aplicación.
Posición de log.
Tiempo y causa de failover.
Lecturas stale reportadas.
Tamaño y QPS por shard.
Distribución de claves.
Hot keys.
Fan-out por consulta.
Latencia por shard.
Tráfico de rebalanceo.
RPO real.
RTO medido.
Éxito de restore.
Integridad posterior.
Detener líder durante escritura.
Introducir lag.
Cortar red entre nodos.
Promover réplica.
Verificar read-your-writes.
Ejecutar consultas durante rebalanceo.
Simular tenant grande.
Restaurar desde backup.
Comparar checksums entre origen y destino.
Cuando necesitas tolerancia de nodo, lecturas adicionales o cercanía, aceptando semántica de lag.
Cuando volumen, retención o mantenimiento justifican dividir datos sin distribuir todavía.
Cuando un nodo o cluster no satisface capacidad o aislamiento y existe una clave de distribución compatible con el dominio.
Replicar copia; particionar divide; sharding distribuye.
Réplicas asíncronas permiten stale reads y pérdida durante failover.
La shard key es una decisión arquitectónica difícil de cambiar.
Consultas y transacciones cross-shard añaden coordinación.
Backups y réplicas resuelven riesgos diferentes.
Rebalanceo y resharding son procesos operativos que necesitan verificación.
¿Por qué una réplica de lectura puede romper la experiencia después de crear un pedido?
¿Qué riesgo introduce promover una réplica retrasada?
¿Cómo afecta tenant_id como shard key a un tenant muy grande?
¿Por qué una consulta global aumenta tail latency?
¿Qué diferencia existe entre consistent hashing y una garantía de consistencia?
Ver respuestas orientativas
Porque el cambio puede no haberse aplicado aún y la lectura devuelve estado antiguo.
Puede perder escrituras confirmadas o retroceder el estado visible.
Mantiene locality, pero concentra todo su tamaño y carga en un único shard.
Porque espera y combina respuestas de múltiples shards, incluyendo el más lento.
Consistent hashing distribuye claves y facilita rebalanceo; no define qué valores observan lecturas concurrentes.
Consistencia, CAP y PACELC profundiza en las garantías observables cuando existen réplicas, particiones y decisiones entre coordinación, latencia y disponibilidad.