Job Queue

Redis Queues

Eine Redis Job-Queue verschiebt langsame, wiederholbare Aufgaben hinter einen dauerhaften Buffer. BullMQ verwandelt Redis Lists und Streams in verzögerte Jobs, Retries, Prioritäten und wiederholbare Zeitpläne.

intermediate14 min readUpdated 16. Sept. 2026
queue.ts
ts
// queue.ts
import { Queue, Worker } from "bullmq";
import IORedis from "ioredis";

const connection = new IORedis(process.env.REDIS_URL!, {
  maxRetriesPerRequest: null,
});

export const emails = new Queue("emails", { connection });

new Worker(
  "emails",
  async (job) => {
    await sendEmail(job.data.to, job.data.template);
    return { sentAt: Date.now() };
  },
  { connection, concurrency: 10 },
);
Backing store
Redis lists und streams
Gängige Library
BullMQ
Delivery
At-least-once
Job-Status
waiting, active, completed, failed, delayed
Retry-Strategie
attempts plus exponential backoff
Bestens geeignet für
Verzögerte, wiederholbare Aufgaben pro Job

Warum es wichtig ist

Was eine Redis-basierte Queue bietet

Verzögerte und wiederholbare Jobs

Ein Job kann einmal in fünf Minuten oder dauerhaft nach einem Cron-Muster laufen. Der Zeitplan liegt in Redis, sodass ein Neustart ihn nicht vergisst.

Retries mit Backoff

Versuche, exponentielle Verzögerungen und das 'failed set' sind deklarative Job-Optionen anstelle von try/catch-Logik in jedem Handler.

Concurrency und Rate Limits

Jeder Worker gibt an, wie viele Jobs gleichzeitig und wie schnell ausgeführt werden, sodass eine nachgelagerte API niemals durch einen Burst überlastet wird.

Das Gesamtbild

Producer, Queue, Worker

Jede Redis Queue besteht aus denselben drei Rollen: Etwas fügt Arbeit hinzu, Redis hält sie dauerhaft fest und etwas anderes verarbeitet sie nach eigenem Zeitplan.

Producer

Enqueue

Ein Request-Handler fügt einen benannten Job mit einem kleinen Payload hinzu und kehrt sofort zurück. Er wartet nie darauf, dass die Arbeit abgeschlossen wird.

Queue

Buffer

Redis hält wartende, verzögerte, aktive und fehlgeschlagene Jobs, gibt sie atomar aus und übersteht einen Worker-Neustart.

Worker

Process

Ein langlebiger Node.js-Prozess zieht Jobs, führt den Handler unter einem Concurrency-Limit aus und meldet Erfolg, Fehler oder Fortschritt.

HTML5 auf einen Blick

Der Funktionsumfang von BullMQ

Queue

new Queue(name) ist der Handle, den Producer nutzen, um Jobs hinzuzufügen und zu inspizieren.

Worker

new Worker(name, handler) führt Jobs aus und emittiert completed, failed und progress Events.

Verzögerte Jobs

{ delay: 60_000 } führt einen Job in einer Minute aus, ohne dass ein separater Scheduler nötig ist.

Retries

{ attempts, backoff } wiederholt fehlgeschlagene Versuche und parkt erschöpfte Jobs im failed set.

Prioritäten

{ priority } lässt dringende Jobs vor Jobs mit niedriger Priorität springen.

Flows

Parent- und Child-Jobs setzen Pipelines und Fan-out-Bäume zusammen.

Ablauf

Der Lebenszyklus eines Jobs in Redis

Die gleichen sechs Schritte gelten, egal ob der Job beim ersten Mal erfolgreich ist oder seine Versuche erschöpft und im failed set wartet.

  1. 1

    Enqueue

    Der Producer fügt einen benannten Job mit Payload und Optionen hinzu. Redis speichert ihn in der Warteliste und der Aufruf kehrt sofort zurück.

  2. 2

    Claim

    Ein freier Worker blockiert an der Queue, verschiebt den nächsten fälligen Job atomar nach 'active' und setzt einen Lock, damit kein anderer Worker ihn beanspruchen kann.

  3. 3

    Process

    Der Handler wird mit den Job-Daten ausgeführt und kann den Fortschritt melden. Er gibt ein Ergebnis zurück oder wirft einen Fehler.

  4. 4

    Retry mit Backoff

    Bei einem Fehler steigt die Anzahl der Versuche und der Job wird nach einer wachsenden, mit Jitter versehenen Verzögerung erneut geplant und verschoben in das delayed set.

  5. 5

    Permanenter Fehler

    Nach dem letzten Versuch wird der Job mit seinem Stack-Trace in das failed set geschrieben, wo ein Mensch ihn inspizieren und erneut ausführen kann.

  6. 6

    Abschluss

    Bei Erfolg wird der Job entfernt, für eine begrenzte Anzahl behalten oder archiviert, je nach Einstellung von removeOnComplete.

