effidevFlutter · Edge de Cloudflare · Optimización de costes en la nube

Guía práctica de sharding en Cloudflare D1: cómo superar el límite de 10GB y el cuello de botella del writer único

Guía práctica de sharding en Cloudflare D1: cómo superar el límite de 10GB y el cuello de botella del writer único

Si alguna vez has montado un SaaS multitenant o un servicio de registro de eventos con mucha carga de escritura sobre Cloudflare D1, tarde o temprano chocas con este muro: una restricción estructural doble formada por un máximo de 10GB por base de datos y el hecho de que todas las escrituras se procesan secuencialmente a través de un único writer. En Postgres o MySQL bastaría con añadir más read replicas o recurrir a una extensión de particionado, pero D1 parte de una premisa distinta —es, en esencia, “múltiples instancias de SQLite distribuidas en el edge”—, así que el enfoque tiene que ser otro.

Este artículo cubre, con el nivel de detalle necesario para llevarlo realmente a producción, desde los criterios para distinguir cuándo basta con optimizar índices y cuándo hace falta sharding de verdad, pasando por el código para implementar enrutamiento basado en tenant ID con Worker + KV, el procedimiento para rebalancear shards sin downtime, hasta las estrategias para sortear la agregación entre shards.

Resumen clave

  • D1 tiene un límite máximo de 10GB (Paid) / 500MB (Free) por base de datos, y de 1TB (Paid) / 5GB (Free) a nivel de cuenta completa; además, cada base de datos procesa las escrituras secuencialmente a través de un writer único.
  • Un cuello de botella orientado a lectura debe resolverse primero con read replicas basadas en la Sessions API, no con sharding — las read replicas no aumentan en absoluto el throughput de escritura.
  • El sharding solo debe considerarse cuando (1) el tamaño de la base de datos se acerca a los 10GB o (2) la contención de escrituras concurrentes dispara la latencia. Antes de eso hay que revisar índices, escrituras en lote y archivado de datos fríos.
  • Para el mapeo de shards, un mapeo de directorio basado en KV + una base de datos de control favorece mucho más el rebalanceo que un esquema basado en hash.
  • Los JOIN y las transacciones entre shards no están soportados, así que la agregación debe resolverse mediante consultas fan-out o un almacén de rollup independiente.

El cuello de botella real que crean el límite de 10GB de D1 y su arquitectura de writer único

Las restricciones de D1 se dividen básicamente en dos niveles: el límite de tamaño por base de datos individual y el límite de capacidad de almacenamiento a nivel de cuenta. Según la documentación oficial de Cloudflare, las cifras son las siguientes.

Elemento Plan Free Plan Paid
Tamaño máximo por base de datos 500MB 10GB
Capacidad de almacenamiento total de la cuenta 5GB 1TB
Consultas por invocación de Worker 50 1.000
Longitud máxima de una sentencia SQL 100.000 bytes (100KB) Igual
Parámetros de binding por consulta 100 Igual
Tiempo máximo de ejecución de una consulta 30 segundos Igual
Tamaño máximo de una fila 2.000.000 bytes (2MB) Igual
Número máximo de columnas por tabla 100 Igual
Tamaño máximo de importación de archivo 5GB Igual

Lo importante aquí es la asimetría: “la capacidad total de la cuenta es holgada, pero cada base de datos individual se llena rápido”. Por ejemplo, en el plan Paid se puede usar hasta 1TB a nivel de cuenta, pero si se meten 100 tenants en una sola base de datos, esa base de datos concreta sigue topando en 10GB. Es decir, el cuello de botella real no es el límite de la cuenta, sino el muro de “10GB por base de datos”. Los servicios con tablas que solo crecen —logs de eventos, tablas de auditoría— llegan a este muro antes de lo que cabría esperar.

El segundo problema es la arquitectura de writer único. D1 funciona sobre el motor SQLite y, tal como indica la documentación, “cada base de datos procesa las consultas secuencialmente en un único hilo”. Para consultas cortas de 1ms se suele citar un techo teórico de alrededor de 1.000 consultas por segundo, pero las consultas de producción reales mezclan escaneos de índice, joins y transacciones de escritura, así que la latencia empieza a dispararse mucho antes de llegar a ese techo. El problema es que esto no es solo un problema de throughput, sino de concurrencia. Mientras se ejecuta una transacción de escritura pesada del tenant A, las solicitudes del tenant B que usa la misma base de datos quedan esperando en cola. En una arquitectura multitenant, esto se traduce directamente en un problema de noisy neighbor: el pico de tráfico de un tenant eleva la latencia de todos los demás.

