Become a member!

PostgreSQL als Job Queue: Braucht man wirklich Redis oder RabbitMQ?

🌐
Dieser Artikel ist auch in anderen Sprachen verfügbar:
🇪🇸 Español  •  🇧🇷 Português  •  🇬🇧 English

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.

Offenes Meer unter bewölktem Himmel bei Sonnenuntergang, Titelfoto von Daniele Teti


Inhaltsverzeichnis

  1. Warum Postgres als Job Queue?
  2. Grundbegriffe und Definitionen
  3. Datenbankschema-Design
  4. Locking in Postgres: FOR UPDATE SKIP LOCKED
  5. Neue Tasks einfügen
  6. Polling nach zu bearbeitenden Tasks
  7. Task-Status aktualisieren
  8. Parallele Worker und Nebenläufigkeit
  9. Wiederholungen und Fehlerbehandlung
  10. Vor- und Nachteile
  11. Eine Beispielimplementierung Schritt für Schritt
  12. Praktische Tipps
  13. Fazit und eine Frage zum Nachdenken

PATREON-Unterstützer haben Zugriff auf ein vollständiges Delphi-Projekt, das ein Job-Queue-System implementiert.

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:

  1. 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.
  2. Datenkonsistenz: Postgres ist eine robuste, ACID-konforme Datenbank. Wenn die Tasks in derselben Datenbank verwaltet werden, sind starke Datenkonsistenz und transaktionale Integrität gesichert.
  3. 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.
  4. 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:

  1. WITH cte_task: Diese Common Table Expression sperrt genau einen Task im Status ‘pending’. ORDER BY id sorgt nur für eine deterministische Reihenfolge; du könntest auch nach Erstellungsdatum oder Priorität sortieren.
  2. FOR UPDATE SKIP LOCKED: Sperrt die Zeile und überspringt sie, wenn ein anderer Worker dieselbe Zeile bereits gesperrt hat.
  3. LIMIT 1: Stellt sicher, dass wir immer nur einen Task auf einmal sperren.
  4. UPDATE: Nachdem die CTE die Zeile gefunden hat, setzen wir ihren Status auf in_progress und weisen den Task damit diesem Worker zu.
  5. RETURNING: Wir bekommen id und payload des 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:

  1. 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.
  2. Automatisch eine begrenzte Anzahl von Malen wiederholen. Du könntest eine Spalte retry_count in der Tabelle tasks führen. Liegt sie unter einem Schwellwert, setzt du den Task wieder auf ‘pending’ und erhöhst retry_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 LOCKED ist 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.

  1. Die Tabelle tasks anlegen:

    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()
    );
    
  2. 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!"}');
    
  3. 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 versenden
    

    Im obigen Skript:

    • verbinden wir uns mit der Datenbank und starten eine Transaktion,
    • versuchen wir, mit der Technik FOR UPDATE SKIP LOCKED einen einzelnen pending-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.

  1. Fehlerbehandlung: Tritt beim Versand der E-Mail ein Fehler auf, kannst du die Exception abfangen, den Task als failed markieren 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:


Quellen


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 PATREON Community 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.

PATREON-Unterstützer haben Zugriff auf ein vollständiges Delphi-Projekt, das ein Job-Queue-System implementiert.

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