Der vollständige Leitfaden

Redis Queues: Alles was Sie wissen müssen

Was ist eine Job-Queue?

Eine Job-Queue ist ein persistenter Buffer, der zwischen dem Code, der eine Aufgabe anfordert, und dem Code, der diese ausführt, sitzt. Der Producer schreibt einen kleinen Datensatz, der eine Arbeitseinheit beschreibt, und kehrt sofort zurück. Ein Worker liest diesen Datensatz später aus und erledigt den zeitintensiven Teil. Die Queue ist das Element, das einen Neustart übersteht, Lastspitzen abfängt und bei Fehlern als Auffangbecken dient.

Redis wird seit mehr als einem Jahrzehnt auf diese Weise eingesetzt. Seine Listen boten frühen Queue-Libraries atomare LPUSH- und BRPOP-Primitive, und Blocking Pops ermöglichten es einem Worker, zu schlafen, bis Arbeit eintraf, anstatt in einer Schleife zu pollen. Streams fügten später Consumer Groups, Acknowledgements und Replays hinzu. Auf Basis dieser Primitive kapselt BullMQ den gesamten Lebenszyklus — verzögerte Jobs, Retries, Prioritäten, Rate Limiting, wiederholbare Zeitpläne und Events — in einer kompakten TypeScript API.

Wenn Redis bereits Ihren Cache oder Ihre Sessions verwaltet, ist die Implementierung einer Queue nur ein kleiner Schritt. Diese Bequemlichkeit ist der Hauptgrund, warum es die Standardwahl für Node-Teams ist, und zugleich der wichtigste Grund, sowohl zu verstehen, was es bietet als auch, was nicht.

Eine Queue besteht aus drei Rollen, nicht aus einer

Jedes Queue-System, unabhängig vom Broker, hat dieselben drei Rollen. Diese strikt voneinander zu trennen, ist der Schlüssel zur Wartbarkeit des Systems.

Der Producer ist jeder Code, der einen Job hinzufügt. Er kennt den Namen des Jobs und die Struktur des Payloads, aber sonst nichts. Er muss schnell sein, da er normalerweise innerhalb eines Requests ausgeführt wird, und es sollte sicher sein, ihn zweimal aufzurufen.

app.post("/signup", async (req, res) => {
  const user = await db.user.create({ data: req.body });

  await emails.add("welcome", { userId: user.id, to: user.email });

  res.status(201).json({ id: user.id });
});

Die Queue ist der gemeinsame Zustand (shared state) in der Mitte. In einem Redis-Deployment ist dies die Menge der Keys, die BullMQ verwaltet: die Warteliste, das Delayed Set, das Active Set und das Failed Set. Sie persistiert Jobs über Neustarts hinweg, gibt sie atomar aus und verfolgt die Versuche.

Der Worker ist ein langlebiger Prozess, der Jobs abruft und Handler ausführt. Er ist aus guten Gründen von der API getrennt: Er kann auf CPU-optimierter Hardware bereitgestellt, an die Queue-Tiefe angepasst und neu gestartet werden, ohne dass Requests verloren gehen. Ein Prozess kann Worker für mehrere Queues hosten, und eine Queue kann von vielen Worker-Prozessen bedient werden. Die Queue ist der einzige gemeinsame Zustand, weshalb sie so sauber horizontal skaliert.

Warum Redis ein gängiger Backing Store ist

Drei Eigenschaften machen Redis zu einer natürlichen Queue.

Erstens: Die Primitiven existieren bereits. Eine Liste mit LPUSH und BRPOP ist eine Queue, und BRPOP blockiert die Verbindung, anstatt die CPU auszulasten. Das ist der einzige Grund, warum die ersten Node-Queue-Bibliotheken nur wenige hundert Zeilen lang waren.

LPUSH queue:emails "welcome:42"
BRPOP queue:emails 30

Zweitens: Befehle sind atomar. Ein Job von „wartend“ auf „aktiv“ zu verschieben, die Anzahl der Versuche zu erhöhen und den Lock zu setzen, kann alles ohne Race Conditions geschehen, da Redis Befehle nacheinander ausführt. Genau das benötigt eine Queue, um zu verhindern, dass zwei Worker denselben Job beanspruchen.

Drittens: Redis ist wahrscheinlich sowieso schon vorhanden. Es ist der Standard-Cache und Session-Store für Node-Services. Die Wiederverwendung für eine Queue vermeidet ein zweites Infrastruktur-Komponente, einen zweiten Satz an Zugangsdaten und ein zweites Runbook.

Streams erweitern dieses Konzept weiter. Während eine Liste nur gepusht und gepoppt werden kann, ist ein Stream ein Append-only-Log mit Consumer-Gruppen, Bestätigungen pro Nachricht (Acknowledgements) und einem wiederholbaren Verlauf.

