Outbox-und-Idempotenz-Lab

22 Java-Dateien. 7 zentrale Dateien werden direkt mit echtem Quellcode und ihrem Zusammenspiel erklärt.

Zurück zu Code-Labs

Kapitel 04 · Outbox, Idempotenz und Crash-Fälle

Was dieses Lab zeigt

Macht kritische Crash-Fenster zwischen Datenbank und Broker sichtbar. Publisher-Retries, Event-Envelopes und idempotente Consumer verhindern doppelte oder verlorene Verarbeitung.

Lernziele

  • Crash-Fenster analysieren
  • Idempotente Consumer bauen
  • Retry und Deduplication testen

Technik und Schwerpunkte

JDK-Lab22 Java-Dateien1 Tests/RunnerTransactional OutboxIdempotenzRetry
Echter Quellcode aus diesem Lab

Geführter Codepfad

Hier wird der Lebenszyklus eines Outbox-Events bis zum Consumer verfolgt. Entscheidend sind Absturzfenster, erneute Zustellung und die idempotente Rechnungserzeugung.

Run4C3DemoPlaceOrderServiceInMemoryOutboxRepositoryOutboxPublisherFakeBrokerBillingConsumerRun4C3Tests
Lesereihenfolge der zentralen Klassen. Die Pfeile zeigen den didaktischen Weg durch den realen Quellcode, nicht zwingend jeden Laufzeitaufruf.
1. Run4C3DemoVerdrahtet Order-Service, Outbox, Broker und Billing-Consumer zu einem ausführbaren End-to-End-Szenario.
2. PlaceOrderServiceSpeichert die Bestellung zusammen mit der zugehörigen Outbox-Nachricht.
3. InMemoryOutboxRepositoryHält Status und Zustellzustand der Outbox-Einträge über mehrere Publisher-Läufe hinweg fest.
4. OutboxPublisherVeröffentlicht offene Nachrichten und behandelt Brokerfehler sowie erneute Versuche.
5. FakeBrokerSimuliert erfolgreiche und fehlerhafte Broker-Zustellungen für reproduzierbare Crash-Szenarien.
6. BillingConsumerErzeugt Rechnungen und schützt die Verarbeitung mit einem IdempotencyStore vor Duplikaten.
7. Run4C3TestsPrüft Crash-Recovery, Mehrfachzustellung und genau-einmalige fachliche Wirkung.

1. Run4C3Demo

src/main/java/com/example/outboxdeepdive/demo/Run4C3Demo.java
Java-Datei öffnen
Rolle im Ablauf

Verdrahtet Order-Service, Outbox, Broker und Billing-Consumer zu einem ausführbaren End-to-End-Szenario.

Im Lesepfad folgt PlaceOrderService: Speichert die Bestellung zusammen mit der zugehörigen Outbox-Nachricht.

Typ
class Run4C3Demo
Verwendet
PlaceOrderService, InMemoryOutboxRepository, OutboxPublisher, FakeBroker, BillingConsumer
Verwendet von
Einstiege
main(String[] args)
package com.example.outboxdeepdive.demo;
import com.example.outboxdeepdive.billing.*;
import com.example.outboxdeepdive.broker.*;
import com.example.outboxdeepdive.order.*;
import com.example.outboxdeepdive.outbox.*;
import com.example.outboxdeepdive.shared.*;
import java.time.Instant;

public final class Run4C3Demo {
    public static void main(String[] args) {
        var orders = new InMemoryOrderRepository();
        var outbox = new InMemoryOutboxRepository();
        var broker = new FakeBroker();
        var service = new PlaceOrderService(orders, outbox);
        var publisher = new OutboxPublisher(outbox, broker, 10, 3);
        Instant now = Instant.parse("2026-01-01T10:00:00Z");
        service.placeOrder("O-100", Money.eur("129.90"), now);
        broker.failNextPublishes(1);
        publisher.publishDue(now); // failed, retry scheduled
        publisher.publishDue(now.plusSeconds(5)); // still not due
        publisher.publishDue(now.plusSeconds(20)); // sent
        var idem = new InMemoryIdempotencyStore();
        var invoices = new InvoiceRepository();
        var consumer = new BillingConsumer(idem, invoices);
        var event = broker.published().getFirst();
        consumer.onMessage(event);
        consumer.onMessage(event); // duplicate delivery must not create second invoice
        System.out.println("RUN4C3_DEMO_OK orders=" + orders.size() + " outbox=" + outbox.all().size() + " broker=" + broker.published().size() + " invoices=" + invoices.invoices().size());
    }
}

