Saltar al contenido

Estrategias de sharding en bases de datos: ventajas y compromisos

Guía práctica de sharding: estrategias de particionado, elección de la clave, consultas entre shards y la complejidad operativa de distribuir datos.

5 min de lectura
Diagrama de un clúster de bases de datos que muestra datos distribuidos entre múltiples shards con una capa de enrutamiento

El sharding es el último recurso para escalar una base de datos. El escalado vertical (más hardware) y las réplicas de lectura cubren la mayor parte del crecimiento. Pero cuando una sola instancia de base de datos no puede seguir el ritmo de escritura, o el conjunto de datos supera lo que cabe en una máquina, el sharding se vuelve necesario. Divides los datos entre varias instancias de base de datos, cada una con un subconjunto del conjunto total.

La compensación es brutal: el sharding elimina la simplicidad de una base de datos única. Los joins entre shards son lentos o imposibles. Las transacciones que abarcan varios shards requieren coordinación distribuida. Los cambios de esquema hay que aplicarlos en cada shard. Estás cambiando simplicidad por escalabilidad; asegúrate de que realmente necesitas la escalabilidad antes de pagar el coste de complejidad.

Estrategias de sharding

Hay tres estrategias principales para distribuir datos entre shards. Cada una hace distintas concesiones entre flexibilidad de consultas, distribución de datos y complejidad operativa.

tstypescript
// Strategy 1: Hash-based sharding
// Distribute rows based on a hash of the shard key
function getShardByHash(shardKey: string, totalShards: number): number {
  let hash = 0;
  for (let i = 0; i < shardKey.length; i++) {
    hash = ((hash << 5) - hash + shardKey.charCodeAt(i)) | 0;
  }
  return Math.abs(hash) % totalShards;
}
 
// Pro: Even distribution of data across shards
// Con: Range queries (e.g., "orders from last month") hit ALL shards
 
// Strategy 2: Range-based sharding
// Distribute rows based on value ranges of the shard key
interface ShardRange {
  shardId: number;
  minValue: string;
  maxValue: string;
}
 
const dateRanges: ShardRange[] = [
  { shardId: 0, minValue: '2020-01', maxValue: '2020-06' },
  { shardId: 1, minValue: '2020-07', maxValue: '2020-12' },
  { shardId: 2, minValue: '2021-01', maxValue: '2021-06' },
  { shardId: 3, minValue: '2021-07', maxValue: '2021-12' },
];
 
// Pro: Range queries hit only relevant shards
// Con: Hot spots — recent data concentrates on latest shard
 
// Strategy 3: Directory-based sharding
// Lookup table maps each entity to its shard
const shardDirectory = new Map<string, number>([
  ['tenant-a', 0],
  ['tenant-b', 1],
  ['tenant-c', 0],
  ['tenant-d', 2],
]);
 
// Pro: Full control over placement — can rebalance anytime
// Con: Directory itself becomes a bottleneck and single point of failure

Elección de la clave de shard

La clave de shard determina cómo se distribuyen los datos y qué consultas puede resolver un único shard. Una mala clave de shard lo complica todo. Una buena clave hace que la mayoría de las consultas sean rápidas.

tstypescript
// ❌ Bad shard key: user's country
// Problem: Uneven distribution — 50% of users might be in one country
// One shard gets half the traffic, others sit idle
interface BadSharding {
  shardKey: 'country';
  distribution: Record<string, number>;
}
const badDistribution: BadSharding = {
  shardKey: 'country',
  distribution: {
    US: 500_000,    // Shard 0 is overwhelmed
    UK: 80_000,     // Shard 1 is underutilized
    DE: 60_000,     // Shard 2 is underutilized
    other: 40_000,  // Shard 3 barely used
  },
};
 
// ✅ Good shard key: user_id (for user-centric applications)
// Even distribution — UUIDs hash uniformly
// Most queries are per-user — served by single shard
// User's orders, sessions, preferences all on one shard
interface GoodSharding {
  shardKey: 'user_id';
  properties: string[];
}
const goodDistribution: GoodSharding = {
  shardKey: 'user_id',
  properties: [
    'UUID/hash distributes evenly across shards',
    'Per-user queries hit exactly one shard',
    'User data locality — all related records co-located',
    'No hot spots from geographic concentration',
  ],
};
ymlyaml
# Shard key selection checklist
shard_key_requirements:
  high_cardinality:
    why: "Enough distinct values to distribute evenly"
    good: "user_id, order_id, tenant_id"
    bad: "country, status, boolean flags"
 
  query_alignment:
    why: "Most queries should include the shard key in WHERE clause"
    good: "Shard by user_id when 90% of queries filter by user"
    bad: "Shard by user_id when most queries filter by date range"
 
  even_distribution:
    why: "Prevent hot shards that receive disproportionate traffic"
    good: "UUIDs, auto-increment IDs with hash"
    bad: "Sequential timestamps, geographic codes"
 
  stability:
    why: "Changing shard key after data is distributed is extremely painful"
    good: "user_id doesn't change — data stays on its shard"
    bad: "Shard by subscription_tier — upgrades require data migration"

Consultas entre shards

Algunas consultas necesitan inherentemente datos de varios shards. Son caras, pero a veces inevitables.