XADD jobs:emails '*' type welcome userId 42
XREADGROUP GROUP workers alice COUNT 10 STREAMS jobs:emails '>'

Der Haken ist die Durability. Redis ist primär ein In-Memory-Store, sodass ein bestätigter, aber noch nicht auf die Festplatte geschriebener Job verloren gehen kann, wenn die Instanz abstürzt. Man kann dieses Zeitfenster mit AOF und einem Replica verkleinern, aber man kann Redis nicht so durable machen wie eine Datenbank mit Write-Ahead-Logging. Betrachten Sie einen Job als wiederherstellbare Arbeit, nicht als Ihr System of Record.

Erste Schritte mit BullMQ

BullMQ benötigt eine Redis-Verbindung, die unendlich viele Wiederholungsversuche zulässt. Das Standardverhalten von ioredis gibt einen Befehl nach einigen Fehlversuchen auf, was für einen Worker, der kurze Aussetzer überbrücken muss, nicht korrekt ist. Setzen Sie daher maxRetriesPerRequest: null.

import { Queue, Worker } from "bullmq";
import IORedis from "ioredis";

const connection = new IORedis(process.env.REDIS_URL!, {
  maxRetriesPerRequest: null,
});

const emails = new Queue("emails", { connection });

const worker = new Worker(
  "emails",
  async (job) => {
    await sendEmail(job.data.to, job.data.template);
  },
  { connection, concurrency: 10 },
);

Die Queue ist der Producer-Handle. worker ist der Consumer. Der Name "emails" bezeichnet die Queue, und ein Worker sieht nur Jobs, die seiner eigenen Queue hinzugefügt wurden. Separate Queues dienen als Isolationseinheiten: Eine langsame Import-Queue kann so eine Passwort-Reset-Queue nicht verzögern.

Eine Queue ist leichtgewichtig und kann überall dort erstellt werden, wo Sie einen Job hinzufügen möchten. Ein Worker ist ein langlebiger Prozess und sollte einmal pro Prozess gestartet werden, nicht pro Request. Nutzen Sie das Connection-Objekt gemeinsam für alle Queues und Worker innerhalb des Prozesses.

Jobs mit Daten und Optionen hinzufügen

add erwartet einen Namen, ein Payload und ein Options-Objekt. Der Name routet den Job an einen Handler; der Payload enthält alles, was der Worker benötigt.

await emails.add("welcome", { userId, to }, {
  jobId: `welcome:${userId}`,
});

Der Payload sollte eine Momentaufnahme der Absicht sein, kein Live-Objekt. Wenn ein Benutzer seine E-Mail-Adresse zwischen dem Enqueue und der Ausführung ändert, sollte der Job immer noch an die Adresse gesendet werden, mit der er erstellt wurde. Speichern Sie IDs und die wenigen Werte, die die Arbeit definieren, anstatt der gesamten Datenbankzeile, da die Queue jeden wartenden Job im Speicher hält.

Im Options-Objekt liegt der Großteil des Mehrwerts von BullMQ:

  • attempts — die maximale Anzahl an Versuchen, bevor der Job als fehlgeschlagen gilt.
  • backoff — die Verzögerungsstrategie zwischen den Versuchen, wie zum Beispiel exponential.
  • delay — darf erst nach Ablauf dieser Millisekunden ausgeführt werden.
  • priority — eine niedrigere Zahl wird zuerst ausgeführt, wenn Jobs warten.
  • jobId — eine stabile ID, die doppelte Enqueues derselben logischen Arbeit verhindert.
  • removeOnComplete — wie viele abgeschlossene Jobs behalten werden sollen, oder true, um sie sofort zu entfernen.
  • removeOnFail — ob Fehler zur Überprüfung behalten werden sollen. Behalten Sie diese.

Jobs verarbeiten und Ergebnisse zurückgeben

Ein Handler empfängt den Job, führt die Arbeit aus und kann einen Wert zurückgeben. Dieser Rückgabewert wird im Job gespeichert und kann später ausgelesen werden, wodurch die Queue zu einem einfachen asynchronen RPC wird.

const worker = new Worker(
  "reports",
  async (job) => {
    const pdf = await renderPdf(job.data.reportId);
    return { url: pdf.url, bytes: pdf.bytes };
  },
  { connection },
);

Ein Aufrufer kann dann den Job auf sein Ergebnis prüfen (Polling) oder über den Event-Stream auf den Abschluss warten.

const job = await reports.getJob(jobId);

if (await job?.isCompleted()) {
  return job.returnvalue;
}

if (await job?.isFailed()) {
  throw new Error(job.failedReason);
}

Geben Sie nur kleine Werte zurück. Ein Ergebnis wird in Redis wie jeder andere Datensatz gespeichert; das Zurückgeben eines Megabyte-großen PDFs würde die Queue unnötig aufblähen. Geben Sie stattdessen eine URL oder eine ID zurück und lassen Sie den Aufrufer die Daten aus dem Object Storage abrufen.