Las señales típicas de que este cuello de botella se está manifestando en producción son:

Cómo decidir si hace falta sharding: cómo distinguir cargas orientadas a lectura vs. a escritura

El sharding tiene un coste alto. Implica construir una capa de enrutamiento, un procedimiento de rebalanceo y lógica para sortear la agregación entre shards, así que abordarlo con la mentalidad de “dividir de entrada por si acaso” termina en sobreingeniería. En la práctica, lo primero es diagnosticar la naturaleza de la carga de trabajo.

Si el cuello de botella es de lectura, el sharding no es la respuesta

Cuando la lentitud viene de un exceso de tráfico de lectura, D1 ofrece read replicas basadas en la Sessions API. Distribuyen las solicitudes hacia réplicas de lectura geográficamente cercanas para reducir la latencia, y como varias réplicas atienden lecturas en paralelo, el throughput de lectura también aumenta. Sin embargo, tal como especifica la documentación oficial, “todas las consultas de escritura siguen enviándose únicamente a la base de datos primary”, y las read replicas no aportan ninguna mejora al throughput de escritura. En otras palabras, plantearse el sharding antes que esto en una carga orientada a lectura es el orden equivocado: primero hay que agotar las read replicas, el caching con KV/Cache API y la optimización de consultas.

Si el cuello de botella es de escritura, hay pocas alternativas fuera del sharding

Por el contrario, si el problema es contención de escrituras concurrentes o el tamaño de una única base de datos acercándose a los 10GB, la historia es distinta. Las read replicas no ayudan en ninguno de los dos casos: sigue habiendo un único primary procesando todas las escrituras secuencialmente, y sigue habiendo un único archivo topando contra el muro de los 10GB. En este escenario, el sharding manual por tenant o por entidad es, en la práctica, la única vía de escalado horizontal.

Lista de verificación para decidir:

Si al menos tres de estos cuatro puntos se cumplen, es el momento de llevar el sharding a la fase de diseño. Si, por el contrario, la mayoría son “no”, la lista de verificación de indexación de la última sección más abajo puede dar varios meses más de margen.

Diseño de enrutamiento basado en tenant ID: implementación de una tabla de mapeo de shards con Worker + KV

Existen dos enfoques para la clave de shard: basado en hash y basado en directorio (tabla de mapeo). El enfoque basado en hash (shard = hash(tenant_id) % N) es simple de implementar, pero cambiar el número de shards o migrar de forma aislada a un único tenant requiere un rehashing masivo. El mapeo basado en directorio, en cambio, registra en una tabla separada el hecho de que “este tenant está ahora mismo en shard-3”, lo que permite un rebalanceo parcial: mover un único tenant concreto a otro shard. En entornos multitenant, el enfoque basado en directorio es casi siempre la opción recomendada.

La arquitectura se compone de tres capas.

  1. Base de datos de control (control DB): la fuente única de verdad (source of truth) del mapeo de shards. Se aloja en una base de datos D1 pequeña y separada, con una sola tabla shard_map.
  2. Caché en KV: consultar la control DB en cada solicitud convertiría a D1 en el propio cuello de botella, así que el resultado del mapeo se guarda como caché read-through en Workers KV.
  3. Bases de datos de shard: las N bases de datos D1 que contienen los datos reales de los tenants.
# wrangler.toml
name = "multitenant-api"
main = "src/index.ts"
compatibility_date = "2025-01-01"

[[d1_databases]]
binding = "CONTROL_DB"
database_name = "control-db"
database_id = "xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx"

[[d1_databases]]
binding = "SHARD_0"
database_name = "tenant-shard-0"
database_id = "xxxxxxxx-0000-xxxx-xxxx-xxxxxxxxxxxx"

[[d1_databases]]
binding = "SHARD_1"
database_name = "tenant-shard-1"
database_id = "xxxxxxxx-1111-xxxx-xxxx-xxxxxxxxxxxx"

[[kv_namespaces]]
binding = "SHARD_MAP_KV"
id = "yyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyy"

Esquema de la base de datos de control:

CREATE TABLE shard_map (
  tenant_id   TEXT PRIMARY KEY,
  shard_id    TEXT NOT NULL,       -- 'SHARD_0', 'SHARD_1' ...
  status      TEXT NOT NULL DEFAULT 'ACTIVE', -- ACTIVE | MIGRATING | DONE
  target_shard_id TEXT,            -- 리밸런싱 중일 때만 채워짐
  updated_at  INTEGER NOT NULL
);
CREATE INDEX idx_shard_map_status ON shard_map(status);

