Das Transactional-Outbox-Pattern gegen Dual-Write
Wie das Transactional-Outbox-Pattern das Dual-Write-Problem löst und zuverlässiges Event-Publishing garantiert, ohne DB und Broker zu koppeln.

In fast jeder verteilten Architektur steckt dieser Bug irgendwo im Code: Die Bestellung wird in der Datenbank gespeichert, und danach wird ein order.created-Event an Kafka publiziert. Wenn dieser Publish fehlschlägt — und früher oder später wird er das —, existiert die Bestellung zwar, aber alle nachgelagerten Services bekommen für immer nichts davon mit. Man kann Retry-Logik ergänzen, aber das verschiebt das Problem nur. Die eigentliche Lösung ist architektonisch: Datenbankschreibvorgang und Broker-Publish nicht länger als zwei getrennte Operationen behandeln.
Das Transactional-Outbox-Pattern löst dieses Problem auf Infrastrukturebene. Es tauscht die Illusion von Atomarität über zwei Systeme hinweg gegen einen echten atomaren Schreibvorgang in einem einzigen System und propagiert die Änderung anschließend asynchron über einen separaten Relay-Prozess. Das Ergebnis ist konsistent, gut debugbar und horizontal skalierbar.
Das Dual-Write-Problem
Sobald man nacheinander in zwei getrennte Systeme schreibt, riskiert man einen partiellen Fehler. Der klassische Übeltäter: eine relationale Datenbank und ein Message-Broker werden im selben Request-Handler beschrieben:
// ❌ Dual-write — the publish can fail after the DB commit
async function createOrder(data: CreateOrderInput): Promise<Order> {
const order = await db.orders.create(data);
await kafka.publish("order.created", order); // what if this throws?
return order;
}
// ✅ Single transaction — the event is part of the same write
async function createOrder(data: CreateOrderInput): Promise<Order> {
return db.transaction(async (tx) => {
const order = await tx.orders.create(data);
await tx.outbox.insert({
aggregateType: "order",
aggregateId: order.id,
eventType: "order.created",
payload: order,
});
return order;
});
}Die Outbox-Zeile und die Bestellzeile werden atomar committet. Macht die Transaktion ein Rollback, existieren weder die Bestellung noch der Outbox-Eintrag. Wird sie committet, existieren beide. Der Message-Broker liegt nie im Hot Path der Transaktion — er erfährt erst später vom Event, über einen Relay-Prozess, der committete Outbox-Zeilen ausliest.
Das Schema der Outbox-Tabelle
Die Outbox-Tabelle ist das Herzstück des Patterns. Sie sollte von Anfang an für geordnete, zuverlässige Verarbeitung ausgelegt sein:
CREATE TABLE outbox_events (
id UUID DEFAULT gen_random_uuid() PRIMARY KEY,
aggregate_type TEXT NOT NULL,
aggregate_id TEXT NOT NULL,
event_type TEXT NOT NULL,
payload JSONB NOT NULL,
created_at TIMESTAMPTZ DEFAULT now() NOT NULL,
processed_at TIMESTAMPTZ,
attempts INT DEFAULT 0 NOT NULL,
last_error TEXT
);
CREATE INDEX idx_outbox_unprocessed
ON outbox_events (created_at)
WHERE processed_at IS NULL;Der partielle Index auf processed_at IS NULL ist entscheidend. Sobald sich Millionen bereits verarbeiteter Zeilen angesammelt haben, durchsucht der Relay-Prozess nur noch die kleine Teilmenge der unverarbeiteten Einträge. Ohne diesen Index wird aus jedem Poll ein vollständiger sequenzieller Scan, der mit der Zeit immer langsamer wird.
Die Spalten attempts und last_error sind kein optionales Extra — mit ihnen erkennt man festhängende Events in Produktion, ohne sich mühsam durch Logs zu graben.
Den Relay-Prozess bauen
Der Relay fragt unverarbeitete Events per Polling ab und publiziert sie an den Broker. Am Anfang zählt vor allem Einfachheit und Korrektheit:
async function processOutboxBatch(batchSize = 100): Promise<number> {
const events = await db.$transaction(async (tx) => {
// Lock rows to prevent concurrent relays from double-publishing
const rows = await tx.$queryRaw<OutboxEvent[]>`
SELECT * FROM outbox_events
WHERE processed_at IS NULL
ORDER BY created_at ASC
LIMIT ${batchSize}
FOR UPDATE SKIP LOCKED
`;
if (rows.length === 0) return [];
await Promise.all(
rows.map((event) =>
kafka.publish(event.event_type, event.payload, {
key: event.aggregate_id, // same aggregate → same partition → ordered delivery
})
)
);
const ids = rows.map((r) => r.id);
await tx.$executeRaw`
UPDATE outbox_events
SET processed_at = now()
WHERE id = ANY(${ids}::uuid[])
`;
return rows;
});
return events.length;
}
// Relay loop — run one per service replica
async function startRelay(): Promise<void> {
while (true) {
const processed = await processOutboxBatch();
// Back off when idle to avoid hammering the DB unnecessarily
await sleep(processed === 0 ? 500 : 50);
}
}FOR UPDATE SKIP LOCKED ist das Detail, das horizontale Skalierung erst möglich macht. Mehrere Relay-Instanzen können parallel laufen — jede schnappt sich die Zeilen, die kein anderer Relay gerade gesperrt hält. So bekommt man sichere Nebenläufigkeit, ganz ohne verteiltes Lock und ohne Engpass durch eine einzelne Replica.
Polling versus Change Data Capture
Für die meisten Teams ist Polling der richtige Ausgangspunkt. Mit dem partiellen Index sorgt ein Poll-Intervall von 200–500ms für eine Latenz unter einer Sekunde bei vernachlässigbarer Datenbanklast.
| Ansatz | Latenz | DB-Last | Betriebliche Komplexität |
|---|---|---|---|
| Polling (500ms-Intervall) | ~500ms | Niedrig | Minimal — eine Schleife pro Replica |
| Polling (100ms-Intervall) | ~100ms | Moderat | Minimal |
| CDC (Debezium + Kafka Connect) | Nahezu in Echtzeit | Sehr niedrig | Hoch — Connect-Cluster, Schema Registry, Connector-Verwaltung |
Change Data Capture liest direkt aus dem WAL-Stream von PostgreSQL und macht Polling komplett überflüssig. Das ist die richtige Wahl, wenn im großen Maßstab eine Latenz unter 100ms nötig ist oder wenn die Outbox-Tabelle unter extremem Schreibdruck steht. Für die überwiegende Mehrheit der Services bleibt Polling der richtige Ausgangspunkt und sollte die Implementierung bleiben, bis Produktionsdaten den betrieblichen Mehraufwand von Debezium tatsächlich rechtfertigen.
Versieh deinen Relay mit einer Metrik: dem Alter des ältesten unverarbeiteten Outbox-Events. Bleibt dieser Wert im stationären Betrieb unter einer Sekunde, macht Polling seinen Job. Greif erst dann zu CDC, wenn dieses SLO mit Polling wirklich nicht erreichbar ist.
Idempotenz beim Consumer ist Pflicht, keine Option
Das Outbox-Pattern garantiert at-least-once-Zustellung (mindestens einmal) — niemals exactly-once. Publiziert der Relay ein Event erfolgreich, schlägt aber der Commit von processed_at fehl, bevor die Transaktion abgeschlossen ist (Netzwerkaussetzer, Prozess mitten im Ablauf beendet), wird genau dieses Event im nächsten Poll-Zyklus erneut publiziert. Das ist erwartetes Verhalten, kein Bug.
Jeder Consumer muss mit Duplikaten umgehen können:
async function handleOrderCreated(event: OrderCreatedEvent): Promise<void> {
const alreadyProcessed = await db.processedEvents.findUnique({
where: { eventId: event.id },
});
if (alreadyProcessed) {
logger.debug({ eventId: event.id }, "skipping duplicate event");
return;
}
// Transactionally apply the effect and record the deduplication key
await db.$transaction(async (tx) => {
await tx.inventory.reserveItems(event.payload.lineItems);
await tx.processedEvents.create({ data: { eventId: event.id } });
});
}Die Tabelle processedEvents fungiert als Deduplizierungs-Log. Die Prüfung und der anschließende Schreibvorgang laufen innerhalb einer Transaktion mit einem Unique Constraint auf eventId — laufen zwei Consumer um dasselbe Event, gelingt es dem einen, während der andere eine Constraint-Verletzung erhält, die man gefahrlos als Duplikat abfangen kann.
Aufbewahrung und Bereinigung
Verarbeitete Outbox-Zeilen sind totes Gewicht. Ein geplanter Cleanup-Job verhindert, dass die Tabelle unbegrenzt wächst:
-- Run daily via pg_cron or your job scheduler
DELETE FROM outbox_events
WHERE processed_at IS NOT NULL
AND processed_at < now() - INTERVAL '7 days';Sieben Tage Aufbewahrung geben dir ein Replay-Fenster, um Incidents zu debuggen, ohne dass der Audit Trail unbegrenzt wächst. Löst dein Bereitschaftsdienst Probleme normalerweise innerhalb von 24 Stunden, reichen auch drei Tage. Der partielle Index hält die Query-Performance stabil, egal wie viele historische Daten aufbewahrt werden — aber kleinere Tabellen bedeuten schnellere Vacuums und niedrigere Storage-Kosten.
Die wichtigsten Erkenntnisse
- Nie Dual-Write betreiben — erst in die Datenbank schreiben und dann an einen Broker publizieren ist ein Konsistenz-Bug, der nur darauf wartet, im ungünstigsten Moment aufzutauchen
- Die Outbox-Zeile ist dein Event — committe sie atomar zusammen mit deinem Domain-Schreibvorgang; sonst darf nichts die Transaktionsgrenze überschreiten
FOR UPDATE SKIP LOCKEDermöglicht sicheres horizontales Skalieren des Relays — mehrere Instanzen koordinieren sich allein über die Datenbank, ganz ohne externe Locks- Erst Polling, dann messen, dann über CDC nachdenken — eine nahezu echtzeitfähige Latenz rechtfertigt selten den Betriebsaufwand von Debezium, so früh im Lebenszyklus eines Systems
- At-least-once-Zustellung ist ein Vertrag, kein Bug — jeder Consumer sollte von Anfang an idempotent entworfen werden, nicht nachträglich umgebaut
- Das Alter des ältesten unverarbeiteten Events verfolgen — das ist die nützlichste Kennzahl für die Relay-Gesundheit, aussagekräftiger als der Durchsatz verarbeiteter Events pro Sekunde