Retries, Exponential Backoff und das Failed Set

Vorübergehende Fehler sind normal. Eine Datenbank führt ein Failover durch, eine API limitiert deine Anfragen (Rate-Limiting), ein Container wird neu geplant. In solchen Fällen ist ein Retry die richtige Reaktion, aber ein sofortiger Retry nicht.

await queue.add("sync", { accountId }, {
  attempts: 5,
  backoff: { type: "exponential", delay: 2_000 },
});

Dies erzeugt Verzögerungen von etwa 2s, 4s, 8s und 16s, wobei BullMQ einen eigenen Jitter anwendet, damit viele Jobs, die gleichzeitig fehlgeschlagen sind, nicht gleichzeitig erneut versucht werden. Ohne Jitter würde eine Flotte von Workern, die alle auf denselben Ausfall stoßen, genau in dem Moment, in dem das System wieder verfügbar ist, eine neue Lastspitze verursachen.

Wenn ein Job alle seine Versuche aufgebraucht hat, löscht BullMQ ihn nicht. Er wird zusammen mit der Fehlermeldung und dem Stack Trace in das failed set verschoben. Dieses Set ist die operative Schnittstelle der Queue: Richte Alarme ein, wenn es anwächst, und baue einen Pfad für den Replay auf.

const failed = await queue.getFailed(0, 20);

for (const job of failed) {
  // Fix the underlying cause first, then requeue.
  await job.retry();
}

Ein Job, der jedes Mal fehlschlägt, ist eine sogenannte Poison Message. Ihn ewig zu wiederholen, ist schlimmer, als ihn gar nicht zu wiederholen, da er bei jedem Versuch einen Worker-Slot belegt und gesunde Aufgaben blockiert. Begrenze die Anzahl der Versuche und analysiere das failed set, anstatt es zu ignorieren.

Verzögerte und wiederholbare Jobs

Ein verzögerter Job wird ohne separaten Scheduler für einen späteren Zeitpunkt geplant. Ein wiederholbarer Job läuft nach einem Cron-Muster und wird von der Queue verwaltet.

await queue.add("reminder", { userId }, { delay: 600_000 });

await queue.add(
  "digest",
  { region: "eu" },
  {
    repeat: { pattern: "0 7 * * *", tz: "Europe/Berlin" },
    jobId: "digest:eu",
  },
);

Die stabile jobId eines wiederholbaren Jobs ist entscheidend. Sie verhindert, dass der Scheduler eine neue Kopie stapelt, wenn noch eine Instanz läuft, und ermöglicht es jeder App-Instanz, denselben Zeitplan zu registrieren, ohne Duplikate zu erstellen. Verwenden Sie für das Muster vorzugsweise UTC oder eine explizite Zeitzone: Ein Digest, der um 07:00 UTC versendet wird, ist nicht dasselbe wie einer um 07:00 Ortszeit – und dieser Unterschied führt bei jeder Zeitumstellung zu einem Support-Ticket.

Verzögerte Jobs werden in einem Sorted Set gespeichert, das über die Ausführungszeit (run-at time) indiziert ist, sodass ein verzögerter Job keinen Worker belegt. Er verbleibt in Redis, bis der Zeitpunkt erreicht ist, was Verzögerungen von Stunden oder Tagen sehr ressourcensparend macht.

Concurrency, Rate Limiting und Backpressure

Die Concurrency eines Workers bestimmt, wie viele Jobs er gleichzeitig verarbeitet. Eine Erhöhung steigert den Durchsatz, bis der Worker keine CPU-Kapazitäten, Datenbankverbindungen oder keinen Arbeitsspeicher mehr hat – ab diesem Punkt verschlechtert sich die Situation.

const worker = new Worker("sync", handler, {
  connection,
  concurrency: 10,
  limiter: { max: 50, duration: 1_000 },
});

Das limiter begrenzt, wie viele Jobs der Worker pro Zeitfenster startet. Dies ist die erste Verteidigungslinie, um eine downstream API nicht zu überlasten. Wenn ein Anbieter 50 Anfragen pro Sekunde erlaubt, wird ein Worker mit einer Concurrency von 200 zu einem Rate-Limit führen; ein Limiter bei 50 pro Sekunde hingegen nicht.

Backpressure verhindert, dass die Queue selbst unbegrenzt anwächst. Wenn Producer Jobs schneller hinzufügen, als Worker sie abarbeiten können, wird die Queue zu einem ständig wachsenden Backlog und die Latenz steigt auf Stunden an. Überwachen Sie die Queue-Tiefe, pausieren Sie Producer ab einem bestimmten Schwellenwert und skalieren Sie Worker automatisch. Eine Queue, die nur noch wächst, ist ein Ausfall, den bisher noch niemand bemerkt hat.

Fortschritt und Events

