Messaging-Retry-und-DLQ-Lab

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

Zurück zu Code-Labs

Kapitel 09 · Messaging-Fehlerbetrieb

Was dieses Lab zeigt

Behandelt den Fehlerpfad als Teil des Designs: transiente und permanente Fehler, Retry-Entscheidungen, Dead-Letter Queue, Quarantäne und kontrolliertes Reprocessing.

Lernziele

  • Fehler klassifizieren
  • Retry-Budgets anwenden
  • DLQ und Replay sicher betreiben

Technik und Schwerpunkte

JDK-Lab28 Java-Dateien1 Tests/RunnerRetryDLQReprocessing
Echter Quellcode aus diesem Lab

Geführter Codepfad

Das Fehlerbetriebs-Lab verfolgt einen Record von der Verarbeitung über Fehlerklassifikation und Retry bis zu DLQ, Parking Lot und gezieltem Reprocessing. Lag und Metriken liefern die betriebliche Sicht.

MessagingRun6BDemoReliableConsumerExecutorDefaultErrorClassifierExponentialBackoffRetryPolicyRetrySchedulerDeadLetterQueueReprocessingServiceRunbookAdvisor
Lesereihenfolge der zentralen Klassen. Die Pfeile zeigen den didaktischen Weg durch den realen Quellcode, nicht zwingend jeden Laufzeitaufruf.
1. MessagingRun6BDemoFührt transiente und permanente Fehler durch den vollständigen Fehlerpfad.
2. ReliableConsumerExecutorKapselt Idempotenz, Consumer-Aufruf, Fehlerklassifikation und die resultierende Retry- oder DLQ-Entscheidung.
3. DefaultErrorClassifierUnterscheidet Fehlerarten, damit nicht jeder Fehler identisch erneut versucht wird.
4. ExponentialBackoffRetryPolicyBerechnet Retry-Zeitpunkte mit wachsendem Abstand und begrenzter Versuchszahl.
5. RetrySchedulerHält fällige Wiederholungen und übergibt sie zum richtigen Zeitpunkt erneut an die Verarbeitung.
6. DeadLetterQueueBewahrt endgültig fehlgeschlagene Records mitsamt Fehlerkontext auf.
7. ReprocessingServiceSteuert die kontrollierte Rückführung von DLQ- oder Parking-Lot-Einträgen.
8. RunbookAdvisorLeitet aus Lag, Fehlern und Retry-Zuständen konkrete Betriebsmaßnahmen ab.

1. MessagingRun6BDemo

src/main/java/com/example/messaging/run6b/MessagingRun6BDemo.java
Java-Datei öffnen
Rolle im Ablauf

Führt transiente und permanente Fehler durch den vollständigen Fehlerpfad.

Im Lesepfad folgt ReliableConsumerExecutor: Kapselt Idempotenz, Consumer-Aufruf, Fehlerklassifikation und die resultierende Retry- oder DLQ-Entscheidung.

Typ
class MessagingRun6BDemo
Verwendet
ReliableConsumerExecutor, DefaultErrorClassifier, ExponentialBackoffRetryPolicy, RetryScheduler, DeadLetterQueue
Verwendet von
Einstiege
main(String[] args)
package com.example.messaging.run6b;

import java.time.*;

public final class MessagingRun6BDemo {
    public static void main(String[] args) {
        InMemoryTopic topic = new InMemoryTopic("order-events", 3);
        ConsumerGroupOffsets offsets = new ConsumerGroupOffsets();
        FlakyBillingProjectionConsumer consumer = new FlakyBillingProjectionConsumer();
        consumer.failTransientTimes("order-200", 1);
        consumer.failFinalDomain("order-404");
        consumer.poison("order-666");
        RetryScheduler retry = new RetryScheduler();
        DeadLetterQueue dlq = new DeadLetterQueue();
        ParkingLot parking = new ParkingLot();
        ProcessingMetrics metrics = new ProcessingMetrics();
        ReliableConsumerExecutor executor = new ReliableConsumerExecutor(consumer, new IdempotencyStore(),
                new ExponentialBackoffRetryPolicy(2, Duration.ofSeconds(10)), retry, dlq, parking,
                new DefaultErrorClassifier(), metrics);
        topic.append("order-100", EventEnvelope.newOrderEvent("order-100", "{total:42}"));
        topic.append("order-200", EventEnvelope.newOrderEvent("order-200", "{total:99}"));
        topic.append("order-404", EventEnvelope.newOrderEvent("order-404", "{total:11}"));
        topic.append("order-666", EventEnvelope.newOrderEvent("order-666", "{total:13}"));
        BatchPoller poller = new BatchPoller(topic, offsets, executor);
        poller.pollOnce(10, Instant.parse("2026-01-01T10:00:00Z"));
        for (EventEnvelope due : retry.due(Instant.parse("2026-01-01T10:01:00Z"))) executor.process(due, Instant.parse("2026-01-01T10:01:00Z"));
        long lag = new ConsumerLagMonitor().totalLag(topic, offsets);
        System.out.println("RUN6B_DEMO_OK metrics=" + metrics + " lag=" + lag + " dlq=" + dlq.size() + " parking=" + parking.size());
    }
}

