Zum Inhalt springen

Eine Task Queue von Grund auf mit Node.js und Redis bauen

Baue Schritt für Schritt eine produktionsreife Task Queue mit Node.js und Redis: zuverlässige Zustellung, Retries, Dead Letter Queues, Nebenläufigkeit.

6 Min. Lesezeit
Architekturdiagramm, das eine Task Queue mit Produzenten, Redis-Broker und konsumierenden Workern zeigt

Jede Anwendung braucht irgendwann Hintergrundverarbeitung. E-Mail-Benachrichtigungen, Bildskalierung, Report-Generierung, Webhook-Zustellung — diese Aufgaben gehören nicht in den Request-Response-Zyklus. Eine Task Queue entkoppelt Produzenten von Konsumenten und ermöglicht zuverlässige asynchrone Verarbeitung.

Auch wenn es Bibliotheken wie BullMQ gibt: Eine Queue von Grund auf zu bauen lehrt dich die Primitive, auf denen verteilte Jobverarbeitung basiert. Du verstehst genau, was passiert, wenn ein Job fehlschlägt, wie Nebenläufigkeit gesteuert wird und warum bestimmte Designentscheidungen wichtig sind.

Die grundlegende Queue-Struktur

Eine Task Queue braucht drei Dinge: eine Möglichkeit, Jobs einzureihen, eine Möglichkeit, sie zuverlässig zu entnehmen, und eine Möglichkeit, ihren Zustand zu verfolgen. Redis liefert alle Bausteine.

tstypescript
import { createClient, RedisClientType } from "redis";
 
interface Job<T = unknown> {
  id: string;
  queue: string;
  payload: T;
  attempts: number;
  maxAttempts: number;
  createdAt: number;
  processAfter: number;
}
 
class TaskQueue {
  private client: RedisClientType;
  private prefix: string;
 
  constructor(client: RedisClientType, prefix: string = "tq") {
    this.client = client;
    this.prefix = prefix;
  }
 
  private key(queue: string, suffix: string): string {
    return `${this.prefix}:${queue}:${suffix}`;
  }
 
  async enqueue<T>(
    queue: string,
    payload: T,
    options: { delay?: number; maxAttempts?: number } = {}
  ): Promise<string> {
    const id = crypto.randomUUID();
    const now = Date.now();
 
    const job: Job<T> = {
      id,
      queue,
      payload,
      attempts: 0,
      maxAttempts: options.maxAttempts ?? 3,
      createdAt: now,
      processAfter: now + (options.delay ?? 0),
    };
 
    const multi = this.client.multi();
 
    // Store job data
    multi.set(
      this.key(queue, `job:${id}`),
      JSON.stringify(job)
    );
 
    // Add to waiting sorted set (scored by processAfter)
    multi.zAdd(this.key(queue, "waiting"), {
      score: job.processAfter,
      value: id,
    });
 
    await multi.exec();
    return id;
  }
}

Ein Sorted Set für die Warteschlange ermöglicht verzögerte Jobs ganz natürlich: Jobs werden nach ihrem processAfter-Zeitstempel bewertet, sodass wir nur Jobs aufnehmen, deren Score den aktuellen Zeitpunkt erreicht hat oder davor liegt.

Zuverlässige Job-Entnahme

Die entscheidende Herausforderung bei jeder Queue ist die Exactly-Once-Verarbeitung. Stürzt ein Worker mitten in einem Job ab, darf der Job nicht verloren gehen. Redis' ZPOPMIN kombiniert mit einem Processing-Set stellt diese Garantie bereit.

tstypescript
// ❌ Unreliable: job lost if worker crashes after pop
async function unsafeDequeue(client: RedisClientType, queue: string) {
  const result = await client.zPopMin(queue);
  // If process crashes here, job is gone forever
  return result;
}
tstypescript
// ✅ Reliable: job tracked in processing set
class TaskQueue {
  // ... previous code
 
  async dequeue(queue: string): Promise<Job | null> {
    const now = Date.now();
    const waitingKey = this.key(queue, "waiting");
    const processingKey = this.key(queue, "processing");
 
    // Atomically move job from waiting to processing
    const script = `
      local result = redis.call('ZRANGEBYSCORE', KEYS[1], '-inf', ARGV[1], 'LIMIT', 0, 1)
      if #result == 0 then return nil end
      local jobId = result[1]
      redis.call('ZREM', KEYS[1], jobId)
      redis.call('ZADD', KEYS[2], ARGV[2], jobId)
      return jobId
    `;
 
    const jobId = await this.client.eval(script, {
      keys: [waitingKey, processingKey],
      arguments: [String(now), String(now)],
    });
 
    if (!jobId) return null;
 
    const jobData = await this.client.get(
      this.key(queue, `job:${jobId}`)
    );
 
    if (!jobData) return null;
    return JSON.parse(jobData) as Job;
  }
 