Langlaufende Jobs sollten ihren Fortschritt melden, damit die UI einen Fortschrittsbalken anzeigen kann und ein Operator den Unterschied zwischen einem langsamen und einem hängengebliebenen Job erkennt.

const worker = new Worker("imports", async (job) => {
  const rows = await loadRows(job.data.fileId);

  for (let i = 0; i < rows.length; i += 500) {
    await insertChunk(rows.slice(i, i + 500));
    await job.updateProgress(Math.round((i / rows.length) * 100));
  }

  return { rows: rows.length };
}, { connection });

Events sind die Art und Weise, wie der Rest des Systems die Queue beobachtet. Ein Worker emittiert completed, failed, progress und stalled für die Jobs, die er ausführt. QueueEvents hört von außerhalb des Workers auf denselben Stream; so kann ein API-Prozess auf den Abschluss eines Jobs reagieren, ohne selbst derjenige zu sein, der ihn ausgeführt hat.

worker.on("progress", (job, progress) => {
  console.log(`job ${job.id} at ${progress}%`);
});

worker.on("completed", (job) => {
  console.log(`job ${job.id} finished`);
});

Graceful Shutdown und stalled Jobs

Ein Worker, der mitten in einem Job beendet wird, hinterlässt diesen in einem mehrdeutigen Zustand. BullMQ löst dies mit einem Lock: Solange ein Job aktiv ist, erneuert der Worker einen Lock in Redis. Wenn der Worker abstürzt und der Lock abläuft, wird der Job als stalled markiert, zurück in die Warteschlange verschoben und erneut aufgegriffen. Dies ist einer der Gründe, warum die Zustellung nach dem Prinzip „at-least-once“ erfolgt.

Da ein Stall zu einem erneuten Durchlauf führt, sollten Sie SIGTERM bewusst handhaben, damit laufende Jobs beendet werden, anstatt unterbrochen zu werden.

process.on("SIGTERM", async () => {
  await worker.close();   // stop accepting, wait for in-flight jobs
  await connection.quit();
  process.exit(0);
});

Geben Sie dem Deployment eine Grace Period, die lang genug für den langsamsten Job ist, und begrenzen Sie die Job-Timeouts, sodass kein einzelner Job länger läuft als diese Zeitspanne. Ein Handler, der eine Stunde lang laufen kann, sollte seinen Fortschritt per Checkpoint speichern, damit ein erneuter Durchlauf dort fortgesetzt wird, wo er unterbrochen wurde, anstatt ganz von vorne zu beginnen.

Redis Streams direkt nutzen

BullMQ ist eine Schicht über den Redis-Primitiven, und manchmal möchte man direkt mit diesen Primitiven arbeiten. Streams sind das richtige Werkzeug, wenn mehrere Consumer jede Nachricht sehen müssen, wenn Sie den Verlauf erneut abspielen müssen oder wenn Sie eine Bestätigung (Acknowledgement) ohne ein Job-Framework benötigen.

XADD jobs:emails '*' type welcome userId 42

XREADGROUP GROUP workers alice COUNT 10 BLOCK 5000 STREAMS jobs:emails '>'

XACK jobs:emails workers 1760000000000-0

Die Consumer-Gruppe verfolgt, welche Nachrichten zugestellt und welche bestätigt wurden. Eine Nachricht, die zugestellt, aber nie bestätigt wurde, bleibt in der Pending-Liste der Gruppe, sodass ein abgestürzter Consumer sie nicht verliert. XAUTOCLAIM weist ausstehende Nachrichten eines toten Consumers einem aktiven zu.

Nutzen Sie Streams direkt, wenn die Daten ein Event-Log sind und die Consumer unabhängige Leser darstellen. Nutzen Sie BullMQ, wenn die Daten eine Arbeitseinheit mit einer Retry-Policy, einer Priorität und einem Ergebnis sind. Die Neuerstellung von Retries, Delays und einem Failed-Set auf Basis von Streams ist genau die Arbeit, die BullMQ bereits erledigt hat.

Idempotenz: Delivery ist at-least-once

Die wichtigste Tatsache über Queues ist, dass die Zustellung at-least-once erfolgt, nicht exactly once. Ein Worker kann abstürzen, nachdem die Arbeit erledigt wurde, aber bevor sie bestätigt wurde – in diesem Fall wird der Job erneut ausgeführt. Ein Stillstand führt systembedingt zu einem erneuten Durchlauf. Exactly-once delivery über ein Netzwerk ist praktisch unmöglich, daher entscheiden sich Queues für at-least-once und übertragen die Verantwortung an dich.

Dein Handler muss idempotent sein: Ein zweifacher Durchlauf muss denselben Endzustand erzeugen wie ein einzelner Durchlauf. Die gängige Technik hierfür ist ein Deduplizierungs-Key, der aus der Aufgabe abgeleitet und atomar vor dem Side Effect geschrieben wird.