2. ReliableConsumerExecutor

src/main/java/com/example/messaging/run6b/ReliableConsumerExecutor.java
Java-Datei öffnen
Rolle im Ablauf

Kapselt Idempotenz, Consumer-Aufruf, Fehlerklassifikation und die resultierende Retry- oder DLQ-Entscheidung.

Im Lesepfad folgt DefaultErrorClassifier: Unterscheidet Fehlerarten, damit nicht jeder Fehler identisch erneut versucht wird.

Typ
class ReliableConsumerExecutor
Verwendet
RetryScheduler, DeadLetterQueue
Verwendet von
MessagingRun6BDemo
Einstiege
process(EventEnvelope event, Instant now)
Facade
package com.example.messaging.run6b;

import java.time.Instant;

// Pattern: Facade - buendelt Retry, DLQ, Idempotenz und Metrics um einen fachlichen Consumer.
public final class ReliableConsumerExecutor {
    private final Consumer consumer;
    private final IdempotencyStore idempotency;
    private final RetryPolicy retryPolicy;
    private final RetryScheduler retryScheduler;
    private final DeadLetterQueue dlq;
    private final ParkingLot parkingLot;
    private final ErrorClassifier classifier;
    private final ProcessingMetrics metrics;
    public ReliableConsumerExecutor(Consumer consumer, IdempotencyStore idempotency, RetryPolicy retryPolicy,
                                    RetryScheduler retryScheduler, DeadLetterQueue dlq, ParkingLot parkingLot,
                                    ErrorClassifier classifier, ProcessingMetrics metrics) {
        this.consumer = consumer; this.idempotency = idempotency; this.retryPolicy = retryPolicy;
        this.retryScheduler = retryScheduler; this.dlq = dlq; this.parkingLot = parkingLot;
        this.classifier = classifier; this.metrics = metrics;
    }
    public ConsumerResult process(EventEnvelope event, Instant now) {
        if (idempotency.alreadyProcessed(event.eventId())) {
            metrics.duplicate();
            return ConsumerResult.duplicate(event.eventId());
        }
        ConsumerResult result;
        try { result = consumer.handle(event); }
        catch (RuntimeException ex) { result = classifier.classify(ex, event); }
        if (result instanceof ConsumerResult.Success) {
            idempotency.markProcessed(event.eventId()); metrics.success(); return result;
        }
        if (result instanceof ConsumerResult.Failure failure) {
            RetryDecision decision = retryPolicy.decide(event, failure);
            if (decision.retry()) {
                retryScheduler.schedule(event, decision.delay(), decision.reason(), now);
                metrics.retryScheduled();
                return failure;
            }
            dlq.add(event, failure, now); metrics.dlq();
            if (failure.kind() == FailureKind.POISON_MESSAGE || failure.kind() == FailureKind.BUG) {
                parkingLot.park(dlq.records().getLast()); metrics.parked();
            }
            return failure;
        }
        return result;
    }
}

3. DefaultErrorClassifier

src/main/java/com/example/messaging/run6b/DefaultErrorClassifier.java
Java-Datei öffnen
Rolle im Ablauf

Unterscheidet Fehlerarten, damit nicht jeder Fehler identisch erneut versucht wird.

Im Lesepfad folgt ExponentialBackoffRetryPolicy: Berechnet Retry-Zeitpunkte mit wachsendem Abstand und begrenzter Versuchszahl.

