Saltar al contenido

Concurrencia asíncrona en Node.js: semáforos y backpressure

Un Promise.all sin límites es un OOM silencioso: así se construyen semáforos, colas con throttling y backpressure de streams en tus aplicaciones Node.js.

5 min de lectura
Diagrama de control de concurrencia en Node.js que muestra un semáforo regulando tareas asíncronas

Promise.all(items.map(fn)) es una de las líneas más peligrosas en un codebase de Node.js. Parece razonable: procesar todo en paralelo y esperar el resultado. Lo que realmente hace es disparar cada promesa simultáneamente sin ninguna consideración por la memoria, los rate limits de APIs externas o los pools de conexiones a la base de datos. Con 50 elementos funciona bien. Con 50.000 es un incidente.

Controlar la concurrencia es una preocupación de primer orden en los servicios Node.js en producción, no una ocurrencia tardía. Los patrones son sencillos de implementar y la recompensa es enorme.

El problema del paralelismo sin límites

La mayoría de los codebases empiezan con Promise.all porque funciona durante el desarrollo. La producción expone el fallo.

tstypescript
// ❌ Fires all 10,000 requests simultaneously
async function syncAllUsers(userIds: string[]): Promise<void> {
  await Promise.all(userIds.map((id) => fetchAndSync(id)));
}
 
// ✅ Processes 20 at a time — predictable memory, no rate-limit bans
async function syncAllUsers(userIds: string[]): Promise<void> {
  await runWithConcurrency(userIds, fetchAndSync, { limit: 20 });
}

La solución no es complicada, pero requiere una primitiva de concurrencia. Construyamos una.

Construyendo un semáforo

Un semáforo es un contador con un máximo. Cuando el contador está lleno, las nuevas tareas esperan hasta que se libere un espacio. Esta es la base sobre la que se construye todo lo demás.

tstypescript
export class Semaphore {
  private queue: Array<() => void> = [];
  private active = 0;
 
  constructor(private readonly limit: number) {}
 
  async acquire(): Promise<void> {
    if (this.active < this.limit) {
      this.active++;
      return;
    }
 
    await new Promise<void>((resolve) => {
      this.queue.push(resolve);
    });
    this.active++;
  }
 
  release(): void {
    this.active--;
    const next = this.queue.shift();
    if (next) next();
  }
 
  async run<T>(fn: () => Promise<T>): Promise<T> {
    await this.acquire();
    try {
      return await fn();
    } finally {
      this.release();
    }
  }
}

El semáforo es reutilizable durante toda la vida de un servicio. Decláralo una vez a nivel de módulo y compártelo entre todos los llamadores que usan el mismo recurso: un pool de base de datos, una API externa, un punto de montaje del sistema de archivos.

tstypescript
// Shared across the entire service — caps total in-flight Stripe calls
const stripeSemaphore = new Semaphore(5);
 
async function chargeCustomer(customerId: string, amount: number) {
  return stripeSemaphore.run(() =>
    stripe.charges.create({ amount, customer: customerId, currency: "usd" })
  );
}

Procesamiento por lotes con throttling

El semáforo maneja la concurrencia instantánea, pero con arrays grandes también conviene evitar encolar miles de promesas a la vez. Un helper runWithConcurrency procesa los elementos en oleadas controladas.

tstypescript
async function runWithConcurrency<T, R>(
  items: T[],
  fn: (item: T) => Promise<R>,
  options: { limit: number }
): Promise<R[]> {
  const semaphore = new Semaphore(options.limit);
  return Promise.all(items.map((item) => semaphore.run(() => fn(item))));
}

Esto sigue creando una promesa por elemento: con arrays extremadamente grandes (millones de registros), esa presión sobre el heap ya puede ser un problema. En esos casos, recurre a un iterador por chunks.

tstypescript
async function* chunks<T>(items: T[], size: number): AsyncGenerator<T[]> {
  for (let i = 0; i < items.length; i += size) {
    yield items.slice(i, i + size);
  }
}
 
async function processInBatches<T>(
  items: T[],
  fn: (item: T) => Promise<void>,
  batchSize = 100
): Promise<void> {
  for await (const batch of chunks(items, batchSize)) {
    await Promise.all(batch.map(fn));
    // Optional: yield to the event loop between batches
    await new Promise((r) => setImmediate(r));
  }
}

El yield con setImmediate permite que los callbacks de I/O (health checks, peticiones entrantes) se ejecuten entre lotes. Sin él, un bucle de procesamiento largo bloquea el event loop aunque cada tarea individual sea asíncrona.

~

Prefiere setImmediate sobre setTimeout(r, 0) entre lotes. setImmediate se dispara después de los callbacks de I/O en la iteración actual del event loop, mientras que setTimeout con delay 0 tiene en realidad un jitter de ~1ms en Node.js.

Backpressure con streams de Node.js