tstypescript
// Single-shard query — fast (shard key in WHERE clause)
// "Get all orders for user_id = 'abc123'"
// → Route to shard hash('abc123'), query locally
async function getUserOrders(userId: string): Promise<Order[]> {
  const shardId = getShardByHash(userId, TOTAL_SHARDS);
  const shard = getShardConnection(shardId);
  return shard.query('SELECT * FROM orders WHERE user_id = $1', [userId]);
}
 
// Cross-shard query — expensive (scatter-gather)
// "Get top 10 orders by amount across all users"
// → Query ALL shards, merge results
async function getTopOrders(limit: number): Promise<Order[]> {
  const allShards = getAllShardConnections();
 
  // Scatter: query each shard in parallel
  const shardResults = await Promise.all(
    allShards.map((shard) =>
      shard.query(
        'SELECT * FROM orders ORDER BY amount DESC LIMIT $1',
        [limit]
      )
    )
  );
 
  // Gather: merge and re-sort
  return shardResults
    .flat()
    .sort((a, b) => b.amount - a.amount)
    .slice(0, limit);
}
tstypescript
// ❌ Cross-shard JOIN — prohibitively expensive
// Joining orders (sharded by user_id) with products (sharded by product_id)
// requires fetching data from potentially every shard on both tables
// "SELECT o.*, p.name FROM orders o JOIN products p ON o.product_id = p.id"
 
// ✅ Denormalize to avoid cross-shard JOINs
// Store product_name directly in the orders table
interface DenormalizedOrder {
  id: string;
  userId: string;
  productId: string;
  productName: string;   // Denormalized — copied from products table
  productCategory: string; // Denormalized — avoids JOIN
  amount: number;
  createdAt: Date;
}
// Trade-off: data duplication, but queries stay on one shard
// Update product name → need to update all orders (async job)

Complejidad operativa

El sharding añade preocupaciones operativas que los sistemas con una sola base de datos no tienen.

ymlyaml
# Operational challenges of sharding
schema_migrations:
  problem: "ALTER TABLE must run on every shard"
  approach: "Rolling migrations — apply to one shard at a time, verify, proceed"
  risk: "Shards temporarily have different schemas during migration window"
 
shard_rebalancing:
  problem: "Some shards grow faster than others"
  approach: "Split hot shards by moving half the key range to a new shard"
  risk: "Data migration during split — reads work, writes need coordination"
 
backup_and_restore:
  problem: "Backups must be coordinated across all shards for consistency"
  approach: "Point-in-time snapshots with global sequence numbers"
  risk: "Uncoordinated backups create inconsistent cross-shard state"
 
monitoring:
  problem: "Each shard has its own metrics — aggregate views needed"
  approach: "Dashboard showing per-shard size, latency, connections, replication lag"
  queries:
    - "SELECT pg_database_size(current_database()) per shard"
    - "Alert if any shard exceeds 80% capacity"
    - "Alert if query latency P99 diverges >2x between shards"
sqlsql
-- Per-shard health check query
-- Run against each shard to detect imbalances
 
-- Shard size and row counts
SELECT
  schemaname,
  tablename,
  pg_size_pretty(pg_total_relation_size(schemaname || '.' || tablename)) AS total_size,
  n_live_tup AS row_count
FROM pg_stat_user_tables
ORDER BY pg_total_relation_size(schemaname || '.' || tablename) DESC
LIMIT 10;
 
-- Active connections and slow queries per shard
SELECT
  count(*) AS active_connections,
  count(*) FILTER (WHERE state = 'active' AND now() - query_start > interval '5 seconds') AS slow_queries
FROM pg_stat_activity
WHERE datname = current_database();

Cuándo hacer sharding (y cuándo no)

El sharding es una estrategia de escalado de último recurso. La mayoría de aplicaciones nunca lo necesitan. Agota primero los enfoques más sencillos.

tstypescript
// Scaling progression — each step is cheaper than sharding
const scalingLadder = [
  '1. Optimize queries — add indexes, rewrite slow queries',
  '2. Vertical scaling — more CPU, RAM, faster disks',
  '3. Read replicas — offload read queries to replicas',
  '4. Caching layer — Redis/Memcached for hot data',
  '5. Table partitioning — split large tables within one database',
  '6. Sharding — distribute data across multiple databases',
] as const;
 
// Only consider sharding when:
// - Single instance cannot handle write throughput (steps 1-4 exhausted)
// - Dataset physically cannot fit on one machine (exceeds disk capacity)
// - Regulatory requirements demand data residence in specific regions

Conclusiones clave

  1. Haz sharding solo cuando hayas agotado los enfoques de escalado más sencillos — la optimización de consultas, el escalado vertical, las réplicas de lectura y el caché resuelven la mayoría de problemas
  2. Elige una clave de shard con alta cardinalidad y alineada con las consultas — la mayoría de consultas deberían incluir la clave de shard, y los valores deberían distribuirse de forma uniforme
  3. El sharding por hash distribuye de forma uniforme, pero hace que las consultas por rango sean caras — el sharding por rango facilita esas consultas, pero crea puntos calientes
  4. Desnormaliza para evitar joins entre shards — la duplicación de datos es el precio de mantener las consultas en un único shard
  5. Cada tarea operativa se complica — las migraciones, copias de seguridad, monitorización y reequilibrio se multiplican por el número de shards
  6. Monitoriza la salud de cada shard — shards desbalanceados anulan los beneficios de escalado y crean cuellos de botella
Wilfredo Rujel

Wilfredo Rujel

Ingeniero de Software Full Stack

Compartir esta publicaciónX