2. PlaceOrderService

src/main/java/com/example/outboxdeepdive/order/PlaceOrderService.java
Java-Datei öffnen
Rolle im Ablauf

Speichert die Bestellung zusammen mit der zugehörigen Outbox-Nachricht.

Im Lesepfad folgt InMemoryOutboxRepository: Hält Status und Zustellzustand der Outbox-Einträge über mehrere Publisher-Läufe hinweg fest.

Typ
class PlaceOrderService
Verwendet
Verwendet von
Run4C3Demo, Run4C3Tests
Einstiege
placeOrder(String id, Money total, Instant now)
Application Service
package com.example.outboxdeepdive.order;
import com.example.outboxdeepdive.outbox.*;
import com.example.outboxdeepdive.shared.*;
import java.time.Instant;

// Pattern: Application Service - definiert Use-Case- und Transaktionsgrenze.
public final class PlaceOrderService {
    private final OrderRepository orders;
    private final OutboxRepository outbox;
    public PlaceOrderService(OrderRepository orders, OutboxRepository outbox) { this.orders = orders; this.outbox = outbox; }

    public OrderId placeOrder(String id, Money total, Instant now) {
        Order order = Order.accept(OrderId.of(id), total, now);
        // In einer echten Datenbank waere dies eine Transaktion:
        // 1) Order speichern
        // 2) OutboxMessage speichern
        // 3) Commit
        orders.save(order);
        outbox.save(OutboxMessage.pending("Order", order.id().value(), "OrderAccepted", 1, payload(order), now));
        return order.id();
    }
    private String payload(Order order) {
        return "{"orderId":"" + order.id().value() + "","amount":"" + order.total().amount() + "","currency":"" + order.total().currency() + ""}";
    }
}

3. InMemoryOutboxRepository

src/main/java/com/example/outboxdeepdive/outbox/InMemoryOutboxRepository.java
Java-Datei öffnen
Rolle im Ablauf

Hält Status und Zustellzustand der Outbox-Einträge über mehrere Publisher-Läufe hinweg fest.

Im Lesepfad folgt OutboxPublisher: Veröffentlicht offene Nachrichten und behandelt Brokerfehler sowie erneute Versuche.

Typ
class InMemoryOutboxRepositoryimplements OutboxRepository
Verwendet
Verwendet von
Run4C3Demo, Run4C3Tests
Einstiege
save(OutboxMessage message), findDue(Instant now, int limit), findById(UUID id), all()
package com.example.outboxdeepdive.outbox;
import java.time.Instant;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;

public final class InMemoryOutboxRepository implements OutboxRepository {
    private final Map<UUID, OutboxMessage> messages = new ConcurrentHashMap<>();
    public void save(OutboxMessage message) { messages.put(message.id(), message); }
    public List<OutboxMessage> findDue(Instant now, int limit) {
        return messages.values().stream()
                .filter(m -> m.dueAt(now))
                .sorted(Comparator.comparing(OutboxMessage::createdAt))
                .limit(limit)
                .toList();
    }
    public Optional<OutboxMessage> findById(UUID id) { return Optional.ofNullable(messages.get(id)); }
    public List<OutboxMessage> all() { return messages.values().stream().sorted(Comparator.comparing(OutboxMessage::createdAt)).toList(); }
}

4. OutboxPublisher

src/main/java/com/example/outboxdeepdive/outbox/OutboxPublisher.java
Java-Datei öffnen
Rolle im Ablauf

Veröffentlicht offene Nachrichten und behandelt Brokerfehler sowie erneute Versuche.

Im Lesepfad folgt FakeBroker: Simuliert erfolgreiche und fehlerhafte Broker-Zustellungen für reproduzierbare Crash-Szenarien.

Typ
class OutboxPublisher
Verwendet
Verwendet von
Run4C3Demo, Run4C3Tests
Einstiege
publishDue(Instant now), recoverInFlightAfterCrash()
Publisher
package com.example.outboxdeepdive.outbox;
import com.example.outboxdeepdive.broker.*;
import java.time.*;

