PostgreSQL als Job Queue: Braucht man wirklich Redis oder RabbitMQ?
Einführung
In diesem Artikel gehen wir durch, warum man Postgres als Queue verwenden kann, sehen uns das Schema-Design an, implementieren nebenläufig sicheres Task-Polling mit SELECT … FOR UPDATE SKIP LOCKED und stellen sicher, dass zu jedem Zeitpunkt nur ein Worker einen bestimmten Task bearbeitet. Am Ende hast du die Blaupause für ein zuverlässiges, transaktional abgesichertes System zur Verarbeitung von Hintergrund-Tasks, und alles läuft in Postgres.

Inhaltsverzeichnis
- Warum Postgres als Job Queue?
- Grundbegriffe und Definitionen
- Datenbankschema-Design
- Locking in Postgres: FOR UPDATE SKIP LOCKED
- Neue Tasks einfügen
- Polling nach zu bearbeitenden Tasks
- Task-Status aktualisieren
- Parallele Worker und Nebenläufigkeit
- Wiederholungen und Fehlerbehandlung
- Vor- und Nachteile
- Eine Beispielimplementierung Schritt für Schritt
- Praktische Tipps
- Fazit und eine Frage zum Nachdenken
Im PATREON-Artikel findest du 2 Projekte:
- JobProducer.dproj - das Client-Programm, das Jobs anfordert. Im Beispiel kann es 2 Arten von Jobs anfordern (“send_email”, “create_report”), du kannst es aber auf alles erweitern, was du brauchst.
- QueueWorker.dproj - der eigentliche Worker. Du kannst mehrere Instanzen davon starten und prüfen, dass das System jeden einzelnen Job genau einer Worker-Instanz zuweist. Stürzt der Worker ab, nachdem er einen Job übernommen hat, steht dieser Job wieder für andere Worker bereit.
1. Warum Postgres als Job Queue?
Wenn es um Job Queues geht, greifen Entwickler oft zu Werkzeugen wie Redis, RabbitMQ oder spezialisierten Systemen wie Sidekiq (im Ruby-Ökosystem). Jede dieser Optionen ist völlig legitim. Es gibt aber gute Gründe, Postgres für deine Queue in Betracht zu ziehen:
- Einfachheit: Wenn deine Infrastruktur ohnehin auf Postgres setzt, musst du keinen weiteren Dienst installieren oder konfigurieren. Das senkt den Wartungsaufwand und die Zeit, bis deine Queue läuft.
- Datenkonsistenz: Postgres ist eine robuste, ACID-konforme Datenbank. Wenn die Tasks in derselben Datenbank verwaltet werden, sind starke Datenkonsistenz und transaktionale Integrität gesichert.
- Zuverlässigkeit: Wenn du Postgres schon deine kritischen Anwendungsdaten anvertraust, kannst du ihm auch die Verwaltung deiner Hintergrund-Jobs anvertrauen. Da alles an einem Ort liegt, schaffst du dir keine zusätzlichen Fehlerquellen.
- Einfaches Backup: Du sicherst deine Queue zusammen mit deinen Anwendungsdaten in einem einzigen Vorgang. Keine separate Backup-Strategie für ein weiteres System nötig.
Allerdings solltest du wissen, dass Postgres nicht die perfekte Queue-Lösung für Szenarien mit riesigem Volumen ist, in denen Zehntausende Jobs pro Sekunde eingestellt werden. Es ist ziemlich leistungsfähig, aber du solltest immer die Performance messen und entsprechend optimieren oder in Extremfällen spezialisiertere Systeme in Betracht ziehen. Für die meisten kleinen bis mittleren Lasten kommt Postgres gut zurecht.
🔔 Wenn du eine Job Queue mit noch mehr Features brauchst, die einfacher zu nutzen ist und eine saubere, einfache REST-API bietet, sieh dir DMSContainer mit seinem Event-Stream-Modul an. Es ist sofort einsatzbereit, und mit DMSContainer kannst du in wenigen Minuten Queues nutzen, E-Mails versenden, Excel- und PDF-Berichte erzeugen (und vieles mehr). Schau es dir an!
2. Grundbegriffe und Definitionen
Bevor wir in die Details gehen, hier die Begriffe, die wir verwenden:
- Job/Task: Wir verwenden die Begriffe Job und Task gleichbedeutend für eine Arbeitseinheit, die verarbeitet werden muss.
- Producer: Eine Instanz (meist Teil deiner Anwendung), die Tasks erzeugt und in die Queue einfügt.
- Consumer/Worker: Ein Hintergrundprozess, der ausstehende Tasks aus der Queue holt, verarbeitet und ihren Status aktualisiert.
FOR UPDATE SKIP LOCKED: Ein Postgres-Feature, mit dem du Zeilen beim Auswählen sperrst und dabei Zeilen überspringst, die bereits von einer anderen Transaktion gesperrt sind. Das ist der Schlüssel, um Konflikte zu vermeiden, wenn mehrere Worker nach Tasks pollen.
3. Datenbankschema-Design
Eine Tabelle für Tasks zu entwerfen ist relativ einfach. Hier ein minimales Schema, das du verwenden könntest:
CREATE TABLE tasks (
id SERIAL PRIMARY KEY,
status VARCHAR(50) NOT NULL DEFAULT 'pending',
payload JSONB NOT NULL,
created_at TIMESTAMP WITH TIME ZONE DEFAULT now(),
updated_at TIMESTAMP WITH TIME ZONE DEFAULT now()
);
Das machen die einzelnen Spalten:
- id: Eine eindeutige, automatisch erzeugte Kennung für jeden Task.
- status: Gibt den Zustand des Tasks an, etwa ‘pending’, ‘in_progress’, ‘completed’ oder ‘failed’.
- payload: Enthält die Daten, die zur Ausführung des Tasks nötig sind. JSONB erlaubt es, flexibel die unterschiedlichsten Datenstrukturen abzulegen.
- created_at: Zeitstempel, wann der Task angelegt wurde.
- updated_at: Zeitstempel, der bei jeder Statusänderung des Tasks aktualisiert wird.
Wir halten dieses Schema schlank. In der Praxis würdest du vielleicht Indizes auf bestimmte JSON-Felder oder weitere Spalten wie eine Priorität oder Verweise auf andere Tabellen hinzufügen.
4. Locking in Postgres: FOR UPDATE SKIP LOCKED
Der Star der Show ist das Row-Level-Locking von Postgres, genauer SELECT … FOR UPDATE SKIP LOCKED. SKIP LOCKED wurde mit PostgreSQL 9.5 eingeführt und sorgt dafür, dass, wenn mehrere Transaktionen dieselben Zeilen sperren wollen, eine Transaktion sie sperrt, während die anderen die bereits gesperrten Zeilen überspringen.
Du kannst also dieselbe SELECT-Anweisung parallel in mehreren Workern ausführen, ohne Duplikate zu riskieren. Jeder Worker schnappt sich eine Zeile (also einen Task) und sperrt sie, sodass niemand sonst sie bekommt. Das ist ideal, um Tasks auf mehrere Worker zu verteilen, weil es Kollisionen vermeidet.
Kurz gesagt:
FOR UPDATE: Setzt eine Sperre auf die Zeile(n), sodass keine andere Transaktion sie ändern kann, bis die Sperre aufgehoben wird.SKIP LOCKED: Weist Postgres an, nicht auf gesperrte Zeilen zu warten, sondern sie zu überspringen und nur die Zeilen zu sperren, die gerade nicht gesperrt sind.
5. Neue Tasks einfügen
Wenn deine Anwendung einen neuen Task anlegt, muss sie nur eine Zeile in die Tabelle tasks einfügen. Ein einfaches Beispiel:
INSERT INTO tasks (payload)
VALUES ('{"task_type": "send_email", "recipient": "[email protected]"}');
Der status ist standardmäßig ‘pending’, created_at standardmäßig now(), und auch updated_at wird auf now() gesetzt. Das ist ganz normales DML, wie du es kennst, keine Überraschungen.
6. Polling nach zu bearbeitenden Tasks
Jetzt kommt der wichtige Teil: Wie holen wir Tasks so ab, dass immer nur ein Worker einen bestimmten Task bearbeitet? Nehmen wir an, wir haben mehrere Worker-Prozesse (oder Threads), die jeweils eine SELECT-Anweisung ausführen, um den nächsten ausstehenden Task zu finden.
Hier eine vereinfachte Version:
BEGIN;
WITH cte_task AS (
SELECT id
FROM tasks
WHERE status = 'pending'
ORDER BY id
FOR UPDATE SKIP LOCKED
LIMIT 1
)
UPDATE tasks
SET status = 'in_progress',
updated_at = now()
FROM cte_task
WHERE tasks.id = cte_task.id
RETURNING tasks.id, tasks.payload;
All das läuft in einer einzigen Transaktion:
- WITH cte_task: Diese Common Table Expression sperrt genau einen Task im Status ‘pending’.
ORDER BY idsorgt nur für eine deterministische Reihenfolge; du könntest auch nach Erstellungsdatum oder Priorität sortieren. - FOR UPDATE SKIP LOCKED: Sperrt die Zeile und überspringt sie, wenn ein anderer Worker dieselbe Zeile bereits gesperrt hat.
- LIMIT 1: Stellt sicher, dass wir immer nur einen Task auf einmal sperren.
- UPDATE: Nachdem die CTE die Zeile gefunden hat, setzen wir ihren Status auf
in_progressund weisen den Task damit diesem Worker zu. - RETURNING: Wir bekommen
idundpayloaddes Tasks zurück, den wir bearbeiten wollen.
Liefert die Abfrage keine Zeile, gibt es keine Tasks im Status ‘pending’. Der Worker kann einfach committen und eine Weile schlafen, bevor er es erneut versucht.
Liefert sie eine Zeile, kann der Worker den Task bearbeiten. Ist die Arbeit erledigt, setzt er mit einem weiteren UPDATE den Task auf ‘completed’ oder auf ‘failed’, falls bei der Verarbeitung ein Fehler aufgetreten ist.
So wird jeder Task in einem atomaren Schritt gesperrt und aktualisiert, und andere Worker können sich nicht gleichzeitig denselben Task schnappen.
7. Task-Status aktualisieren
Hat ein Worker den Task erledigt, muss er die Tabelle tasks entsprechend aktualisieren:
UPDATE tasks
SET status = 'completed',
updated_at = now()
WHERE id = <task_id>;
Oder, wenn die Verarbeitung fehlschlägt:
UPDATE tasks
SET status = 'failed',
updated_at = now()
WHERE id = <task_id>;
Vielleicht willst du auch Fehlermeldungen, die Anzahl der Wiederholungen oder andere relevante Metadaten festhalten. Wenn du damit rechnest, dass Tasks fehlschlagen und erneut verarbeitet werden müssen, nimm diese Felder in dein Schema auf.
8. Parallele Worker und Nebenläufigkeit
Mit FOR UPDATE SKIP LOCKED können mehrere Worker dieselbe Polling-Abfrage gleichzeitig ausführen, ohne sich in die Quere zu kommen. Jeder Worker sperrt am Ende andere Zeilen. Sperrt ein Worker eine Zeile, überspringen die anderen sie. Dieser Ansatz ist ideal, um Tasks auf einen Pool von Workern zu verteilen.
Du musst allerdings steuern, wie oft jeder Worker die Tabelle abfragt. Zu häufiges Polling kann deine Datenbank stark belasten, vor allem bei vielen Worker-Prozessen. Es gilt, abzuwägen zwischen schneller Übernahme von Tasks und einer Datenbank, die nicht mit zu vielen SELECT-Anweisungen bombardiert wird.
9. Wiederholungen und Fehlerbehandlung
Job-Verarbeitung in der Praxis ist nie zu 100 % fehlerfrei. Du brauchst eine Strategie für fehlgeschlagene Tasks. Angenommen, ein Worker übernimmt einen Task, aber die Ausführung scheitert an einem Netzwerkfehler oder einem anderen Problem. Du kannst:
- Den Task als ‘failed’ markieren und die Fehlerdetails protokollieren. Später kann ein anderes Subsystem oder ein eigenes Skript nach ‘failed’-Tasks suchen und entscheiden, ob sie erneut versucht werden.
- Automatisch eine begrenzte Anzahl von Malen wiederholen. Du könntest eine Spalte
retry_countin der Tabelletasksführen. Liegt sie unter einem Schwellwert, setzt du den Task wieder auf ‘pending’ und erhöhstretry_count.
Hier ein Beispiel, das einen Task als fehlgeschlagen markiert und retry_count erhöht:
ALTER TABLE tasks ADD COLUMN retry_count INT NOT NULL DEFAULT 0;
BEGIN;
UPDATE tasks
SET status = 'failed',
retry_count = retry_count + 1,
updated_at = now()
WHERE id = <task_id>;
COMMIT;
Von dort aus kann ein separater Prozess oder ein Cron-Job nach failed-Tasks mit retry_count < 5 suchen und sie wieder auf ‘pending’ setzen, sie also für einen weiteren Versuch erneut einreihen. In einfacheren Szenarien kann der Worker retry_count auf retry_count + 1 setzen und die Verarbeitung des Tasks nach einer festgelegten Zahl von Versuchen einstellen.
10. Vor- und Nachteile
Postgres als Job Queue zu verwenden kann sehr bequem sein, ist aber keine Universallösung. Hier eine Liste von Vor- und Nachteilen, die dem allgemeinen Konsens unter Entwicklern entspricht:
Vorteile
- Keine zusätzliche Infrastruktur: Postgres hast du schon, es gibt nichts weiter zu installieren oder zu warten.
- Transaktionale Sicherheit: Wenn dein Code eng mit den Daten in Postgres verbunden ist, vereinfacht es die Konsistenz, wenn die Tasks am selben Ort liegen.
- Row-Level-Locking: Postgres hat robuste Mechanismen für Nebenläufigkeit.
FOR UPDATE SKIP LOCKEDist gut getestet und zuverlässig. - Backup und Restore: Alle Daten, einschließlich der Tasks, lassen sich gemeinsam sichern.
Nachteile
- Skalierbarkeit: Bei extrem hohem Durchsatz oder wenn du fortgeschrittene Queue-Features brauchst (etwa ausgefeiltes Routing, Priority Queues usw.), sind spezialisierte Systeme vielleicht besser.
- Datenbanklast: Das Polling nach Tasks kann deine primäre Datenbank zusätzlich belasten. Mit durchdachter Architektur oder einer Replica lässt sich das abmildern, aber es ist ein weiterer Punkt, an den man denken muss.
- Eingeschränkte Features: Systeme wie RabbitMQ oder Kafka bieten ausgefeiltes Message-Routing, Fan-out und mehr. Postgres ist in dieser Hinsicht einfacher, aber für bestimmte Muster auch weniger flexibel.
11. Eine Beispielimplementierung Schritt für Schritt
Gehen wir ein hypothetisches Szenario mit einem kleinen Tech-Startup durch: Acme Email Services. Die Firma verschickt transaktionale E-Mails. Sie muss E-Mail-Versand-Tasks in eine Queue stellen, damit ihr SMTP-Server nicht überlastet wird.
Die Tabelle
tasksanlegen:CREATE TABLE tasks ( id SERIAL PRIMARY KEY, status VARCHAR(50) NOT NULL DEFAULT 'pending', payload JSONB NOT NULL, created_at TIMESTAMP WITH TIME ZONE DEFAULT now(), updated_at TIMESTAMP WITH TIME ZONE DEFAULT now() );Tasks einfügen (auf Producer-Seite, z. B. aus einem API-Endpunkt):
INSERT INTO tasks (payload) VALUES ('{"task_type": "send_email", "subject": "Welcome!", "recipient": "[email protected]", "body": "Hello and welcome!"}');Worker-Prozess (zum Beispiel als Pseudocode in Python):
import psycopg2 import time def poll_and_process_tasks(): while True: conn = psycopg2.connect("dbname=acme user=postgres password=postgres") conn.autocommit = False try: with conn.cursor() as cur: cur.execute(""" WITH cte_task AS ( SELECT id, payload FROM tasks WHERE status = 'pending' ORDER BY id FOR UPDATE SKIP LOCKED LIMIT 1 ) UPDATE tasks SET status = 'in_progress', updated_at = now() FROM cte_task WHERE tasks.id = cte_task.id RETURNING tasks.id, tasks.payload; """) row = cur.fetchone() if row: task_id, task_payload = row # Den Task verarbeiten process_email(task_payload) # deine eigene Funktion # Als erledigt markieren cur.execute(""" UPDATE tasks SET status = 'completed', updated_at = now() WHERE id = %s """, (task_id,)) conn.commit() else: conn.rollback() # Keine Tasks, kurz schlafen und dann erneut prüfen time.sleep(5) finally: conn.close() def process_email(payload): # Pseudocode für den E-Mail-Versand # payload['task_type'] == 'send_email' # payload['recipient'], payload['subject'] usw. verwenden print(f"Sending email to {payload['recipient']}...") # Hier die E-Mail tatsächlich versendenIm obigen Skript:
- verbinden wir uns mit der Datenbank und starten eine Transaktion,
- versuchen wir, mit der Technik
FOR UPDATE SKIP LOCKEDeinen einzelnenpending-Task zu holen, - markieren wir ihn, wenn das klappt, sofort als
in_progress, - verarbeiten wir dann die E-Mail (der eigentliche Versandcode fehlt der Kürze halber),
- setzen wir bei Erfolg seinen Status auf
completed, - machen wir, wenn keine Tasks ausstehen, ein Rollback (damit Sperren oder halbe Updates freigegeben werden) und schlafen eine Weile.
⭐ Die vollständige Delphi-Version des Job-Queue-Systems erscheint in den nächsten Tagen für die PATREON-Unterstützer.
- Fehlerbehandlung: Tritt beim Versand der E-Mail ein Fehler auf, kannst du die Exception abfangen, den Task als
failedmarkieren und eventuell den Fehler protokollieren. Bei Bedarf kommt dann eine Strategie zum erneuten Einreihen hinzu.
12. Praktische Tipps
- Indizierung: Bei sehr vielen Tasks lohnt sich ein Index auf
(status)oder(status, id), um ausstehende Tasks schneller zu finden. - Polling-Frequenz begrenzen: Baue eine kurze Pause oder eine Backoff-Strategie in deine Worker ein, um die Datenbanklast zu senken.
- Eine separate Tabelle verwenden: Hat deine Anwendung mehrere Arten von Tasks oder eine riesige Zahl davon, kannst du sie auf verschiedene Tabellen oder sogar verschiedene Schemas verteilen. Das hilft bei Organisation und Performance-Tuning.
- Monitoring: Behalte die Zahl der Tasks in jedem Status im Blick. Werkzeuge wie Grafana oder eigene Metriken können dich alarmieren, wenn sich ein Rückstau bildet (z. B. wenn Tasks zu lange auf ‘pending’ stehen).
- Transaktionsgröße: Achte auf die Transaktionsgrenzen. Wenn du Hunderte Tasks in einer einzigen Transaktion sperren und verarbeiten willst, hältst du Sperren womöglich zu lange. Normalerweise verarbeitest du Tasks einzeln oder in kleinen Batches.
13. Fazit und eine Frage zum Nachdenken
Inzwischen hast du gesehen, dass Postgres als Job Queue praktisch und elegant sein kann, vor allem wenn deine Umgebung Postgres ohnehin für andere Aufgaben nutzt. Du vermeidest die Komplexität zusätzlicher Systeme, bekommst robuste transaktionale Garantien und nutzt die Nebenläufigkeits-Features von Postgres, um Tasks ohne Duplikate auf mehrere Worker zu verteilen.
Mit PostgreSQL kannst du einen Cluster von Worker-Prozessen bauen, der zuverlässig Tausende Tasks pro Tag abarbeitet. Lege alle Logs und Job-Metadaten in Postgres ab, dann kannst du historische Daten ganz einfach abfragen, die Performance analysieren und die Fehlerbehebung steuern. Lösungen wie DMSContainer/EventStream!, Redis, RabbitMQ oder Kafka mögen in bestimmten hochskalierten oder spezialisierten Szenarien besser passen, aber der Ansatz mit der Postgres Job Queue ist für viele kleine bis mittlere Anforderungen ein starker Kandidat.
Frage zum Nachdenken: Brauchst du bei der Last und Infrastruktur deiner Organisation wirklich die speziellen Features eines dedizierten Queue-Systems, oder kann Postgres deine Job-Verarbeitung gut genug stemmen, um deine Architektur zu vereinfachen? Schreib uns eine E-Mail, wenn du spezialisierte Beratung und Entwicklung brauchst.
Weitere Artikel zu PostgreSQL:
- PostgreSQL Composite Types: Das Feature zwischen einfachen Spalten und JSON - Fortgeschrittene Datenmodellierung für komplexe Queues
- PostgreSQL: SERIAL oder IDENTITY? - Best Practices für Primärschlüssel
- LIKE-Abfragen mit pg_trgm beschleunigen - Techniken zur Abfrageoptimierung
Quellen
Offizielle PostgreSQL-Dokumentation:
https://www.postgresql.org/docs/current/sql-select.html#SQL-FOR-UPDATE-SHARE
Details zu Verwendung und Syntax vonFOR UPDATE SKIP LOCKED.PostgreSQL-Wiki zum Thema Queueing:
https://wiki.postgresql.org/wiki/Category:Queueing
Verschiedene Diskussionen der Community und Erweiterungen für Queue-ähnliche Funktionalität in Postgres.ACID-Eigenschaften:
https://en.wikipedia.org/wiki/ACID
Erklärung der Transaktionsgarantien, die Postgres von Haus aus bietet.
Das war’s. Ob Nebenprojekt, Aufwertung eines internen Tools oder eine schnelle, aber zuverlässige Lösung für die Hintergrundverarbeitung: Postgres kann der Held deiner Job Queue sein. Mit einem sorgfältig entworfenen Schema, der Magie von FOR UPDATE SKIP LOCKED und ein wenig Worker-Logik bekommst du ein System, das deine Tasks ordnet, deine Worker beschäftigt und deine Infrastruktur vereinfacht. Viel Freude mit der Einfachheit und Zuverlässigkeit deines neuen Queue-Systems, angetrieben vom guten alten Postgres!
Du willst mehr?
Tritt der
bei, um das Projekt zu unterstützen und Zugang zu Premium-Inhalten wie Artikeln, Videos und weiteren Einblicken zu bekommen.
Denk daran, der PATREON-Community beizutreten, um wertvolle Informationen, Tutorials und Einblicke sowie bevorzugten Support und mehr zu bekommen. Du kannst das Projekt auch über Buy Me a Coffe unterstützen und bekommst dieselben Vorteile.
Im PATREON-Artikel findest du 2 Projekte:
- JobProducer.dproj - das Client-Programm, das Jobs anfordert. Im Beispiel kann es 2 Arten von Jobs anfordern (“send_email”, “create_report”), du kannst es aber auf alles erweitern, was du brauchst.
- QueueWorker.dproj - der eigentliche Worker. Du kannst mehrere Instanzen davon starten und prüfen, dass das System jeden einzelnen Job genau einer Worker-Instanz zuweist. Stürzt der Worker ab, nachdem er einen Job übernommen hat, steht dieser Job wieder für andere Worker bereit.
Viel Spaß!
– Daniele Teti
Comments