export async function handleCharge(job) {
  const key = `charged:${job.data.orderId}`;
  const first = await connection.set(key, "1", "NX", "EX", 86_400);

  if (first === null) return { skipped: true, reason: "already_processed" };

  await stripe.charges.create(
    { amount: job.data.amount, source: job.data.token },
    { idempotencyKey: job.data.orderId },
  );
}

Beachte, dass der Key von der Order-ID stammt und nicht von der Job-ID. Das macht den Handler sicher, selbst wenn der Producer dieselbe logische Aufgabe zweimal in die Queue schreibt. Kombiniere den Deduplizierungs-Key mit dem eigenen Idempotency-Key des Providers, da dieser Key nur den API-Aufruf schützt, nicht aber die drumherum liegende Logik.

Die Queue überwachen

Eine Queue ist unsichtbar, sofern man sie nicht sichtbar macht. Vier Signale decken die meisten Anforderungen ab.

  • Queue depth — wie viele Jobs warten. Eine steigende Depth bedeutet, dass die Worker nicht hinterherkommen.
  • Oldest waiting job — die Depth sagt aus, wie viele es sind, das Alter sagt aus, wie schlimm es ist. Zehntausend Jobs, die innerhalb einer Sekunde abgearbeitet werden, sind kein Problem; zehn Jobs, die seit einer Stunde warten, hingegen schon.
  • Failure rate — Fehler pro Minute, aufgeschlüsselt nach Job-Name. Ein Spike nach einem Deploy deutet auf die Änderung hin.
  • Job duration — ein Histogramm pro Job-Name. Ein steigender p95-Wert bedeutet, dass eine Dependency langsamer wird.

BullMQ stellt diese Zähler direkt bereit, und ein QueueEvents Listener oder ein Exporter kann diese an Ihr Metrics-System weiterleiten.

const counts = await queue.getJobCounts(
  "wait", "active", "completed", "failed", "delayed",
);

console.log(counts);

Für eine visuelle Ansicht setzen Bull Board oder das BullMQ Dashboard ein kleines Web-UI über dieselben Keys auf. Fügen Sie jedem Payload eine Correlation ID hinzu und nehmen Sie diese in die Logs auf, damit ein Job von der Request, die ihn erstellt hat, über jeden Retry hinweg zurückverfolgt werden kann.

Payloads klein und versioniert halten

Eine Queue ist eine persistente Schnittstelle zwischen zwei Deployments. Ein Producer, der Version 1 des Codes ausführt, kann einen Job schreiben, den ein Worker mit Version 2 lesen muss. Dies ist dasselbe Kompatibilitätsproblem wie bei einer API, und es wird oft ignoriert, bis ein Deployment einen Backlog beschädigt.

Zwei Gewohnheiten helfen dabei, Payloads kompatibel zu halten. Erstens: Füge Felder hinzu, anstatt sie umzubenennen oder zu entfernen, und gib neuen Feldern im Handler einen sinnvollen Standardwert. Ein Worker, der ein fehlendes locale toleriert, kann Jobs verarbeiten, die bereits vor der Einführung des Feldes in die Queue eingereiht wurden. Zweitens: Füge eine Version in den Payload ein, wenn sich die Struktur signifikant ändern kann, und verzweige im Handler entsprechend.

await queue.add("import", { version: 2, importId, mapping });

Halte Payloads aus einem zweiten Grund klein: Redis hält jeden wartenden Job im Arbeitsspeicher. Ein Payload, der eine komplette Datenbankzeile einbettet, multipliziert sich über tausende von Jobs und macht die Queue teuer. Referenziere Daten per ID und lass den Worker diese abrufen. Die einzige Ausnahme ist ein Wert, der zum Zeitpunkt des Enqueueing eingefroren werden muss – wie etwa der Empfänger einer E-Mail oder der einem Kunden angebotene Preis. Diese gehören genau deshalb in den Payload, weil sie sich nicht ändern dürfen.

Arbeitsabläufe mit Flows strukturieren

Einige Jobs bestehen in Wahrheit aus mehreren Einzelschritten. Ein Bericht könnte beispielsweise Daten abrufen, ein PDF rendern und dieses per E-Mail versenden – wobei Sie vielleicht möchten, dass jeder dieser Schritte unabhängig voneinander erneut versucht wird. BullMQ flows bilden dies als einen Baum aus Parent- und Child-Jobs ab.

import { FlowProducer } from "bullmq";

const flow = new FlowProducer({ connection });

await flow.add({
  name: "report",
  queueName: "reports",
  data: { reportId },
  children: [
    { name: "fetch", queueName: "reports", data: { reportId } },
    { name: "render", queueName: "reports", data: { reportId } },
  ],
});

