Was ist Batch-Processing?
Batch-Processing ist die Praxis, Aufgaben, die zu langsam, zu fehleranfällig oder zu schwankend in ihrer Auslastung sind, um sie während des Wartens eines Benutzers auszuführen, an einen separaten Prozess zu übergeben, der diese nach eigenem Zeitplan abarbeitet. Die Anfrage des Benutzers erledigt nur das Nötigste – Eingaben validieren, eine Zeile schreiben, einen Job in die Queue stellen – und kehrt dann zurück. Alles Weitere geschieht „out of band“.
Der Name stammt aus der Mainframe-Ära, in der Stapel von Datensätzen in einem einzigen Durchlauf verarbeitet wurden, und das Grundkonzept hat sich nicht geändert. Anstelle eines Magnetbands gibt es heute eine Queue. Anstelle eines nächtlichen Zeitfensters gibt es Worker, die kontinuierlich Daten abrufen. Was gleich geblieben ist, ist die Trennung: Die Komponente, die die Arbeit annimmt, und die Komponente, die sie ausführt, sind verschieden, und dazwischen liegt ein Buffer.
Diese Trennung ist der entscheidende Punkt. Sie sorgt dafür, dass der Request-Pfad schnell und vorhersagbar bleibt, während der langsame Pfad so viel Zeit in Anspruch nehmen kann, wie nötig ist, bei Fehlern Retries durchführt und nach einer anderen Kurve skaliert.
Warum rechenintensive Aufgaben nicht in einen Request gehören
Ein synchroner HTTP-Request ist aus drei Gründen, die sich gegenseitig verstärken, ein schlechter Ort für langsame Prozesse.
Erstens: Timeouts. Proxies, Load Balancer und Clients setzen alle Zeitlimits. Ein Request, der ein 200-seitiges PDF generiert, eine instabile Drittanbieter-API aufruft und eine E-Mail versendet, kann diese leicht überschreiten. Wenn das passiert, sieht der Client einen Fehler, obwohl die Arbeit auf dem Server möglicherweise bereits zur Hälfte abgeschlossen wurde.
Zweitens: Ein blockierter Event Loop. Node führt JavaScript in einem einzigen Thread aus. Eine synchrone, CPU-intensive Aufgabe – wie Bildverarbeitung, das Parsen eines riesigen JSON-Objekts oder eine kryptografische Schleife – verhindert, dass jeder andere Request bedient wird, solange sie läuft. Selbst asynchrone Arbeit bindet Ressourcen: eine offene Datenbankverbindung, ein File-Handle oder Speicher für die Antwort.
Drittens: Retries sind unmöglich. Wenn der Mail-Provider einen 503-Fehler zurückgibt, hat ein Inline-Handler keine guten Optionen. Er kann entweder den gesamten Request fehlschlagen lassen und den Benutzer zum erneuten Versuch zwingen, oder den Fehler verschlucken und die E-Mail verlieren. Eine Queue bietet für denselben Fehler eine dritte Lösung: Versuche es später automatisch erneut, ohne dass der Benutzer davon erfährt.
app.post("/reports", async (req, res) => {
const report = await db.report.create({ data: { userId: req.user.id } });
await reportsQueue.add("generate", { reportId: report.id });
res.status(202).json({ id: report.id, status: "queued" });
});
Der Status 202 Accepted ist hier die ehrlichste Antwort: Der Server hat den Request akzeptiert, aber die Arbeit noch nicht abgeschlossen. Dies ist das Standardmuster für jeden gut implementierten Batch-Endpoint.
Die Anatomie eines Jobs
Ein Job ist ein kleiner, selbsterklärender Datensatz. Die meisten Queues speichern etwas in dieser Form:
{
"id": "welcome:user_42",
"name": "welcome",
"data": { "userId": "user_42", "to": "[email protected]", "template": "welcome" },
"status": "queued",
"attemptsMade": 0,
"maxAttempts": 5,
"runAt": 1760000000000,
"createdAt": 1759999100000
}
Jedes Feld hat seine Berechtigung. Die id macht den Job adressierbar und ermöglicht – sofern sie aus der Arbeit selbst und nicht zufällig generiert wird – eine automatische Deduplizierung. Der name routet den Job an einen Handler. Der payload enthält alles, was der Worker benötigt – und nichts mehr, denn ein Payload, der auf eine Datenbankzeile verweist, ist kleiner und aktueller als einer, der diese kopiert. Der status verfolgt den Job durch seinen Lebenszyklus. attemptsMade und maxAttempts steuern die Retries. runAt plant verzögerte und wiederkehrende Aufgaben.
Der Payload sollte ein Snapshot der Absicht sein, kein Live-Objekt. Wenn ein Nutzer seine E-Mail-Adresse zwischen dem Enqueue und der Ausführung aktualisiert, sollte der Job immer noch an die Adresse gesendet werden, mit der er erstellt wurde. Es ist jedoch ein Fehler, den gesamten User-Datensatz zu speichern: Das bläht die Queue auf und die Daten veralten. Speichern Sie IDs und nur die wenigen Werte, die die eigentliche Arbeit definieren.
Job-Payloads und Versionierung
Eine Queue ist eine persistente Schnittstelle zwischen zwei Deploys. 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 Deploy einen Backlog zerschießt.
Zwei Gewohnheiten halten Payloads kompatibel. 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 maßgeblich ändern kann:
await queue.add("import", { version: 2, importId, mapping });
Der Handler verzweigt dann basierend auf version und weiß genau, wie der Rest zu interpretieren ist. Das kostet fast nichts und verwandelt eine ganze Klasse von Post-Deploy-Incidents in eine einfache Lookup-Tabelle.
Halte Payloads klein. Eine Queue speichert jeden wartenden Job; ein Payload, der ein großes Objekt einbettet, vervielfacht sich über tausende Zeilen und verlangsamt jeden Scan. Referenziere die Daten per ID und lass den Worker sie abrufen. Die einzige Ausnahme ist ein Wert, der zum Zeitpunkt des Enqueue-Vorgangs eingefroren werden muss – etwa der Empfänger einer E-Mail oder der einem Kunden angebotene Preis –, welcher genau deshalb in den Payload gehört, weil er sich nicht ändern darf.
Queues, Cron und eventgesteuerte Trigger
Nicht jede Hintergrundaufgabe benötigt eine Queue, und die Wahl des falschen Triggers macht ein einfaches Problem unnötig kompliziert.
Cron führt einen Handler nach einem festen Zeitplan aus: jede Nacht um 02:00 Uhr, jeden Montag um 09:00 Uhr. Es ist das richtige Werkzeug für periodische Abstimmungen, die Generierung von Berichten und Cleanup-Aufgaben. Die Schwäche ist, dass Cron kein Konzept von Arbeitseinheiten besitzt. Ein Cron-Job, der länger dauert als sein Intervall, wird sich mit sich selbst überschneiden, und ein verpasster Durchlauf ist schlichtweg verpasst.
Event-driven triggers reagieren auf Ereignisse, die gerade eingetreten sind: ein Webhook ist eingegangen, eine Zeile wurde eingefügt, eine Datei wurde im Storage abgelegt. Sie sind unmittelbar und intuitiv, bieten jedoch keine Pufferung, keine Retries und keinen Backpressure. Ein Webhook-Handler, der tatsächlich rechenintensive Arbeit verrichtet, ist im Grunde nur ein getarnter Inline-Request.
Queues liegen genau dazwischen. Ein Job ist eine persistente, wiederholbare Arbeitseinheit, die für sofort oder für einen späteren Zeitpunkt geplant werden kann. Die meisten Produktionssysteme nutzen alle drei Ansätze: Cron stellt einen Fan-out-Job in die Queue, Events stellen Jobs als Reaktion auf Benutzeraktionen in die Queue und Worker leeren die Queue. Die Faustregel lautet: Cron und Events entscheiden, wann gearbeitet wird, und die Queue entscheidet, wie gearbeitet wird.
Producer, Worker und die Queue
Das System besteht aus drei Rollen. Diese strikt getrennt zu halten, sorgt für eine gute Wartbarkeit.
Der Producer ist jeder Code, der einen Job hinzufügt. Er kennt lediglich den Job-Namen und die Struktur des Payloads, mehr nicht. Der Aufruf sollte schnell erfolgen und idempotent sein – falls der Producer selbst nach einem Timeout einen Retry ausführt, sollen nicht zwei Jobs erstellt werden.
Die Queue ist ein persistenter, geordneter Speicher. Redis mit BullMQ, Amazon SQS, RabbitMQ und Google Cloud Tasks erfüllen diese Rolle. Die Queue speichert Jobs über Neustarts hinweg, gibt sie atomar aus, verfolgt die Versuche und verschiebt fehlgeschlagene Jobs in eine separate Liste. Man kann eine Queue auch auf Basis einer Datenbanktabelle aufbauen, was bei geringem Volumen ein vernünftiger Start ist, sobald das Volumen jedoch steigt, wird dies zum Problem.
Der Worker ist ein langlebiger Prozess, der Jobs abruft und die entsprechenden Handler ausführt. Er ist aus guten Gründen von der API getrennt: Er kann auf CPU-optimierter Hardware bereitgestellt, basierend auf der Queue-Tiefe skaliert und neu gestartet werden, ohne dass Anfragen verloren gehen. In BullMQ ist diese Trennung explizit:
import { Worker } from "bullmq";
import { connection } from "./queue.js";
const worker = new Worker(
"emails",
async (job) => {
await sendEmail(job.data);
},
{ connection, concurrency: 10 },
);
Ein Prozess kann mehrere Worker für verschiedene Queues hosten, und eine Queue kann von vielen Worker-Prozessen bedient werden. Die Queue ist der einzige gemeinsame Zustand, weshalb das System so sauber horizontal skaliert.
Den richtigen Broker auswählen
Die drei Broker, denen Sie am häufigsten begegnen werden, bewegen sich in einem Spektrum zwischen Funktionsumfang und operativem Aufwand.
Redis mit BullMQ ist der Standard für Node-Teams. Redis ist wahrscheinlich bereits Teil Ihres Stacks, der Client ist ausgereift und BullMQ ergänzt dies um verzögerte Jobs, wiederholbare Jobs, Prioritäten, Rate Limiting, Retries und ein UI. Der Haken ist, dass Redis primär ein In-Memory-Store ist; die Persistenz hängt also davon ab, wie Sie die Konfiguration vornehmen. Ein bestätigter Job, der noch nicht auf die Festplatte geschrieben wurde, kann verloren gehen, wenn die Instanz abstürzt.
Amazon SQS ist vollständig verwaltet und bietet einen faktisch unbegrenzten Durchsatz. Sie erhalten Persistenz und Verfügbarkeit, ohne selbst etwas betreiben zu müssen, und zahlen pro Request. Im Gegenzug verzichten Sie auf einige ergonomische Vorteile: Die verzögerte Zustellung ist begrenzt, es gibt keinen integrierten Scheduler für Cron-ähnliche Jobs und die API ist niedrigschwelliger als ein dediziertes Job-Framework.
RabbitMQ ist die flexibelste Option. Exchanges und Routing Keys ermöglichen es, eine Nachricht an viele Consumer zu verteilen (Fan-out), und Bestätigungen pro Nachricht, Prioritäten sowie Dead-Letter-Exchanges sind First-Class-Features. Der Betrieb ist aufwendiger als bei Redis und die Lernkurve steiler, aber für komplexes Routing ist es die richtige Wahl.
Die ehrlichste Empfehlung lautet: Beginnen Sie mit dem, was Sie bereits im Einsatz haben. Eine Queue auf Redis ist weitaus besser als gar keine Queue, nur weil Sie noch auf die Evaluierung verschiedener Broker gewartet haben. Wechseln Sie erst, wenn eine spezifische Einschränkung – sei es bei der Persistenz, dem Durchsatz oder dem Routing – tatsächlich zum Problem wird.
Idempotenz und At-Least-Once-Delivery
Die wichtigste Erkenntnis über Job-Queues ist, dass die Zustellung nach dem Prinzip At-Least-Once (mindestens einmal) und nicht Exactly-Once (genau einmal) erfolgt. Ein Worker kann abstürzen, nachdem er die Arbeit erledigt, aber bevor er diese bestätigt hat – in diesem Fall wird die Queue den Job erneut zustellen. Ein kurzer Netzwerkfehler kann dazu führen, dass eine Bestätigung verloren geht. Die Arbeit wird erneut ausgeführt.
Dies ist kein Bug, den man umgehen muss, sondern Teil des Vertrags. Eine Exactly-Once-Zustellung über ein Netzwerk hinweg 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 Aufruf muss denselben Endzustand erzeugen wie ein einzelner Aufruf.
Es gibt drei praktische Techniken.
Natürliche Idempotenz. Einige Operationen sind bereits von Natur aus sicher zu wiederholen. Den Status eines Benutzers zweimal auf active zu setzen, ist dasselbe wie ein einziges Mal. Das Löschen einer Zeile per ID ist beim zweiten Mal eine No-Op. Bevorzuge diese Ansätze, wann immer möglich.
Ein Deduplizierungs-Key. Schreibe vor der Ausführung der Arbeit einen Marker, der über die Arbeit indiziert ist, und überspringe den Vorgang, falls der Marker bereits existiert. Ein Unique Constraint oder ein Redis SET NX macht die Prüfung atomar, sodass zwei gleichzeitige Worker nicht beide erfolgreich sein können.
export async function handleCharge(job) {
const key = `charged:${job.data.orderId}`;
const inserted = await connection.set(key, "1", "NX", "EX", 86_400);
if (inserted === null) return { skipped: true };
await stripe.charges.create(
{ amount: job.data.amount, source: job.data.token },
{ idempotencyKey: job.data.orderId },
);
}
Idempotenz-Keys des Providers. Payment-Gateways und viele andere APIs akzeptieren einen Idempotenz-Key. Übergib die stabile ID des Jobs, und der Provider wird das ursprüngliche Ergebnis zurückgeben, anstatt den Betrag doppelt abzubuchen. Kombiniere dies immer mit deiner eigenen Deduplizierung, da der Key nur den API-Aufruf schützt, nicht aber die drumherum liegende Logik.
Beachte, dass der Deduplizierungs-Key von der eigentlichen Arbeit – etwa der Order-ID – abgeleitet wird und nicht von der Job-ID. Das ist beabsichtigt: So bleibt der Handler sicher, selbst wenn der Producer dieselbe logische Arbeit zweimal in die Queue schreibt.
Retries mit Exponential Backoff und Jitter
Vorübergehende Fehler sind normal. Eine Datenbank führt ein Failover durch, eine API limitiert Ihre Anfragen (Rate-Limiting), ein Container wird neu geplant. Ein Retry ist hier die richtige Reaktion, aber ein sofortiger Retry nicht.
Exponential Backoff erhöht die Verzögerung mit jedem Versuch: etwa 2s, 4s, 8s, 16s, 32s. Dies gibt einer überlasteten Dependency Zeit, sich zu erholen, anstatt sie weiter zu bombardieren. Jitter fügt jeder Verzögerung einen Zufallswert hinzu, damit viele gleichzeitig fehlgeschlagene Jobs nicht gleichzeitig einen Retry ausführen. Letzteres würde die Lastspitze, die den Fehler ursprünglich verursacht hat, einfach nur reproduzieren.
Die meisten Queues unterstützen dies deklarativ. In BullMQ ist es eine Eigenschaft des Jobs:
await queue.add(
"sync",
{ accountId },
{
attempts: 5,
backoff: { type: "exponential", delay: 2_000 },
},
);
Dies erzeugt Verzögerungen von etwa 2s, 4s, 8s, 16s und 32s, wobei der eigene Jitter der Queue angewendet wird. Wenn Sie Backoff manuell implementieren, fügen Sie den Jitter selbst hinzu und begrenzen die maximale Verzögerung, damit ein Job nicht einen ganzen Tag lang schläft:
function nextDelay(attempt: number, base = 1_000, cap = 60_000) {
const exponential = Math.min(cap, base * 2 ** attempt);
const jitter = Math.random() * exponential * 0.5;
return Math.round(exponential + jitter);
}
Legen Sie ein Limit für die Versuche fest und entscheiden Sie bewusst, was passiert, wenn dieses erreicht wird. Einige Jobs sollten ewig mit niedriger Frequenz wiederholt werden – zum Beispiel ein Reconciliation-Task –, aber die meisten sollten stoppen und eine Fehlermeldung ausgeben.
Dead-letter queues und poison messages
Eine poison message ist ein Job, der bei jedem Ausführen fehlschlägt: etwa durch fehlerhafte Daten, einen Bug im Handler oder eine fehlende referenzierte Zeile. Da er immer fehlschlägt, belegt er bei jedem Versuch einen Worker-Slot und kann dadurch gesunde Aufgaben blockieren. Ihn endlos zu wiederholen, ist schlimmer, als ihn gar nicht erst zu wiederholen.
Die Lösung ist eine dead-letter queue (DLQ). Sobald das Limit für die Versuche erreicht ist, verschiebt die Queue den Job – inklusive Payload, Versuchsanzahl und letztem Fehler – in einen separaten Haltebereich. Worker greifen diesen Bereich nie an, sodass nichts blockiert werden kann. Ein Operator kann den Job untersuchen, die Ursache beheben und ihn erneut ausführen (replay).
const worker = new Worker("emails", handler, { connection });
worker.on("failed", async (job, err) => {
if (job && job.attemptsMade >= (job.opts.attempts ?? 1)) {
await deadLetter.add("failed-email", {
payload: job.data,
error: err.message,
failedAt: new Date().toISOString(),
});
await alerting.notify(`job ${job.id} exhausted retries`);
}
});
Betrachten Sie die DLQ als operative Schnittstelle, nicht als Friedhof. Richten Sie Alarme ein, wenn die Anzahl der Nachrichten darin steigt, erstellen Sie ein Dashboard und bauen Sie einen Pfad für den Replay auf. Eine DLQ, die niemand liest, ist der Ort, an dem sich Bugs verstecken.
Concurrency, Rate Limiting und Backpressure
Die Concurrency eines Workers bestimmt, wie viele Jobs gleichzeitig verarbeitet werden. 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. Der richtige Wert ist die höchste Zahl, bei der jede Abhängigkeit komfortabel unter ihrem Limit bleibt.
Concurrency ist zudem die erste Verteidigungslinie gegen einen Thundering Herd-Effekt. Wenn eine nachgelagerte API nur 50 Anfragen pro Sekunde zulässt, wird ein Worker mit einer Concurrency von 200 diese überlasten. Viele Queues bieten genau für diesen Fall einen Rate Limiter an:
const worker = new Worker("sync", handler, {
connection,
concurrency: 10,
limiter: { max: 50, duration: 1_000 },
});
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. Zu den Optionen gehören das Pausieren von Producern, sobald eine bestimmte Tiefe überschritten wird, das Ablehnen von Jobs mit niedriger Priorität oder das automatische Skalieren der Worker. Eine Queue, die nur noch wächst, ist im Grunde ein Ausfall, der nur noch nicht bemerkt wurde.
Halten Sie Queues nach Workload getrennt. Ein langsamer nächtlicher Import und ein zeitkritischer Passwort-Reset sollten nicht in derselben Queue liegen.
Datenbank-Schreibvorgänge bündeln
Die Datenbank ist in einem Batch-Job meist der Flaschenhals, und die häufigste Ursache dafür ist die Kommunikation in Form von Einzelzeilen. Jeder Roundtrip verursacht einen fixen Overhead – Netzwerk, Parsing, Planning –, der die Kosten der eigentlichen Zeile bei weitem übersteigt. Eine Schleife aus Einzel-Inserts verbringt die meiste Zeit mit Warten.
Ein Multi-Row Insert überträgt dieselben Daten in einem einzigen Statement:
const values = chunk
.map((_, n) => `($${n * 3 + 1}, $${n * 3 + 2}, $${n * 3 + 3})`)
.join(",");
await pool.query(
`INSERT INTO orders (user_id, status, total_cents)
VALUES ${values}
ON CONFLICT (external_id) DO UPDATE
SET status = EXCLUDED.status`,
chunk.flatMap((r) => [r.userId, r.status, r.totalCents]),
);
Die ON CONFLICT ... DO UPDATE-Klausel macht das Insert zu einem Upsert, wodurch der Bulk-Schreibvorgang idempotent wird. Ein erneuter Durchlauf des Jobs aktualisiert bestehende Zeilen, anstatt Duplikate zu erstellen, sodass ein Retry sicher ist.
Zwei Warnhinweise: Halten Sie die Chunks begrenzt – einige hundert bis einige tausend Zeilen –, da ein parametrisiertes Statement eine Grenze für Parameter hat und ein sehr großes Statement Locks und Speicher länger belegt. Und kapseln Sie einen Chunk in eine Transaction, wenn die Zeilen gemeinsam geschrieben werden müssen, aber halten Sie die Transaction kurz, damit sie andere Writer nicht blockiert.
Für sehr große Datenmengen ist ein dedizierter Bulk-Pfad wie COPY von Postgres noch schneller, wie im PostgreSQL-Guide beschrieben.
Große Datensätze in Chunks aufteilen
Ein Batch-Job, der „alle Benutzer“ verarbeitet, kann diese nicht alle gleichzeitig in den Speicher laden. Die Lösung ist das Chunking: Verarbeite eine begrenzte Seite, führe die Änderungen aus (commit) und hole dann die nächste Seite. Dadurch bleibt der Speicherverbrauch konstant und der Job kann an der Stelle fortgesetzt werden, an der er gestoppt wurde.
Keyset pagination ist der robuste Weg, dies umzusetzen. Anstatt OFFSET zu verwenden – was mit zunehmender Datenmenge langsamer wird und bei Datenänderungen im Hintergrund Zeilen überspringen oder doppelt ausgeben kann –, merkt man sich den letzten gesehenen Key:
let cursor: string | null = null;
while (true) {
const batch = await pool.query(
`SELECT id, email FROM users
WHERE ($1::text IS NULL OR id > $1)
ORDER BY id
LIMIT 1000`,
[cursor],
);
if (batch.rowCount === 0) break;
await processBatch(batch.rows);
cursor = batch.rows[batch.rows.length - 1].id;
}
Jeder Chunk ist unabhängig. Wenn es also mitten im Durchlauf zu einem Absturz kommt, geht nur der aktuelle Chunk verloren, und der Job kann durch Übergabe des letzten Cursors wieder aufgenommen werden. Dies lässt sich ideal mit einer Queue kombinieren: Setze pro Chunk einen Job in die Queue, sodass ein einzelner Fehler nicht den Neustart des gesamten Durchlaufs erzwingt.
Wiederkehrende Jobs planen
Wiederkehrende Aufgaben – wie tägliche Digests, Cleanup-Prozesse oder Abstimmungen – lassen sich am besten als repeatable job ausdrücken, anstatt einen Cron-Eintrag zu nutzen, der einen HTTP-Endpunkt aufruft. Die Queue übernimmt dadurch die Verwaltung des Zeitplans, den Schutz vor Überlappungen und die Retry-Policy.
await reportsQueue.add(
"daily-digest",
{ region: "eu" },
{
repeat: { pattern: "0 7 * * *" },
jobId: "daily-digest:eu",
attempts: 3,
},
);
Die stabile jobId ist hierbei entscheidend: Sie verhindert, dass der Scheduler eine neue Instanz startet, während noch eine läuft, und ermöglicht es jeder App-Instanz, denselben Zeitplan zu registrieren, ohne Duplikate zu erzeugen. Ein Distributed Lock um die eigentliche Arbeit ist dennoch ratsam für Jobs, die unter keinen Umständen zweimal gleichzeitig laufen dürfen.
Verwenden Sie für Zeitpläne vorzugsweise UTC und lassen Sie den Job die Zeitzonen berücksichtigen, wenn er die Ausgabe formatiert. 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.
Observability
Ein Background Job ist unsichtbar, sofern Sie ihn nicht sichtbar machen. Vier Signale decken den Großteil dessen ab, was Sie benötigen.
- Queue depth — wie viele Jobs warten. Eine steigende Tiefe bedeutet, dass die Worker nicht hinterherkommen, was das früheste Warnsignal für einen Ausfall ist.
- Job duration — ein Histogramm pro Job-Name. Ein p95-Wert, der im Laufe der Zeit ansteigt, signalisiert, dass eine Abhängigkeit langsamer wird.
- Failure rate — Fehler pro Minute, aufgeteilt nach Job-Name. Ein Peak nach einem Deploy weist direkt auf die Änderung hin.
- Age of the oldest waiting job — die Tiefe sagt aus, wie viele es sind, das Alter sagt aus, wie schlimm es ist. Zehntausend Jobs, die in einer Sekunde abgearbeitet werden, sind kein Problem; zehn Jobs, die seit einer Stunde warten, hingegen schon.
Fügen Sie jedem Job-Payload eine correlation id hinzu und nehmen Sie diese in die Logs auf, sodass ein Job von der Request, die ihn erstellt hat, über jeden Retry hinweg zurückverfolgt werden kann. Ohne diese ID bedeutet das Debugging eines Worker-Fehlers, Zeitstempel zu greppen und zu raten.
await queue.add("import", { importId, correlationId: req.id });
Stellen Sie die Metriken im selben Dashboard wie Ihre API bereit und richten Sie Alerts für die Tiefe und das Alter des ältesten Jobs ein, anstatt für einzelne Fehler, da diese zu erwarten sind.
Graceful Shutdown
Ein Worker, der mitten in einem Job beendet wird, hinterlässt diesen in einem undefinierten Zustand. Die Queue wird den Job zwar schließlich erneut zustellen – was zwar korrekt, aber ineffizient ist –, und ein hartes Beenden kann eine Datenbanktransaktion im ungünstigsten Moment unterbrechen.
Behandeln Sie SIGTERM und beenden Sie den Worker gezielt:
process.on("SIGTERM", async () => {
await worker.close();
await connection.quit();
process.exit(0);
});
worker.close() stoppt die Annahme neuer Jobs und wartet, bis die aktuell laufenden abgeschlossen sind. Kombinieren Sie dies mit einer Deployment-Grace-Period, die lang genug für den langsamsten Job ist, und begrenzen Sie die Job-Timeouts, damit kein einzelner Job diese Zeit überschreitet. Wenn ein Job tatsächlich sehr lange läuft, implementieren Sie Checkpoints für den Fortschritt, sodass er fortgesetzt und nicht komplett neu gestartet werden muss.
Die gleiche Disziplin gilt für die Verbindungen: Schließen Sie den Datenbank-Pool und den Broker-Client, damit der Prozess sauber beendet wird, anstatt an offenen Sockets hängen zu bleiben.
Best Practices
- Halten Sie Request-Handler auf einen Schreibvorgang und ein Enqueue beschränkt; geben Sie
202zurück, wenn die Arbeit aufgeschoben wird. - Machen Sie jeden Handler idempotent, da die Zustellung nach dem Prinzip „at-least-once“ erfolgt.
- Leiten Sie Deduplizierungs-Keys aus der eigentlichen Arbeit ab, nicht aus einer zufälligen Job-ID.
- Verwenden Sie Exponential Backoff mit Jitter und begrenzen Sie die maximale Verzögerung.
- Legen Sie ein Limit für Versuche fest und leiten Sie erschöpfte Jobs an eine Dead-Letter-Queue weiter.
- Trennen Sie Queues nach Workload, damit langsame Jobs dringende nicht blockieren können.
- Begrenzen Sie die Concurrency auf das Maß, das Ihre Datenbank und die nachgelagerten APIs bewältigen können.
- Bündeln Sie Datenbank-Schreibvorgänge mit Multi-Row-Inserts oder Upserts in begrenzten Chunks.
- Pagen Sie große Datensätze mit Keyset-Pagination, nicht mit
OFFSET. - Überwachen Sie die Queue-Tiefe, die Job-Dauer, die Fehlerrate und das Alter des ältesten Jobs.
- Fahren Sie Worker kontrolliert herunter (graceful shutdown) und räumen Sie Deployments eine entsprechende Grace Period ein.
Häufige Fehler
- Langsame Third-Party-Aufrufe innerhalb des Requests auszuführen und dies als akzeptabel zu betrachten, nur weil es lokal funktioniert.
- Davon auszugehen, dass ein Job exakt einmal ausgeführt wird, und dadurch einen Kunden doppelt zu belasten.
- Sofortige Retries ohne Backoff durchzuführen und so den ursprünglichen Fehler zu verstärken.
- Eine Poison Message endlos zu wiederholen und so die Queue zu blockieren.
- Einen einzigen riesigen Job auszuführen, der jede Zeile in einer einzigen Transaction verarbeitet.
- Zeilen einzeln pro Statement einzufügen und dann der Datenbank die Schuld zu geben.
- Die Concurrency so hoch anzusetzen, dass die Datenbank ihr Connection Limit erreicht.
- Wiederkehrende Jobs mit einer zufälligen ID zu planen und so Duplikate zu stapeln.
- Die Dead-Letter Queue niemals zu prüfen.
- Worker mit
SIGKILLzu beenden und dadurch laufende Arbeit zu verlieren. - Die Queue Depth erst dann auf das Dashboard zu nehmen, wenn der erste Incident auftritt.
Wie geht es weiter?
Eine Queue ist nur so gut wie der Store, der hinter ihr steht. Daher ist der Redis-Guide die logische nächste Lektüre für den Broker, den die meisten Node-Teams einsetzen. Der Guide zu Caching zeigt, wie man Arbeit komplett vermeidet, und Connection Pooling erklärt, wie man verhindert, dass eine Flotte von Workern Ihre Datenbank überlastet. Wenn es um Bulk-Writes geht, behandelt der PostgreSQL-Guide COPY, Upserts und die Transaktionsstrukturen, die für maximale Geschwindigkeit sorgen.