Typ
class DefaultErrorClassifierimplements ErrorClassifier
Verwendet
Verwendet von
MessagingRun6BDemo
Einstiege
package com.example.messaging.run6b;

public final class DefaultErrorClassifier implements ErrorClassifier {
    @Override public ConsumerResult.Failure classify(RuntimeException exception, EventEnvelope event) {
        String m = exception.getMessage() == null ? exception.getClass().getSimpleName() : exception.getMessage();
        if (m.contains("timeout") || m.contains("temporarily")) return ConsumerResult.transientFailure(m);
        if (m.contains("unknown customer")) return ConsumerResult.temporaryDomain(m);
        if (m.contains("invalid schema")) return ConsumerResult.poison(m);
        return ConsumerResult.bug(m);
    }
}

4. ExponentialBackoffRetryPolicy

src/main/java/com/example/messaging/run6b/ExponentialBackoffRetryPolicy.java
Java-Datei öffnen
Rolle im Ablauf

Berechnet Retry-Zeitpunkte mit wachsendem Abstand und begrenzter Versuchszahl.

Im Lesepfad folgt RetryScheduler: Hält fällige Wiederholungen und übergibt sie zum richtigen Zeitpunkt erneut an die Verarbeitung.

Typ
class ExponentialBackoffRetryPolicyimplements RetryPolicy
Verwendet
Verwendet von
MessagingRun6BDemo
Einstiege
Strategy
package com.example.messaging.run6b;

import java.time.Duration;

// Pattern: Strategy - kapselt Retry-Regeln, ohne Consumer-Code zu veraendern.
public final class ExponentialBackoffRetryPolicy implements RetryPolicy {
    private final int maxAttempts;
    private final Duration baseDelay;
    public ExponentialBackoffRetryPolicy(int maxAttempts, Duration baseDelay) {
        this.maxAttempts = maxAttempts; this.baseDelay = baseDelay;
    }
    @Override public RetryDecision decide(EventEnvelope event, ConsumerResult.Failure failure) {
        if (!failure.retryable()) return RetryDecision.stop("failure is not retryable: " + failure.kind());
        if (event.attempt() >= maxAttempts) return RetryDecision.stop("max attempts reached: " + event.attempt());
        long factor = 1L << Math.min(event.attempt(), 8);
        return RetryDecision.retryAfter(baseDelay.multipliedBy(factor), "retryable " + failure.kind());
    }
}

5. RetryScheduler

src/main/java/com/example/messaging/run6b/RetryScheduler.java
Java-Datei öffnen
Rolle im Ablauf

Hält fällige Wiederholungen und übergibt sie zum richtigen Zeitpunkt erneut an die Verarbeitung.

Im Lesepfad folgt DeadLetterQueue: Bewahrt endgültig fehlgeschlagene Records mitsamt Fehlerkontext auf.

Typ
class RetryScheduler
Verwendet
Verwendet von
MessagingRun6BDemo, ReliableConsumerExecutor
Einstiege
schedule(EventEnvelope event, Duration delay, String reason, Instant now), due(Instant now), size(), entries()
Scheduler
package com.example.messaging.run6b;

import java.time.*;
import java.util.*;

// Pattern: Scheduler - entkoppelt Retry-Zeitpunkt vom fachlichen Consumer.
public final class RetryScheduler {
    private final List<RetryEntry> entries = new ArrayList<>();
    public void schedule(EventEnvelope event, Duration delay, String reason, Instant now) {
        entries.add(new RetryEntry(event.nextAttempt(), now.plus(delay), reason));
    }
    public List<EventEnvelope> due(Instant now) {
        List<EventEnvelope> due = new ArrayList<>();
        Iterator<RetryEntry> it = entries.iterator();
        while (it.hasNext()) {
            RetryEntry entry = it.next();
            if (!entry.dueAt().isAfter(now)) { due.add(entry.event()); it.remove(); }
        }
        return due;
    }
    public int size() { return entries.size(); }
    public List<RetryEntry> entries() { return List.copyOf(entries); }
}

6. DeadLetterQueue

src/main/java/com/example/messaging/run6b/DeadLetterQueue.java
Java-Datei öffnen
Rolle im Ablauf