Ein Child-Job wird zuerst ausgeführt; der Parent-Job wird erst dann ausführbar, wenn alle seine Child-Jobs abgeschlossen sind. Dies ermöglicht Fan-out und Fan-in, ohne dass Sie den Status in Ihrer eigenen Datenbank koordinieren müssen. Schlägt ein Child-Job fehl, wartet der Parent-Job oder schlägt ebenfalls fehl, je nach den Flow-Optionen. So bleibt die Retry-Policy an dem Schritt gebunden, der tatsächlich den Fehler verursacht hat. Halten Sie Flows flach und explizit; ein tiefer Baum aus interdependenten Jobs ist schwerer nachzuvollziehen als eine kleine Pipeline mit klar definierten Phasen.

Redis Queues vs. Kafka und RabbitMQ: Die richtige Wahl

Redis ist der richtige Standard, wenn die Arbeitseinheit ein Job ist: eine benannte Aufgabe mit einem Payload, einer Retry-Policy und einem Ergebnis. Es ist schnell, vertraut und in den meisten Stacks bereits vorhanden.

Kafka ist die richtige Wahl, wenn die Daten ein Log sind. Wenn mehrere unabhängige Consumer jedes Event lesen müssen, wenn Sie den Verlauf ab einem bestimmten Offset erneut abspielen müssen oder wenn der Durchsatz in Millionen von Nachrichten pro Sekunde gemessen wird, ist ein partitioniertes Log besser geeignet als eine Job-Queue.

RabbitMQ ist die richtige Wahl, wenn das Routing die größte Herausforderung darstellt. Exchanges und Binding Keys ermöglichen es, eine Nachricht über Muster an viele Queues zu verteilen, wobei pro-Nachricht-Acknowledgements, Prioritäten und Dead-Letter-Exchanges als First-Class-Features zur Verfügung stehen. Der Betrieb ist aufwendiger als bei Redis, was sich für Teams auszahlt, die diese Flexibilität benötigen.

Die ehrlichste Empfehlung lautet: Starten Sie mit dem, was Sie bereits im Einsatz haben. Eine Redis-Queue ist weitaus besser als gar keine Queue, nur weil man noch auf die Evaluierung von Brokern gewartet hat. Wechseln Sie erst, wenn eine spezifische Einschränkung – sei es Durability, Replay oder Routing – tatsächlich zum Problem wird.

Best Practices

  • Halten Sie Request-Handler auf einen Schreibvorgang und ein Enqueue beschränkt und geben Sie 202 Accepted zurück.
  • Nutzen Sie pro Prozess eine einzige Redis-Verbindung mittels maxRetriesPerRequest: null.
  • Leiten Sie jobId aus der Arbeit ab, sodass doppelte Enqueues zu einem einzigen Job zusammengefasst werden.
  • Gestalten Sie jeden Handler idempotent, da die Zustellung nach dem At-Least-Once-Prinzip erfolgt.
  • Verwenden Sie Exponential Backoff mit Jitter und legen Sie ein Limit für die Versuche fest.
  • Behalten Sie Fehler durch das Setzen von removeOnFail: false bei und begrenzen Sie abgeschlossene Jobs mit removeOnComplete.
  • Trennen Sie Queues nach Workload, damit langsame Jobs dringende Aufgaben nicht blockieren.
  • Begrenzen Sie die Concurrency und fügen Sie für jede rate-limited Dependency einen Limiter hinzu.
  • Melden Sie den Fortschritt von langen Jobs, damit man langsame und hängengebliebene Jobs unterscheiden kann.
  • Fahren Sie Worker bei SIGTERM herunter und geben Sie Deployments eine entsprechende Grace Period.
  • Überwachen Sie die Queue-Tiefe, das Alter des ältesten Jobs, die Fehlerrate und die Dauer auf einem Dashboard.

Häufige Fehler

  • Davon ausgehen, dass ein Job genau einmal ausgeführt wird, und dadurch einen Kunden doppelt belasten.
  • Eine zufällige Job-ID verwenden und so doppelte, wiederholbare Zeitpläne stapeln.
  • Eine “Poison Message” endlos wiederholen und so die Queue blockieren.
  • Rechenintensive Aufgaben direkt im API-Prozess auszuführen und dies als Queue zu bezeichnen.
  • Die Concurrency so hoch anzusetzen, dass die Datenbank ihr Connection-Limit erreicht.
  • Ein riesiges Objekt als Job-Ergebnis zurückgeben und so Redis aufblähen.
  • maxRetriesPerRequest: null zu vergessen und dadurch Jobs bei einer Wiederverbindung zu verlieren.
  • Das Failed-Set auf jedem Dashboard auszublenden, bis es zum ersten Vorfall kommt.
  • Worker mit SIGKILL zu beenden und so jeden laufenden Job zum Stillstand zu bringen.
  • Eine komplette Datenbankzeile in den Payload zu packen, anstatt nur eine ID.

Wie geht es weiter?

