sequenceDiagram
participant A as Anwendung
participant D as Datenbank
participant R as Relay
participant B as Broker
A->>D: BEGIN, INSERT bestellung, INSERT outbox, COMMIT
R->>D: lies Zeilen mit gesendet = 0
R->>B: veröffentliche
R->>D: markiere gesendet = 1
Note over R,B: Stirbt das Relay nach dem Senden, aber vor dem Markieren, geht die Nachricht doppelt raus
Outbox, Dead Letter Queue, Pub/Sub und Event-Log
Track Konzepte · Messaging und Zuverlässigkeit · ca. 45 Min.
Was du aus Teil a brauchst
Eine Queue liefert at-least-once: Fehlt das Ack, kommt die Nachricht erneut, also kann dieselbe Nachricht mehrfach ankommen. Ein idempotenter Konsument merkt sich Nachrichten-IDs in einer Dedup-Tabelle, in derselben Transaktion wie die Wirkung. Falls das neu ist: Teil 1.
Worum es geht
Auf der Empfängerseite hast du jetzt Ruhe vor Duplikaten. Die Senderseite hat ein eigenes Problem: Bestellung speichern und Nachricht senden sind zwei Systeme und lassen sich nicht mit einem Commit schreiben (Dual Write). Das Outbox-Pattern löst das. Dazu kommt, was mit einer Nachricht passiert, die immer wieder scheitert (Dead Letter Queue, DLQ), und welches Modell zu welchem Bedarf passt: Queue, Pub/Sub oder Event-Log. Am Ende von Teil 2 kannst du die Outbox reparieren, eine DLQ schreiben und die Architektur begründet wählen.
Typische Situation aus JS/TS
Du kennst das aus Node-Code:
await db.insert("bestellung", bestellung); // 1. Datenbank
await queue.publish("bestellt", bestellung); // 2. Nachricht
// Prozess stirbt zwischen 1 und 2: Bestellung gespeichert, Nachricht nie gesendet.
// Reihenfolge vertauschen: Nachricht gesendet, Bestellung nie gespeichert.Beide Reihenfolgen sind falsch. Die Datenbank-Transaktion kann das Senden nicht einschließen. Die Lösung ist, die Nachricht in derselben Transaktion in eine Tabelle zu schreiben und von dort zuverlässig zu versenden.
Konzept
Schritt 1: Das Outbox-Pattern gegen das Dual-Write-Problem
Für die Beispiele brauchst du sqlite3 (wie in Teil 1):
Dasselbe Problem tritt beim Senden auf. Eine Bestellung wird in der Datenbank gespeichert und ein Ereignis an die Queue gesendet. Das sind zwei Systeme, also kein gemeinsamer Commit. Man nennt das Dual Write (doppeltes Schreiben).
Ausgabe: Prozess stirbt, dann 1 0. Die Bestellung existiert, das Ereignis nie. Der Versand erfährt nichts. Dreht man die Reihenfolge um (erst publishen, dann committen), ist es umgekehrt: Ein Ereignis für eine Bestellung, die es nicht gibt.
Die Lösung ist das Outbox-Pattern: Das Ereignis wird nicht gesendet, sondern in einer Tabelle outbox in derselben Transaktion wie die Bestellung gespeichert. Ein getrennter Prozess, der Relay (auch Publisher genannt), liest unversendete Zeilen und sendet sie an den Broker. Erst nach erfolgreichem Senden markiert er die Zeile als gesendet.
Beachte die Notiz: Die Outbox macht aus “vielleicht verloren” ein at-least-once. Verloren geht nichts mehr, doppelt kann es werden. Darum braucht der Konsument am Ende trotzdem die Dedup-Tabelle aus Schritt 4. Beide Muster ergänzen sich.
Schritt 2: Dead Letter Queue, Pub/Sub und Event-Log
Dead Letter Queue (DLQ). Manche Nachrichten sind “giftig” (poison message): Der Konsument wirft bei jeder Zustellung eine Exception, etwa wegen eines kaputten Formats. Ohne Gegenmaßnahme kommt sie endlos wieder und blockiert Kapazität. Darum zählt man die Zustellversuche. Nach N Fehlversuchen wandert die Nachricht in eine separate Queue, die DLQ, und die normale Arbeit läuft weiter. In der DLQ liegen die Fälle, die ein Mensch oder ein Reparaturjob ansehen muss (und nach der Reparatur wieder einspeist). Eine DLQ ohne Alarm ist nur ein Friedhof: Überwache ihre Länge.
Drei Messaging-Modelle. Sie unterscheiden sich darin, wer eine Nachricht bekommt und was nach dem Lesen mit ihr passiert:
| Modell | Wer bekommt die Nachricht? | Nach dem Lesen | Typischer Einsatz |
|---|---|---|---|
| Queue (Work Queue, competing consumers) | Ein Konsument aus der Gruppe der Worker | gelöscht (nach Ack) | Arbeit verteilen: Mails, Bildverarbeitung |
| Pub/Sub (Publish/Subscribe) | Jeder Abonnent bekommt eine Kopie | beim klassischen Pub/Sub ohne Speicherung danach weg (konkrete Systeme unterscheiden sich, bitte prüfen) | Ereignisse an mehrere Interessenten melden |
| Event Streaming (Log-basiert) | Jede Consumer Group liest den ganzen Log in eigenem Tempo | bleibt im Log (für eine Aufbewahrungszeit) | Ereignisse mit Reihenfolge, Wiederverarbeitung (replay) |
Beim Event Streaming ist die Nachricht ein Eintrag in einem append-only Log, der in Partitionen (partitions) aufgeteilt ist. Jede Nachricht hat einen Schlüssel (key). Aus dem Schlüssel wird die Partition berechnet, darum landen alle Nachrichten mit demselben Schlüssel in derselben Partition. Reihenfolge ist nur innerhalb einer Partition garantiert, nicht über Partitionen hinweg. Jeder Eintrag hat einen Offset (laufende Nummer in der Partition). Ein Konsument merkt sich seinen Offset, und eine Consumer Group teilt sich die Partitionen: pro Partition liest höchstens ein Mitglied der Gruppe. Verschiedene Gruppen sind voneinander unabhängig und haben eigene Offsets.
Ausgabe (crc32 ist hier nur eine stabile Hashfunktion, echte Systeme nutzen eigene):
kunde-A bestellt -> (0, 0)
kunde-B bestellt -> (1, 0)
kunde-A bezahlt -> (0, 1)
kunde-C bestellt -> (2, 0)
kunde-A versendet -> (0, 2)
kunde-B bezahlt -> (1, 1)
Alle Ereignisse von kunde-A liegen in Partition 0 mit den Offsets 0, 1, 2, also in der richtigen Reihenfolge. Wer neu liest, setzt den Offset zurück. Das ist Replay: Eine Queue kann das nicht, ihre Nachrichten sind nach dem Ack gelöscht.
Weil die Zustellung auch hier at-least-once ist (Offset wird nach der Verarbeitung gespeichert, bei Absturz wird ab dem alten Offset neu gelesen), braucht jeder Consumer einen idempotenten Konsum. Das ist derselbe Gedanke wie in Teil 1, Schritt 4.
Falle
- Dual Write. Zwei Systeme lassen sich nicht mit einem Commit schreiben (Übung 1).
- Retry ohne Obergrenze. Eine giftige Nachricht blockiert ewig, wenn es keine DLQ gibt (Übung 2).
- Reihenfolge über Partitionen erwarten. Sie gilt pro Schlüssel, wenn der Schlüssel richtig gewählt ist.
Übungen
Übung 1: Outbox reparieren (ca. 15 Min.)
Du hast in Schritt 1 gesehen, wie Dual Write ein Ereignis verliert. Baue die Outbox. Die Tabellen bestellung und outbox gibt es schon (outbox.payload ist JSON-Text, outbox.gesendet ist 0 oder 1). Broker.veroeffentliche(payload) wirft ConnectionError, wenn der Broker nicht erreichbar ist.
Schreibe zwei Funktionen:
bestellen(con, betrag): legt die Bestellung an und schreibt in derselben Transaktion eine Outbox-Zeile mit dem Payload{"bestellung_id": id, "betrag": betrag}als JSON-Text. Rückgabe: die Bestellungs-ID.relay(con, broker): sendet alle Zeilen mitgesendet = 0in der Reihenfolge ihreridanbroker.veroeffentliche(...)(als Python-Dict, nicht als Text) und markiert jede Zeile erst nach erfolgreichem Senden als gesendet. Broker-Fehler dürfen nach außen fliegen.
Am Ende gibst du beide als Paar zurück. Die Prüfung testet den Normalfall, einen Absturz mitten in bestellen, einen Broker-Ausfall mitten im Relay und einen zweiten Relay-Lauf.
Bei bestellen zählt, dass es keinen Zustand gibt, in dem nur eines der beiden Dinge gespeichert ist. Bei relay überlege, was passiert, wenn der Broker beim zweiten von drei Nachrichten ausfällt: Was darf als gesendet markiert sein, was nicht?
def bestellen(con, betrag):
with con:
cur = con.execute("INSERT INTO bestellung (betrag) VALUES (?)", (betrag,))
bid = cur.lastrowid
payload = json.dumps({"bestellung_id": bid, "betrag": betrag})
con.execute("INSERT INTO outbox (payload) VALUES (?)", (payload,))
return bid
def relay(con, broker):
zeilen = con.execute(
"SELECT id, payload FROM outbox WHERE gesendet = 0 ORDER BY id").fetchall()
for zeile_id, payload in zeilen:
broker.veroeffentliche(json.loads(payload)) # erst senden ...
with con:
con.execute("UPDATE outbox SET gesendet = 1 WHERE id = ?", (zeile_id,)) # ... dann markieren
(bestellen, relay)Fällt der Broker aus, bleiben die noch nicht gesendeten Zeilen auf gesendet = 0 und gehen beim nächsten Lauf raus. Stirbt das Relay genau zwischen Senden und Markieren, geht eine Nachricht doppelt raus (at-least-once), darum braucht der Konsument die Dedup-Tabelle.
Übung 2: Dead Letter Queue selbst schreiben (ca. 10 Min.)
Schreibe verarbeite_alle(q, handler, max_versuche, dlq). Die Queue ist die aus Teil 1, Schritt 2 (Feld versuche zählt die Zustellungen). Ablauf, solange q.offen() > 0:
m = q.receive(). Ist das ErgebnisNone(alles noch im Visibility Timeout), rufeq.tick()auf und mache weiter.- Rufe
handler(m["body"])auf. Bei Erfolg:q.ack(m["id"]). - Wirft der Handler eine Exception: nicht bestätigen, es sei denn,
m["versuche"] >= max_versuche. Dann hänge{"body": ..., "versuche": ..., "fehler": str(exc)}an die Listedlqan und bestätige die Nachricht (sie liegt jetzt in der DLQ und blockiert die Queue nicht mehr).
Eine nicht bestätigte Nachricht kommt von allein zurück, sobald ihr Timeout abläuft. Woher weißt du beim Fehler, ob es der letzte erlaubte Versuch war?
def verarbeite_alle(q, handler, max_versuche, dlq):
while q.offen() > 0:
m = q.receive()
if m is None:
q.tick()
continue
try:
handler(m["body"])
except Exception as exc:
if m["versuche"] >= max_versuche:
dlq.append({"body": m["body"], "versuche": m["versuche"], "fehler": str(exc)})
q.ack(m["id"])
else:
q.ack(m["id"])
verarbeite_alleÜbung 3: Queue, Pub/Sub oder Event-Log? (ca. 8 Min.)
Beim Bezahlen entsteht pro Bestellung ein Ereignis bezahlt, davor bestellt, danach eventuell storniert. Randbedingungen:
- Versand, Rechnungswesen und Analytics sind drei getrennte Teams. Jedes Team braucht jedes Ereignis.
- Pro Bestellung müssen die Ereignisse in der Reihenfolge
bestellt,bezahlt,storniertverarbeitet werden. Zwischen verschiedenen Bestellungen ist die Reihenfolge egal. - Analytics baut nächste Woche ein neues Dashboard und muss dafür die Ereignisse der letzten 30 Tage noch einmal einlesen.
- Ein Team darf langsam oder wegen eines Deployments eine Stunde offline sein, ohne dass es Ereignisse verpasst.
Welche Architektur passt?
- A Eine gemeinsame Work Queue. Jedes Team betreibt eigene Worker, die aus derselben Queue lesen, und bestätigt jede Nachricht nach der Verarbeitung.
- B Ein partitionierter Event-Log mit der Bestellnummer als Schlüssel. Jedes Team liest als eigene Consumer Group, Aufbewahrung 30 Tage.
- C Pro Team eine eigene Queue. Die Bestell-Anwendung schreibt jedes Ereignis nacheinander in alle drei Queues (ohne Outbox), jedes Team liest dort.
- D Klassisches Pub/Sub ohne Speicherung. Ein Topic verteilt jede Nachricht an alle gerade verbundenen Abonnenten und löscht sie danach sofort.
Gehe die vier Randbedingungen einzeln durch und frage bei jeder Option: Was passiert mit der Nachricht nach dem Lesen? Wer bekommt sie? Wo wird die Reihenfolge festgelegt?
antwort = "B"
antwortB erfüllt alle vier Bedingungen: eigene Consumer Groups liefern jedem Team jedes Ereignis (1), der Schlüssel Bestellnummer legt alle Ereignisse einer Bestellung in dieselbe Partition und damit in Reihenfolge (2), der Log bleibt 30 Tage gespeichert, ein neuer Offset ermöglicht Replay (3), und der Offset jedes Teams hält fest, wo es stehen geblieben ist (4).
Merksatz
Jeder Sender koppelt Schreiben und Senden über eine Outbox, jede Nachricht wandert nach N Fehlversuchen in eine DLQ, und die Wahl zwischen Queue, Pub/Sub und Event-Log folgt aus Reihenfolge, Mehrfachabnehmern und Replay.
Prüfstein
Warum reicht “erst speichern, dann senden” nicht, wie löst das Outbox-Pattern das Problem, und was bleibt trotzdem übrig, das der Konsument selbst lösen muss?
Zurück zu Teil 1.
Quelle: quellen/konzeptuebersicht-software-grundlagen.md, Abschnitt “7. Verteilte Systeme und Betrieb” (Messaging: Queue vs. Pub/Sub vs. Event Streaming, Outbox-Pattern); quellen/konzeptuebersicht-software-fortgeschritten.docx, Schicht 7 “Log-basierte Systeme” (Partitionierung und Reihenfolge, Consumer Groups, Offsets, Dead Letter Queues).
Über die Quelle hinaus (allgemeines Fachwissen, die Quellen sind Stichwortlisten): die Aufteilung des Outbox-Relays in Senden und Markieren, die Mechanik von Partitionen, Offsets und Replay. Die Simulationen (Queue, Log) sind didaktische Vereinfachungen und bilden kein konkretes Produkt (z. B. SQS, RabbitMQ, Kafka) ab. Produktdetails zu Pub/Sub-Speicherung sind bitte zu prüfen.