La lógica de enrutamiento consulta primero KV y solo recurre a la base de datos de control cuando hay un cache miss, escribiendo el resultado de vuelta en la caché (write-through):

// src/shard-router.ts
type Env = {
  CONTROL_DB: D1Database;
  SHARD_MAP_KV: KVNamespace;
  SHARD_0: D1Database;
  SHARD_1: D1Database;
  [key: string]: any;
};

interface ShardEntry {
  shardId: string;
  status: "ACTIVE" | "MIGRATING" | "DONE";
  targetShardId?: string;
}

export async function resolveShard(
  tenantId: string,
  env: Env
): Promise<D1Database> {
  const cacheKey = `shard:${tenantId}`;
  const cached = await env.SHARD_MAP_KV.get<ShardEntry>(cacheKey, "json");

  let entry: ShardEntry | null = cached;

  if (!entry) {
    const row = await env.CONTROL_DB
      .prepare(
        "SELECT shard_id, status, target_shard_id FROM shard_map WHERE tenant_id = ?"
      )
      .bind(tenantId)
      .first<{ shard_id: string; status: string; target_shard_id: string | null }>();

    if (!row) {
      throw new Error(`Unknown tenant: ${tenantId}`);
    }

    entry = {
      shardId: row.shard_id,
      status: row.status as ShardEntry["status"],
      targetShardId: row.target_shard_id ?? undefined,
    };

    // 읽기 전용 캐시. 리밸런싱 중에는 TTL을 짧게 둔다.
    const ttl = entry.status === "MIGRATING" ? 30 : 300;
    await env.SHARD_MAP_KV.put(cacheKey, JSON.stringify(entry), {
      expirationTtl: ttl,
    });
  }

  // MIGRATING 상태면 라이팅은 아직 원본 샤드로 보낸다 (다음 섹션 참고).
  const activeShardId = entry.status === "DONE" && entry.targetShardId
    ? entry.targetShardId
    : entry.shardId;

  const db = env[activeShardId] as D1Database | undefined;
  if (!db) throw new Error(`Shard binding not found: ${activeShardId}`);
  return db;
}

Al dar de alta un tenant nuevo, conviene añadir lógica que elija el shard con más margen disponible según su tamaño y carga actuales. Asignar por simple round-robin puede acumular tenants grandes en el mismo shard, así que es preferible tener en cuenta el volumen de datos esperado (por ejemplo, el nivel de plan) en el momento del alta — eso reduce la frecuencia con la que habrá que rebalancear más adelante.

Procedimiento para mover datos entre shards y rebalancear sin downtime

D1 no ofrece replicación nativa entre bases de datos ni una API de migración online, así que el rebalanceo tiene que implementarse como una máquina de estados explícita en la capa de aplicación. El procedimiento, usando la columna shard_map.status (ACTIVE → MIGRATING → DONE), es el siguiente:

  1. Selección del objetivo de migración: cuando el monitoreo de tamaño y carga muestra que un shard supera un umbral (por ejemplo, 8GB, o que un tenant concreto representa más del 40% de las escrituras de ese shard), se elige el tenant a mover. Poder elegir un único tenant entre varios es precisamente la ventaja del mapeo basado en directorio.
  2. Cambio de estado a MIGRATING: se actualiza el status de ese tenant en shard_map a MIGRATING y se rellena target_shard_id con el shard de destino. A partir de este punto, la aplicación sigue escribiendo en el shard original mientras publica de forma asíncrona los mismos eventos de escritura a través de Cloudflare Queues (dual write).
  3. Copia masiva: se leen paginadamente, con SELECT, las filas del tenant en el shard original y se insertan en el shard de destino con db.batch(). Hay que dimensionar el tamaño de los lotes considerando el límite de consultas por invocación de Worker (50 en Free / 1.000 en Paid) y el límite de 100 parámetros de binding. Si hay mucho volumen de datos y la latencia importa menos, es más sencillo optar por la vía offline: hacer un dump con wrangler d1 export --output=tenant.sql, filtrarlo y luego importarlo al destino con wrangler d1 execute (teniendo en cuenta el límite de 5GB para importación de archivos).
  4. Reproducción (replay) de la cola: una vez terminada la copia masiva, se reproducen en orden, contra el shard de destino, los eventos acumulados en la cola durante la migración (las escrituras ocurridas mientras tanto) para ponerse al día.
  5. Verificación de consistencia: se comparan el número de filas y checksums de tablas clave entre origen y destino (por ejemplo, SELECT COUNT(*), SUM(amount) FROM orders WHERE tenant_id = ?). Si hay discrepancias, se reintentan los pasos 3-4.
  6. Cutover: tras superar la verificación, se cambia shard_map.status a DONE y se actualiza shard_id al destino en una única transacción. A partir de este momento, las escrituras nuevas van únicamente al shard de destino.
  7. Espera a que expire la caché: se mantiene el shard original activo como fallback de solo lectura hasta que expire la entrada de enrutamiento cacheada en KV (de ahí que el código anterior fije un TTL corto de 30 segundos para el estado MIGRATING). Conviene dejar un margen de al menos 1-2 minutos para tener en cuenta el retraso de propagación global de KV.
  8. Limpieza: pasado el margen de gracia, se eliminan los datos de ese tenant del shard original para recuperar espacio.