  async acknowledge(queue: string, jobId: string): Promise<void> {
    const multi = this.client.multi();
 
    multi.zRem(this.key(queue, "processing"), jobId);
    multi.del(this.key(queue, `job:${jobId}`));
    multi.incr(this.key(queue, "stats:completed"));
 
    await multi.exec();
  }
}

Das Lua-Skript läuft atomar in Redis — kein anderer Client kann zwischen der Prüfung des Warte-Sets und dem Verschieben des Jobs ins Processing-Set dazwischenfunken. Damit ist die Race Condition ausgeschlossen, bei der zwei Worker denselben Job greifen könnten.

Retry-Logik und Dead Letter Queues

Jobs schlagen fehl. Netzwerke laufen in Timeouts, externe APIs liefern Fehler, Daten sind fehlerhaft. Eine robuste Queue wiederholt vorübergehende Fehler und leitet dauerhafte Fehler zur Untersuchung in eine Dead Letter Queue weiter.

tstypescript
class TaskQueue {
  // ... previous code
 
  async fail(
    queue: string,
    jobId: string,
    error: string
  ): Promise<"retried" | "dead-lettered"> {
    const jobData = await this.client.get(
      this.key(queue, `job:${jobId}`)
    );
 
    if (!jobData) throw new Error(`Job ${jobId} not found`);
 
    const job: Job = JSON.parse(jobData);
    job.attempts += 1;
 
    // Remove from processing set
    await this.client.zRem(this.key(queue, "processing"), jobId);
 
    if (job.attempts < job.maxAttempts) {
      // Exponential backoff: 1s, 4s, 9s, 16s...
      const delay = Math.pow(job.attempts, 2) * 1000;
      job.processAfter = Date.now() + delay;
 
      const multi = this.client.multi();
      multi.set(
        this.key(queue, `job:${jobId}`),
        JSON.stringify(job)
      );
      multi.zAdd(this.key(queue, "waiting"), {
        score: job.processAfter,
        value: jobId,
      });
      await multi.exec();
 
      return "retried";
    }
 
    // Max attempts exceeded: dead letter queue
    const multi = this.client.multi();
    multi.lPush(
      this.key(queue, "dead"),
      JSON.stringify({ ...job, error, failedAt: Date.now() })
    );
    multi.del(this.key(queue, `job:${jobId}`));
    multi.incr(this.key(queue, "stats:dead-lettered"));
    await multi.exec();
 
    return "dead-lettered";
  }
}

Exponentielles Backoff mit der Formel attempts² × 1000ms ergibt immer längere Wartezeiten: 1 Sekunde, 4 Sekunden, 9 Sekunden. Das verhindert Retry-Stürme, wenn ein nachgelagerter Dienst kämpft, und gibt ihm Zeit zur Erholung.

Die Worker-Schleife bauen

Ein Worker fragt kontinuierlich Jobs ab, verarbeitet sie und behandelt Erfolg oder Fehlschlag. Das Polling-Intervall balanciert Reaktionsfähigkeit gegen Redis-Last aus.

tstypescript
type JobHandler<T = unknown> = (payload: T) => Promise<void>;
 
class Worker {
  private queue: TaskQueue;
  private queueName: string;
  private handler: JobHandler;
  private running: boolean = false;
  private concurrency: number;
  private activeJobs: number = 0;
 
  constructor(
    queue: TaskQueue,
    queueName: string,
    handler: JobHandler,
    concurrency: number = 5
  ) {
    this.queue = queue;
    this.queueName = queueName;
    this.handler = handler;
    this.concurrency = concurrency;
  }
 
  async start(): Promise<void> {
    this.running = true;
    console.log(
      `Worker started for queue "${this.queueName}" ` +
      `(concurrency: ${this.concurrency})`
    );
 
    while (this.running) {
      if (this.activeJobs >= this.concurrency) {
        await this.sleep(100);
        continue;
      }
 
      const job = await this.queue.dequeue(this.queueName);
 
      if (!job) {
        await this.sleep(1000); // No jobs available, wait
        continue;
      }
 
      this.activeJobs++;
      this.processJob(job).finally(() => {
        this.activeJobs--;
      });
    }
  }
 
  private async processJob(job: Job): Promise<void> {
    try {
      await this.handler(job.payload);
      await this.queue.acknowledge(this.queueName, job.id);
    } catch (error) {
      const message =
        error instanceof Error ? error.message : "Unknown error";
      const result = await this.queue.fail(
        this.queueName,
        job.id,
        message
      );
      console.warn(
        `Job ${job.id} failed (${result}): ${message}`
      );
    }
  }
 
  stop(): void {
    this.running = false;
  }
 
  private sleep(ms: number): Promise<void> {
    return new Promise(resolve => setTimeout(resolve, ms));
  }
}