Cuando los volúmenes de datos son demasiado grandes para mantenerlos en memoria —piensa en trabajos ETL, exportaciones de CSV, procesamiento de logs— los streams con backpressure integrado son la herramienta adecuada. La clave es respetar el highWaterMark y no anularlo con operaciones asíncronas sin límites.

tstypescript
import { Transform, TransformCallback } from "node:stream";
 
class ConcurrentTransform extends Transform {
  private semaphore: Semaphore;
  private pending = 0;
  private drainCallback: (() => void) | null = null;
 
  constructor(
    private readonly fn: (chunk: unknown) => Promise<unknown>,
    concurrency: number
  ) {
    super({ objectMode: true, highWaterMark: concurrency * 2 });
    this.semaphore = new Semaphore(concurrency);
  }
 
  _transform(chunk: unknown, _enc: string, callback: TransformCallback): void {
    this.pending++;
    this.semaphore
      .run(() => this.fn(chunk))
      .then((result) => {
        this.push(result);
        this.pending--;
        if (this.pending === 0 && this.drainCallback) {
          this.drainCallback();
          this.drainCallback = null;
        }
      })
      .catch((err) => this.destroy(err));
 
    // Signal readiness for the next chunk immediately
    // so the readable side doesn't stall waiting for us
    callback();
  }
 
  _flush(callback: TransformCallback): void {
    if (this.pending === 0) return callback();
    this.drainCallback = callback;
  }
}

La idea crítica: llama a callback() inmediatamente en _transform para mantener el flujo del stream, pero usa el semáforo para limitar el trabajo realmente en vuelo. El highWaterMark del stream añade un buffer: configúralo en concurrency * 2 para que el lado readable siempre tenga trabajo encolado.

Cómo elegir el límite de concurrencia adecuado

No hay un número universal. Ajústalo según lo que estés protegiendo.

RecursoLímite inicial recomendadoPor qué
API HTTP externa (sin SLA)5–10Evita los 429; la mayoría de los límites gratuitos son por segundo
API HTTP externa (con SLA)Iguala su rate limitLee la documentación: Stripe y Twilio publican límites exactos
Pool de conexiones PostgreSQLpool.max - 2Deja margen para health checks y consultas de administración
Redis50–100Redis es rápido; el cuello de botella suele ser la red, no la BD
Trabajo CPU-bound (worker thread)os.cpus().length - 1Uno por núcleo menos uno para el hilo principal
Escrituras al sistema de archivos10–20Depende del disco: un SSD aguanta más, un disco mecánico mucho menos

Empieza de forma conservadora, haz pruebas de carga en condiciones realistas y luego sube el límite. Es mucho más fácil aumentar el límite de un semáforo que recuperarse de un fallo en cascada de la base de datos.

Evitar la acumulación de la cola bajo carga sostenida

La cola de un semáforo no tiene límite por defecto. Bajo carga sostenida, la cola crece sin tope: acabas de trasladar el problema de memoria del heap al semáforo. Añade un circuit breaker o descarta carga cuando la cola supere un umbral.

tstypescript
export class BoundedSemaphore extends Semaphore {
  constructor(
    limit: number,
    private readonly maxQueue: number
  ) {
    super(limit);
  }
 
  async acquire(): Promise<void> {
    // queue is private in parent — expose via a getter in real code
    if (this.queue.length >= this.maxQueue) {
      throw new Error("Semaphore queue full — shedding load");
    }
    return super.acquire();
  }
}

En el contexto de un servidor HTTP, captura este error en el manejador de la ruta y devuelve 503 Service Unavailable. Esto es load shedding: una degradación deliberada que protege al servicio de un colapso total.

tstypescript
export async function POST(req: Request) {
  try {
    const result = await processingQueue.run(() => handleRequest(req));
    return Response.json(result);
  } catch (err) {
    if (err instanceof Error && err.message.includes("queue full")) {
      return new Response("Service overloaded", { status: 503 });
    }
    throw err;
  }
}

Conclusiones clave

  1. Promise.all sobre arrays grandes no es seguro: crea paralelismo sin límites y provocará un OOM o disparará rate limits en producción.
  2. Un semáforo es la primitiva adecuada: son 30 líneas de código, sin dependencias, y se combina limpiamente con async/await.
  3. Comparte los semáforos en la frontera del recurso: un semáforo por pool de base de datos, por API externa, por punto de montaje de disco. No los crees por petición.
  4. Para conjuntos de datos muy grandes, combina el procesamiento por lotes con un semáforo: los iteradores por chunks evitan la presión sobre el heap de las promesas pendientes; el semáforo limita el trabajo en vuelo.
  5. Los streams de Node.js necesitan backpressure explícito: llama al callback() de _transform inmediatamente, pero condiciona el trabajo real a un semáforo; ajusta el highWaterMark para amortiguar sin inundar.
  6. Añade un límite a la cola y descarta carga: una cola sin límite es un OOM diferido. Rechaza pronto con un 503 cuando la cola esté llena.
Wilfredo Rujel

Wilfredo Rujel

Ingeniero de Software Full Stack

Compartir esta publicaciónX