Bewahrt endgültig fehlgeschlagene Records mitsamt Fehlerkontext auf.

Im Lesepfad folgt ReprocessingService: Steuert die kontrollierte Rückführung von DLQ- oder Parking-Lot-Einträgen.

Typ
class DeadLetterQueue
Verwendet
Verwendet von
MessagingRun6BDemo, ReliableConsumerExecutor, ReprocessingService
Einstiege
add(EventEnvelope event, ConsumerResult.Failure failure, Instant now), records(), removeFirst(), size()
Dead Letter Channel
package com.example.messaging.run6b;

import java.time.Instant;
import java.util.*;

// Pattern: Dead Letter Channel - isoliert nicht verarbeitbare Nachrichten mit Diagnosekontext.
public final class DeadLetterQueue {
    private final List<DeadLetterRecord> records = new ArrayList<>();
    public void add(EventEnvelope event, ConsumerResult.Failure failure, Instant now) {
        records.add(new DeadLetterRecord(event, failure.kind(), failure.message(), now));
    }
    public List<DeadLetterRecord> records() { return List.copyOf(records); }
    public DeadLetterRecord removeFirst() { return records.removeFirst(); }
    public int size() { return records.size(); }
}

7. ReprocessingService

src/main/java/com/example/messaging/run6b/ReprocessingService.java
Java-Datei öffnen
Rolle im Ablauf

Steuert die kontrollierte Rückführung von DLQ- oder Parking-Lot-Einträgen.

Im Lesepfad folgt RunbookAdvisor: Leitet aus Lag, Fehlern und Retry-Zuständen konkrete Betriebsmaßnahmen ab.

Typ
class ReprocessingService
Verwendet
DeadLetterQueue
Verwendet von
Einstiege
replayFromDlq(DeadLetterQueue dlq, int maxRecords, String operator)
Process Manager
package com.example.messaging.run6b;

import java.time.Instant;
import java.util.*;

// Pattern: Process Manager - steuert manuelles Reprocessing mit Limit und Audit-Hinweis.
public final class ReprocessingService {
    private final InMemoryTopic topic;
    public ReprocessingService(InMemoryTopic topic) { this.topic = topic; }
    public List<MessageRecord> replayFromDlq(DeadLetterQueue dlq, int maxRecords, String operator) {
        List<MessageRecord> requeued = new ArrayList<>();
        int count = Math.min(maxRecords, dlq.size());
        for (int i=0; i<count; i++) {
            DeadLetterRecord record = dlq.removeFirst();
            EventEnvelope e = record.event();
            EventEnvelope replay = new EventEnvelope(e.eventId(), e.eventType(), e.aggregateId(), e.schemaVersion(), e.occurredAt(),
                    new java.util.LinkedHashMap<>(e.headers()) {{ put("reprocessedBy", operator); put("reprocessedAt", Instant.now().toString()); }},
                    e.payload(), 0);
            requeued.add(topic.append(replay.aggregateId(), replay));
        }
        return requeued;
    }
}

8. RunbookAdvisor

src/main/java/com/example/messaging/run6b/RunbookAdvisor.java
Java-Datei öffnen
Rolle im Ablauf

Leitet aus Lag, Fehlern und Retry-Zuständen konkrete Betriebsmaßnahmen ab.

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

Typ
class RunbookAdvisor
Verwendet
Verwendet von
Einstiege
advise(long lag, int dlqSize, int retrySize)
package com.example.messaging.run6b;

public final class RunbookAdvisor {
    public String advise(long lag, int dlqSize, int retrySize) {
        if (dlqSize > 0 && retrySize > 0) return "DLQ und Retry gleichzeitig: Fehlerklassifikation pruefen, nicht blind skalieren.";
        if (lag > 1000) return "Lag hoch: Downstream-Latenz, Partition-Hotspot und Consumer-Rebalancing pruefen.";
        if (dlqSize > 0) return "DLQ vorhanden: Samples lesen, Fehlerklasse bestimmen, Reprocessing erst nach Fix.";
        if (retrySize > 100) return "Retry-Sturm: Backoff vergroessern, Rate limitieren, Circuit Breaker pruefen.";
        return "System stabil: SLO, Lag und DLQ weiter beobachten.";
    }
}
Alle Projektdateien öffnen (31 Einträge)
⌂ Cockpit