Die Nebenläufigkeitssteuerung stellt sicher, dass der Worker mehrere Jobs gleichzeitig verarbeitet, ohne das System zu überlasten. Jeder Job läuft unabhängig — ein langsamer Job blockiert die anderen nicht.

Wiederherstellung hängengebliebener Jobs

Stürzt ein Worker ab, bleiben seine Jobs auf unbestimmte Zeit im Processing-Set. Ein separater Recovery-Prozess muss hängengebliebene Jobs erkennen und erneut einreihen.

tstypescript
class TaskQueue {
  // ... previous code
 
  async recoverStalledJobs(
    queue: string,
    stallTimeout: number = 30000
  ): Promise<number> {
    const processingKey = this.key(queue, "processing");
    const cutoff = Date.now() - stallTimeout;
 
    // Find jobs that have been processing longer than stallTimeout
    const stalledIds = await this.client.zRangeByScore(
      processingKey,
      "-inf",
      String(cutoff)
    );
 
    let recovered = 0;
 
    for (const jobId of stalledIds) {
      const jobData = await this.client.get(
        this.key(queue, `job:${jobId}`)
      );
 
      if (!jobData) {
        // Job data missing, just clean up
        await this.client.zRem(processingKey, jobId);
        continue;
      }
 
      const job: Job = JSON.parse(jobData);
      job.attempts += 1;
 
      if (job.attempts >= job.maxAttempts) {
        // Exceeded retries, dead letter it
        const multi = this.client.multi();
        multi.zRem(processingKey, jobId);
        multi.lPush(
          this.key(queue, "dead"),
          JSON.stringify({
            ...job,
            error: "Stalled and exceeded max attempts",
            failedAt: Date.now(),
          })
        );
        multi.del(this.key(queue, `job:${jobId}`));
        await multi.exec();
      } else {
        // Re-enqueue for retry
        const multi = this.client.multi();
        multi.zRem(processingKey, jobId);
        multi.set(
          this.key(queue, `job:${jobId}`),
          JSON.stringify(job)
        );
        multi.zAdd(this.key(queue, "waiting"), {
          score: Date.now(),
          value: jobId,
        });
        await multi.exec();
        recovered++;
      }
    }
 
    return recovered;
  }
}

Führe den Recovery-Lauf per Timer aus — alle 30 Sekunden ist für die meisten Anwendungen ein sinnvoller Wert. Der Stall-Timeout sollte länger sein als die längste erwartete Jobdauer, damit keine Jobs wiederhergestellt werden, die noch aktiv verarbeitet werden.

Alles zusammenführen

So interagieren Produzenten und Konsumenten in einem realen Anwendungsszenario über die Queue.

tstypescript
async function main() {
  const redis = createClient({ url: "redis://localhost:6379" });
  await redis.connect();
 
  const taskQueue = new TaskQueue(redis);
 
  // Producer: enqueue email jobs
  await taskQueue.enqueue("emails", {
    to: "user@example.com",
    subject: "Welcome",
    template: "onboarding",
  });
 
  await taskQueue.enqueue(
    "emails",
    {
      to: "admin@example.com",
      subject: "Daily Report",
      template: "report",
    },
    { delay: 60000 } // Delay 1 minute
  );
 
  // Consumer: process email jobs
  const worker = new Worker(
    taskQueue,
    "emails",
    async (payload) => {
      const emailPayload = payload as {
        to: string;
        subject: string;
        template: string;
      };
      console.log(`Sending email to ${emailPayload.to}`);
      // await sendEmail(emailPayload);
    },
    3 // Process 3 emails concurrently
  );
 
  // Start stall recovery on interval
  setInterval(() => {
    taskQueue.recoverStalledJobs("emails").then(count => {
      if (count > 0) console.log(`Recovered ${count} stalled jobs`);
    });
  }, 30000);
 
  await worker.start();
}

Die wichtigsten Erkenntnisse

Eine Task Queue von Grund auf zu bauen offenbart die Komplexität, die hinter einfach aussehenden Jobverarbeitungs-APIs steckt. Die zentralen Herausforderungen liegen nicht im Einreihen von Daten — sie liegen in den Zuverlässigkeitsgarantien. Atomare Zustandsübergänge von Jobs verhindern Duplikate. Exponentielles Backoff verhindert Retry-Stürme. Dead Letter Queues verhindern stillen Datenverlust. Die Stall-Recovery verhindert, dass Jobs verschwinden, wenn Worker abstürzen.

Das Verständnis dieser Primitive macht dich zu einem besseren Nutzer von produktiven Queue-Bibliotheken und gibt dir die Grundlage, um Probleme zu debuggen, die in verteilten Jobverarbeitungssystemen unweigerlich auftreten.

Wilfredo Rujel

Wilfredo Rujel

Full-Stack-Softwareentwickler

Diesen Beitrag teilenX