// Pattern: Publisher - entkoppelt DB-Transaktion und Broker-Versand.
public final class OutboxPublisher {
    private final OutboxRepository repository;
    private final BrokerPort broker;
    private final int batchSize;
    private final int maxAttempts;

    public OutboxPublisher(OutboxRepository repository, BrokerPort broker, int batchSize, int maxAttempts) {
        this.repository = repository; this.broker = broker; this.batchSize = batchSize; this.maxAttempts = maxAttempts;
    }
    public int publishDue(Instant now) {
        int sent = 0;
        for (OutboxMessage message : repository.findDue(now, batchSize)) {
            message.claim();
            try {
                broker.publish(new EventEnvelope(message.id(), message.eventType(), message.eventVersion(), message.payload()));
                message.markSent();
                sent++;
            } catch (RuntimeException ex) {
                message.markFailed(ex.getClass().getSimpleName() + ": " + ex.getMessage(), nextAttempt(now, message.attempts() + 1), maxAttempts);
            }
        }
        return sent;
    }
    private Instant nextAttempt(Instant now, int attempt) {
        long seconds = Math.min(300, (long)Math.pow(2, Math.max(1, attempt)));
        return now.plusSeconds(seconds);
    }
    public void recoverInFlightAfterCrash() {
        repository.all().forEach(OutboxMessage::requeueAfterPublisherCrash);
    }
}

5. FakeBroker

src/main/java/com/example/outboxdeepdive/broker/FakeBroker.java
Java-Datei öffnen
Rolle im Ablauf

Simuliert erfolgreiche und fehlerhafte Broker-Zustellungen für reproduzierbare Crash-Szenarien.

Im Lesepfad folgt BillingConsumer: Erzeugt Rechnungen und schützt die Verarbeitung mit einem IdempotencyStore vor Duplikaten.

Typ
class FakeBrokerimplements BrokerPort
Verwendet
Verwendet von
Run4C3Demo, Run4C3Tests
Einstiege
failNextPublishes(int count), publish(EventEnvelope envelope), published()
Test Fake
package com.example.outboxdeepdive.broker;
import java.util.*;

// Pattern: Test Fake - kontrollierbarer Broker fuer Crash- und Retry-Tests.
public final class FakeBroker implements BrokerPort {
    private final List<EventEnvelope> published = new ArrayList<>();
    private int failNextPublishes;
    public void failNextPublishes(int count) { this.failNextPublishes = count; }
    public void publish(EventEnvelope envelope) {
        if (failNextPublishes > 0) { failNextPublishes--; throw new BrokerUnavailableException("broker temporarily unavailable"); }
        published.add(envelope);
    }
    public List<EventEnvelope> published() { return List.copyOf(published); }
}

6. BillingConsumer

src/main/java/com/example/outboxdeepdive/billing/BillingConsumer.java
Java-Datei öffnen
Rolle im Ablauf

Erzeugt Rechnungen und schützt die Verarbeitung mit einem IdempotencyStore vor Duplikaten.

Im Lesepfad folgt Run4C3Tests: Prüft Crash-Recovery, Mehrfachzustellung und genau-einmalige fachliche Wirkung.

Typ
class BillingConsumer
Verwendet
Verwendet von
Run4C3Demo, Run4C3Tests
Einstiege
onMessage(EventEnvelope envelope)
Idempotent Consumer
package com.example.outboxdeepdive.billing;
import com.example.outboxdeepdive.broker.EventEnvelope;

// Pattern: Idempotent Consumer - verarbeitet dasselbe Event mehrfach ohne doppelte Rechnung.
public final class BillingConsumer {
    private final IdempotencyStore idempotencyStore;
    private final InvoiceRepository invoices;
    public BillingConsumer(IdempotencyStore idempotencyStore, InvoiceRepository invoices) { this.idempotencyStore = idempotencyStore; this.invoices = invoices; }

    public void onMessage(EventEnvelope envelope) {
        if (idempotencyStore.alreadyProcessed(envelope.messageId())) return;
        String orderId = extractOrderId(envelope.payload());
        invoices.createInvoice(orderId);
        idempotencyStore.markProcessed(envelope.messageId());
    }
    private String extractOrderId(String json) {
        String marker = ""orderId":"";
        int start = json.indexOf(marker);
        if (start < 0) throw new IllegalArgumentException("orderId missing");
        start += marker.length();
        int end = json.indexOf('"', start);
        return json.substring(start, end);
    }
}