La clave de este procedimiento es que no consiste en “parar por completo y mover”, sino en “dual write más un cutover gradual gobernado por el TTL de la caché”. Rebalancear sin downtime implica, en última instancia, asumir el coste de escribir temporalmente por duplicado en dos sitios a la vez.

Límites de las consultas entre shards y estrategias para sortear la agregación

El problema con el que se choca más a menudo tras aplicar sharding es la agregación entre shards — algo tan sencillo como un “total de todos los tenants”. Como los bindings de bases de datos de D1 están físicamente separados, no se admiten ni JOIN ni transacciones entre shards. Es una restricción dura sin forma de sortearla, así que hay que asumirla como premisa desde la propia fase de diseño.

Consultas fan-out

Si se necesita inmediatez y el número de shards es reducido (decenas, no cientos), la forma más simple es lanzar consultas en paralelo a todos los shards y sumar los resultados en el Worker.

async function totalOrdersAcrossShards(env: Env): Promise<number> {
  const shardBindings = ["SHARD_0", "SHARD_1", "SHARD_2"] as const;

  const results = await Promise.all(
    shardBindings.map((key) =>
      (env[key] as D1Database)
        .prepare("SELECT COUNT(*) AS cnt FROM orders")
        .first<{ cnt: number }>()
    )
  );

  return results.reduce((sum, r) => sum + (r?.cnt ?? 0), 0);
}

Este enfoque se ve afectado, a medida que crece el número de shards, por el límite de consultas por invocación de Worker y por la latencia total de la solicitud (el shard más lento determina el tiempo de respuesta global). Cuando el número de shards escala a cientos, el fan-out deja de ser adecuado para una ruta en tiempo real.

Tablas de rollup + ETL asíncrono

Para agregaciones que no requieren inmediatez al segundo —dashboards, informes—, es mucho más estable pre-agregar en el momento de escritura en lugar de hacer fan-out en el momento de la solicitud. Tras finalizar la transacción de escritura de cada shard, se publica un evento de “hace falta actualizar la agregación” mediante Cloudflare Queues; un Worker consumidor separado recoge ese evento y actualiza las filas de rollup en una base de datos D1 dedicada a analítica (o en Analytics Engine, si el carácter de la serie temporal es marcado). De este modo, el dashboard de administración no toca los shards en ningún momento — solo consulta esa única base de datos de rollup.

Las transacciones que cruzan shards se resuelven con el patrón saga

Para el caso, poco frecuente pero real, de necesitar una transacción atómica que abarque dos shards —como transferir un recurso entre tenants—, D1 no admite transacciones distribuidas, así que hay que resolverlo con un patrón saga basado en transacciones compensatorias. Es decir, dividir la operación en pasos idempotentes (“descontar en el shard A → si tiene éxito, sumar en el shard B → si falla, revertir en A”) y asignar a cada paso un ID de operación único y reintentable para evitar ejecuciones duplicadas.

Checklist de indexación y optimización de consultas a probar antes de hacer sharding

Es bastante habitual, una vez diseñada por completo una arquitectura de sharding, darse cuenta de que “en realidad, con arreglar los índices habría bastado”. Merece la pena agotar antes la siguiente checklist.

Si se ha aplicado toda esta checklist y aun así persisten los criterios mencionados antes (tamaño de base de datos acercándose a los 10GB, disparo de la latencia p95 de escritura), ese es el momento en el que realmente toca llevar a la práctica la arquitectura de sharding descrita en este artículo.