1. MessagingRun6BDemo
src/main/java/com/example/messaging/run6b/MessagingRun6BDemo.javaFührt transiente und permanente Fehler durch den vollständigen Fehlerpfad.
- 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());
}
}