7. Run4C3Tests

src/test/java/com/example/outboxdeepdive/Run4C3Tests.java
Java-Datei öffnen
Rolle im Ablauf

Prüft Crash-Recovery, Mehrfachzustellung und genau-einmalige fachliche Wirkung.

Damit ist der zentrale Pfad abgeschlossen; der Test-/Runner-Code und die vollständige Dateiliste darunter zeigen die übrigen Varianten.

Typ
class Run4C3Tests
Verwendet
PlaceOrderService, InMemoryOutboxRepository, OutboxPublisher, FakeBroker, BillingConsumer
Verwendet von
Einstiege
main(String[] args)
package com.example.outboxdeepdive;
import com.example.outboxdeepdive.billing.*;
import com.example.outboxdeepdive.broker.*;
import com.example.outboxdeepdive.order.*;
import com.example.outboxdeepdive.outbox.*;
import com.example.outboxdeepdive.shared.*;
import java.time.Instant;

public final class Run4C3Tests {
    public static void main(String[] args) {
        orderAndOutboxAreWrittenTogether();
        brokerFailureKeepsMessageRetryable();
        duplicateDeliveryCreatesOnlyOneInvoice();
        maxAttemptsMovesToDeadLetter();
        System.out.println("RUN4C3_TESTS_OK");
    }
    static void orderAndOutboxAreWrittenTogether() {
        var orders = new InMemoryOrderRepository(); var outbox = new InMemoryOutboxRepository();
        new PlaceOrderService(orders, outbox).placeOrder("O-1", Money.eur("10.00"), Instant.parse("2026-01-01T00:00:00Z"));
        check(orders.size() == 1, "order saved"); check(outbox.all().size() == 1, "outbox saved");
        check(outbox.all().getFirst().status() == OutboxStatus.PENDING, "pending");
    }
    static void brokerFailureKeepsMessageRetryable() {
        var outbox = new InMemoryOutboxRepository(); var broker = new FakeBroker();
        outbox.save(OutboxMessage.pending("Order", "O-2", "OrderAccepted", 1, "{"orderId":"O-2"}", Instant.parse("2026-01-01T00:00:00Z")));
        broker.failNextPublishes(1);
        var publisher = new OutboxPublisher(outbox, broker, 10, 3);
        publisher.publishDue(Instant.parse("2026-01-01T00:00:00Z"));
        check(outbox.all().getFirst().status() == OutboxStatus.FAILED, "failed after broker error");
        publisher.publishDue(Instant.parse("2026-01-01T00:00:20Z"));
        check(outbox.all().getFirst().status() == OutboxStatus.SENT, "sent after retry");
    }
    static void duplicateDeliveryCreatesOnlyOneInvoice() {
        var idem = new InMemoryIdempotencyStore(); var invoices = new InvoiceRepository(); var consumer = new BillingConsumer(idem, invoices);
        var event = new EventEnvelope(java.util.UUID.randomUUID(), "OrderAccepted", 1, "{"orderId":"O-3"}");
        consumer.onMessage(event); consumer.onMessage(event);
        check(invoices.invoices().size() == 1, "only one invoice");
    }
    static void maxAttemptsMovesToDeadLetter() {
        var outbox = new InMemoryOutboxRepository(); var broker = new FakeBroker(); broker.failNextPublishes(5);
        outbox.save(OutboxMessage.pending("Order", "O-4", "OrderAccepted", 1, "{"orderId":"O-4"}", Instant.parse("2026-01-01T00:00:00Z")));
        var publisher = new OutboxPublisher(outbox, broker, 10, 2);
        publisher.publishDue(Instant.parse("2026-01-01T00:00:00Z"));
        publisher.publishDue(Instant.parse("2026-01-01T00:01:00Z"));
        check(outbox.all().getFirst().status() == OutboxStatus.DEAD_LETTER, "dead letter after max attempts");
    }
    static void check(boolean condition, String message) { if (!condition) throw new AssertionError(message); }
}
Alle Projektdateien öffnen (24 Einträge)
⌂ Cockpit