Der Redis-Guide behandelt den Speicher unter der Queue: Lists, Streams, TTLs, Persistenz und das Single-Threaded-Command-Modell, das atomare Claims ermöglicht. Um mehr über das allgemeine Konzept zu erfahren, Arbeit aus dem Request-Pfad auszulagern, lesen Sie Batch Processing. Wenn Routing und Bestätigungen pro Nachricht (Acknowledgements) komplexer werden, ist RabbitMQ der nächste Schritt. Und wenn Sie ein Replay-fähiges Log anstelle einer Job-Queue benötigen, ist der Kafka-Guide die richtige Adresse.

In der Praxis

Die vier Dateien eines BullMQ-Services

Verbindung, Producer, Worker und ein wiederholbarer Zeitplan. Zusammen decken sie den gesamten Lebenszyklus ab.

jobs/queue.ts
import { Queue } from "bullmq";
import IORedis from "ioredis";

export const connection = new IORedis(process.env.REDIS_URL!, {
  maxRetriesPerRequest: null,
});

export const emails = new Queue("emails", { connection });

export async function enqueueWelcome(userId: string, to: string) {
  await emails.add(
    "welcome",
    { userId, to, template: "welcome" },
    {
      jobId: `welcome:${userId}`,
      attempts: 5,
      backoff: { type: "exponential", delay: 2_000 },
      removeOnComplete: { count: 1_000 },
      removeOnFail: false,
    },
  );
}

BullMQ vs. reine Redis Lists

Eine Liste reicht für Fire-and-Forget-Aufgaben. Sobald man Retries, Verzögerungen oder Sichtbarkeit benötigt, lohnt sich eine Library.

Bevorzugt
import { Queue, Worker } from "bullmq";

const queue = new Queue("imports", { connection });
await queue.add("import", { importId }, {
  attempts: 5,
  backoff: { type: "exponential", delay: 2_000 },
});

new Worker("imports", async (job) => {
  await runImport(job.data.importId);
}, { connection, concurrency: 5 });
Vermeiden
// A bare list has no retries, no delayed jobs,
// no attempt tracking and no way to inspect failures.
await redis.lpush("imports", JSON.stringify({ importId }));

while (true) {
  const [, raw] = await redis.brpop("imports", 0);
  await runImport(JSON.parse(raw).importId);
}

Backoff vs. sofortiger Retry

Sofortige Retries belasten eine Abhängigkeit, die bereits Probleme hat. Exponential Backoff mit Jitter gibt ihr Raum zur Erholung.

Bevorzugt
await queue.add("sync", { accountId }, {
  attempts: 5,
  backoff: { type: "exponential", delay: 2_000 },
});
// ~2s, 4s, 8s, 16s, plus jitter
Vermeiden
// A tight loop makes the outage worse and
// burns every attempt in a few milliseconds.
for (let i = 0; i < 5; i++) {
  try { await sync(accountId); break; }
  catch { /* retry at once */ }
}

Abwägungen

Ist Redis die richtige Queue für dich?

Redis ist der kürzeste Weg von einem bestehenden Cache zu einer funktionierenden Job-Queue. Es ist jedoch nicht das beste Tool für jede Art von Arbeit.

Strengths

  • Bereits im Stack

    Wenn Redis deinen Cache oder deine Sessions unterstützt, fügst du eine Queue hinzu, ohne neue Infrastruktur, neue Credentials oder ein neues On-Call-Runbook zu benötigen.

  • Exzellente Ergonomie

    BullMQ liefert Verzögerungen, wiederholbare Jobs, Prioritäten, Rate Limits, Fortschritt, Flows und ein UI mit. Du schreibst Handler, kein Queue-Plumbing.

  • Schnell und einfach zu betreiben

    Ein Prozess, ein kleiner Befehlssatz und vorhersehbare Latenz machen eine Redis Queue leicht nachvollziehbar und günstig im Betrieb.

Trade-offs

  • Durability ist konfigurierbar, nicht garantiert

    Redis ist primär In-Memory. Ein Absturz zwischen Bestätigung und Festplatte kann einen Job kosten, es sei denn, du optimierst AOF und akzeptierst die Schreibkosten.

  • Nicht für Replay oder Fan-out gebaut

    Eine Redis Queue entfernt einen Job, sobald er erledigt ist. Wenn viele Consumer jede Nachricht sehen müssen oder du den Verlauf wiederholen musst, ist ein Log besser geeignet.

  • Eine 'Hot Queue' kann blockieren

    Ein 'Poison Job' und eine niedrige Concurrency-Einstellung können alles dahinter verzögern. Trenne Queues nach Workload und limitiere die Versuche.

Häufig gestellte Fragen

Häufig gestellte Fragen

Keep learning

Related topics from the roadmap.

$ Lernen Sie jetzt

Bereit, Redis Queues zu lernen?

Unser interaktives Tutorial führt Sie Schritt für Schritt durch Redis Queues — mit Quizzen und echtem Code, den Sie im Browser ausführen können.