Messaging-Grundlagen-Lab

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

Zurück zu Code-Labs

Kapitel 08 · Messaging-Grundlagen

Was dieses Lab zeigt

Erklärt Messaging als Vertrag zwischen Publisher und Consumer. Commands, Events, Event-Envelopes, Topic-Zuordnung und Schema-Entscheidungen werden am Order-Fluss greifbar.

Lernziele

  • Commands und Events unterscheiden
  • Event-Verträge sauber modellieren
  • Publisher und Consumer entkoppeln

Technik und Schwerpunkte

Java 2131 Java-Dateien1 Tests/RunnerEventsEvent EnvelopeSchema
Echter Quellcode aus diesem Lab

Geführter Codepfad

Der Messaging-Grundpfad beginnt bei typisierten Envelopes und Schema-Prüfung, führt durch Topic und Partition und endet in einer idempotenten Billing-Projektion.

MessagingRun6ARebuildDemoEventSchemaRegistryPartitionKeyStrategyInMemoryKafkaBrokerBillingProjectionConsumerIdempotencyStoreMessagingRun6ARebuildTestRunner
Lesereihenfolge der zentralen Klassen. Die Pfeile zeigen den didaktischen Weg durch den realen Quellcode, nicht zwingend jeden Laufzeitaufruf.
1. MessagingRun6ARebuildDemoErzeugt Commands und Events, publiziert sie und lässt den Consumer die Projektion aktualisieren.
2. EventSchemaRegistryRegistriert bekannte Payload-Typen und prüft die Kompatibilität von Event-Schemas.
3. PartitionKeyStrategyBestimmt den fachlichen Schlüssel, über den zusammengehörige Nachrichten in derselben Partition landen.
4. InMemoryKafkaBrokerSimuliert Topics, Partitionen, Consumer-Gruppen und Records des Messaging-Ablaufs.
5. BillingProjectionConsumerVerarbeitet Order-Events und aktualisiert die BillingProjection.
6. IdempotencyStoreMerkt bereits verarbeitete Nachrichten und verhindert doppelte fachliche Wirkung.
7. MessagingRun6ARebuildTestRunnerPrüft Schema, Partitionierung, Reihenfolge, Consumer-Gruppen und Idempotenz.

1. MessagingRun6ARebuildDemo

src/main/java/com/example/messaging/run6a/demo/MessagingRun6ARebuildDemo.java
Java-Datei öffnen
Rolle im Ablauf

Erzeugt Commands und Events, publiziert sie und lässt den Consumer die Projektion aktualisieren.

Im Lesepfad folgt EventSchemaRegistry: Registriert bekannte Payload-Typen und prüft die Kompatibilität von Event-Schemas.

Typ
class MessagingRun6ARebuildDemo
Verwendet
EventSchemaRegistry, InMemoryKafkaBroker, BillingProjectionConsumer, IdempotencyStore
Verwendet von
Einstiege
main(String[] args)
package com.example.messaging.run6a.demo;

import com.example.messaging.run6a.broker.InMemoryKafkaBroker;
import com.example.messaging.run6a.consumer.BillingProjectionConsumer;
import com.example.messaging.run6a.consumer.DeadLetterQueue;
import com.example.messaging.run6a.consumer.IdempotencyStore;
import com.example.messaging.run6a.model.EventFactory;
import com.example.messaging.run6a.schema.EventSchemaRegistry;

public final class MessagingRun6ARebuildDemo {
    public static void main(String[] args) {
        InMemoryKafkaBroker broker = new InMemoryKafkaBroker();
        broker.createTopic("order-events", 3);
        EventFactory factory = new EventFactory();
        broker.publish("order-events", factory.orderPlaced("trace-1", "ORD-100", "C-1", 12_990, "EUR", 2));
        broker.publish("order-events", factory.orderPlaced("trace-2", "ORD-101", "C-2", 7_500, "EUR", 2));

        BillingProjectionConsumer billing = new BillingProjectionConsumer(
                new EventSchemaRegistry().register("OrderPlaced", 1, 2),
                new IdempotencyStore(),
                new DeadLetterQueue());

        broker.poll("order-events", "billing-group", 10).forEach(record -> {
            billing.handle(record.event());
            broker.commit(record.topic(), "billing-group", record.partition(), record.offset() + 1);
        });
        System.out.println("RUN6A_REBUILD_DEMO_OK projectionCount=" + billing.projectionCount() + " lag=" + broker.lag("order-events", "billing-group"));
    }
}

2. EventSchemaRegistry

src/main/java/com/example/messaging/run6a/schema/EventSchemaRegistry.java
Java-Datei öffnen
Rolle im Ablauf

Registriert bekannte Payload-Typen und prüft die Kompatibilität von Event-Schemas.

Im Lesepfad folgt PartitionKeyStrategy: Bestimmt den fachlichen Schlüssel, über den zusammengehörige Nachrichten in derselben Partition landen.

Typ
class EventSchemaRegistry
Verwendet
Verwendet von
MessagingRun6ARebuildDemo, BillingProjectionConsumer, MessagingRun6ARebuildTestRunner
Einstiege
register(String eventType, int... versions), validate(EventEnvelope event)
Registry
package com.example.messaging.run6a.schema;

import com.example.messaging.run6a.model.EventEnvelope;
import java.util.HashMap;
import java.util.HashSet;
import java.util.Map;
import java.util.Set;

// Pattern: Registry - zentralisiert akzeptierte Event-Versionen pro Consumer.
public final class EventSchemaRegistry {
    private final Map<String, Set<Integer>> acceptedVersions = new HashMap<>();

    public EventSchemaRegistry register(String eventType, int... versions) {
        Set<Integer> set = acceptedVersions.computeIfAbsent(eventType, ignored -> new HashSet<>());
        for (int version : versions) set.add(version);
        return this;
    }

    public CompatibilityCheck validate(EventEnvelope event) {
        Set<Integer> versions = acceptedVersions.get(event.eventType());
        if (versions == null) return CompatibilityCheck.unknownType(event.eventType());
        if (!versions.contains(event.schemaVersion())) return CompatibilityCheck.rejectedVersion(event.eventType(), event.schemaVersion());
        return CompatibilityCheck.accepted(event.eventType(), event.schemaVersion());
    }
}

3. PartitionKeyStrategy

src/main/java/com/example/messaging/run6a/broker/PartitionKeyStrategy.java
Java-Datei öffnen
Rolle im Ablauf

Bestimmt den fachlichen Schlüssel, über den zusammengehörige Nachrichten in derselben Partition landen.

Im Lesepfad folgt InMemoryKafkaBroker: Simuliert Topics, Partitionen, Consumer-Gruppen und Records des Messaging-Ablaufs.

Typ
class PartitionKeyStrategy
Verwendet
Verwendet von
InMemoryKafkaBroker
Einstiege
partitionFor(String key, int partitionCount)
Strategy
package com.example.messaging.run6a.broker;

// Pattern: Strategy - kapselt die Zuordnung von Key zu Partition.
public final class PartitionKeyStrategy {
    public int partitionFor(String key, int partitionCount) {
        if (partitionCount <= 0) throw new IllegalArgumentException("partitionCount");
        return Math.floorMod(key.hashCode(), partitionCount);
    }
}

4. InMemoryKafkaBroker

src/main/java/com/example/messaging/run6a/broker/InMemoryKafkaBroker.java
Java-Datei öffnen
Rolle im Ablauf

Simuliert Topics, Partitionen, Consumer-Gruppen und Records des Messaging-Ablaufs.

Im Lesepfad folgt BillingProjectionConsumer: Verarbeitet Order-Events und aktualisiert die BillingProjection.

Typ
class InMemoryKafkaBroker
Verwendet
PartitionKeyStrategy
Verwendet von
MessagingRun6ARebuildDemo, MessagingRun6ARebuildTestRunner
Einstiege
createTopic(String name, int partitions), publish(String topicName, EventEnvelope event), poll(String topicName, String groupId, int maxRecords), commit(String topic, String groupId, int partition, long nextOffset), lag(String topicName, String groupId)
Test Double / Fake Broker
package com.example.messaging.run6a.broker;

import com.example.messaging.run6a.model.EventEnvelope;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;

// Pattern: Test Double / Fake Broker - macht Kafka-Konzepte ohne externe Infrastruktur testbar.
public final class InMemoryKafkaBroker {
    private final Map<String, Topic> topics = new HashMap<>();
    private final Map<String, ConsumerGroupState> groups = new HashMap<>();
    private final PartitionKeyStrategy partitionKeyStrategy = new PartitionKeyStrategy();

    public void createTopic(String name, int partitions) {
        topics.put(name, new Topic(name, partitions));
    }

    public PublishResult publish(String topicName, EventEnvelope event) {
        Topic topic = requireTopic(topicName);
        int partitionId = partitionKeyStrategy.partitionFor(event.partitionKey(), topic.partitionCount());
        long offset = topic.partition(partitionId).append(event);
        return new PublishResult(topicName, partitionId, offset);
    }

    public List<ConsumerRecord> poll(String topicName, String groupId, int maxRecords) {
        Topic topic = requireTopic(topicName);
        ConsumerGroupState group = groups.computeIfAbsent(groupId, ConsumerGroupState::new);
        List<ConsumerRecord> result = new ArrayList<>();
        for (Partition partition : topic.partitions()) {
            if (result.size() >= maxRecords) break;
            long nextOffset = group.nextOffset(topicName, partition.id());
            result.addAll(partition.readFrom(topicName, nextOffset, maxRecords - result.size()));
        }
        return result;
    }

    public void commit(String topic, String groupId, int partition, long nextOffset) {
        groups.computeIfAbsent(groupId, ConsumerGroupState::new).commit(topic, partition, nextOffset);
    }

    public long lag(String topicName, String groupId) {
        Topic topic = requireTopic(topicName);
        ConsumerGroupState group = groups.computeIfAbsent(groupId, ConsumerGroupState::new);
        long lag = 0;
        for (Partition partition : topic.partitions()) {
            lag += partition.size() - group.nextOffset(topicName, partition.id());
        }
        return lag;
    }

    private Topic requireTopic(String topicName) {
        Topic topic = topics.get(topicName);
        if (topic == null) throw new IllegalArgumentException("unknown topic " + topicName);
        return topic;
    }
}

5. BillingProjectionConsumer

src/main/java/com/example/messaging/run6a/consumer/BillingProjectionConsumer.java
Java-Datei öffnen
Rolle im Ablauf

Verarbeitet Order-Events und aktualisiert die BillingProjection.

Im Lesepfad folgt IdempotencyStore: Merkt bereits verarbeitete Nachrichten und verhindert doppelte fachliche Wirkung.

Typ
class BillingProjectionConsumer
Verwendet
EventSchemaRegistry, IdempotencyStore
Verwendet von
MessagingRun6ARebuildDemo, MessagingRun6ARebuildTestRunner
Einstiege
handle(EventEnvelope event), projectionCount(), projection()
package com.example.messaging.run6a.consumer;

import com.example.messaging.run6a.model.EventEnvelope;
import com.example.messaging.run6a.schema.CompatibilityCheck;
import com.example.messaging.run6a.schema.EventSchemaRegistry;

public final class BillingProjectionConsumer {
    private final EventSchemaRegistry schemaRegistry;
    private final IdempotencyStore idempotencyStore;
    private final DeadLetterQueue deadLetterQueue;
    private final BillingProjection projection = new BillingProjection();

    public BillingProjectionConsumer(EventSchemaRegistry schemaRegistry, IdempotencyStore idempotencyStore, DeadLetterQueue deadLetterQueue) {
        this.schemaRegistry = schemaRegistry;
        this.idempotencyStore = idempotencyStore;
        this.deadLetterQueue = deadLetterQueue;
    }

    public ProcessingResult handle(EventEnvelope event) {
        if (idempotencyStore.alreadyProcessed(event.eventId())) {
            return ProcessingResult.duplicateIgnored(event.eventId());
        }
        CompatibilityCheck check = schemaRegistry.validate(event);
        if (!check.accepted()) {
            deadLetterQueue.add(event, check.reason());
            return ProcessingResult.rejected(event.eventId(), check.reason());
        }
        projection.apply(event);
        idempotencyStore.markProcessed(event.eventId());
        return ProcessingResult.processed(event.eventId());
    }

    public int projectionCount() { return projection.size(); }
    public BillingProjection projection() { return projection; }
}

6. IdempotencyStore

src/main/java/com/example/messaging/run6a/consumer/IdempotencyStore.java
Java-Datei öffnen
Rolle im Ablauf

Merkt bereits verarbeitete Nachrichten und verhindert doppelte fachliche Wirkung.

Im Lesepfad folgt MessagingRun6ARebuildTestRunner: Prüft Schema, Partitionierung, Reihenfolge, Consumer-Gruppen und Idempotenz.

Typ
class IdempotencyStore
Verwendet
Verwendet von
MessagingRun6ARebuildDemo, BillingProjectionConsumer, MessagingRun6ARebuildTestRunner
Einstiege
alreadyProcessed(String eventId), markProcessed(String eventId), size()
Idempotent Consumer
package com.example.messaging.run6a.consumer;

import java.util.HashSet;
import java.util.Set;

// Pattern: Idempotent Consumer - verhindert doppelte fachliche Wirkung bei at-least-once.
public final class IdempotencyStore {
    private final Set<String> processedEventIds = new HashSet<>();

    public boolean alreadyProcessed(String eventId) { return processedEventIds.contains(eventId); }
    public void markProcessed(String eventId) { processedEventIds.add(eventId); }
    public int size() { return processedEventIds.size(); }
}

7. MessagingRun6ARebuildTestRunner

src/test/java/com/example/messaging/run6a/MessagingRun6ARebuildTestRunner.java
Java-Datei öffnen
Rolle im Ablauf

Prüft Schema, Partitionierung, Reihenfolge, Consumer-Gruppen und Idempotenz.

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

Typ
class MessagingRun6ARebuildTestRunner
Verwendet
EventSchemaRegistry, InMemoryKafkaBroker, BillingProjectionConsumer, IdempotencyStore
Verwendet von
Einstiege
main(String[] args)
package com.example.messaging.run6a;

import com.example.messaging.run6a.broker.InMemoryKafkaBroker;
import com.example.messaging.run6a.broker.JmsQueue;
import com.example.messaging.run6a.consumer.*;
import com.example.messaging.run6a.model.*;
import com.example.messaging.run6a.schema.EventSchemaRegistry;
import java.time.Instant;
import java.util.Map;

public final class MessagingRun6ARebuildTestRunner {
    public static void main(String[] args) {
        commandEventNaming();
        independentConsumerGroups();
        partitionKeyKeepsOrderTogether();
        idempotentConsumerIgnoresDuplicate();
        unsupportedSchemaGoesToDlq();
        jmsQueueDistributesWork();
        System.out.println("RUN6A_REBUILD_TESTS_OK");
    }

    static void commandEventNaming() {
        EventNamingPolicy policy = new EventNamingPolicy();
        ok(policy.looksLikeCommand("ReserveInventory"));
        ok(policy.looksLikeEvent("InventoryReserved"));
    }

    static void independentConsumerGroups() {
        InMemoryKafkaBroker broker = new InMemoryKafkaBroker();
        broker.createTopic("order-events", 2);
        EventFactory f = new EventFactory();
        broker.publish("order-events", f.orderPlaced("t1", "ORD-1", "C1", 100, "EUR", 1));
        broker.publish("order-events", f.orderPlaced("t2", "ORD-2", "C1", 200, "EUR", 1));
        eq(2, broker.poll("order-events", "billing", 10).size());
        eq(2, broker.poll("order-events", "reporting", 10).size());
    }

    static void partitionKeyKeepsOrderTogether() {
        InMemoryKafkaBroker broker = new InMemoryKafkaBroker();
        broker.createTopic("order-events", 4);
        EventFactory f = new EventFactory();
        int p1 = broker.publish("order-events", f.orderPlaced("t1", "ORD-42", "C1", 100, "EUR", 1)).partition();
        int p2 = broker.publish("order-events", f.paymentAuthorized("t2", "ORD-42", "AUTH-1")).partition();
        eq(p1, p2);
    }

    static void idempotentConsumerIgnoresDuplicate() {
        EventEnvelope event = new EventFactory().orderPlaced("t", "ORD-3", "C3", 300, "EUR", 2);
        BillingProjectionConsumer c = new BillingProjectionConsumer(new EventSchemaRegistry().register("OrderPlaced", 1, 2), new IdempotencyStore(), new DeadLetterQueue());
        ok(c.handle(event) instanceof ProcessingResult.Processed);
        ok(c.handle(event) instanceof ProcessingResult.DuplicateIgnored);
        eq(1, c.projectionCount());
    }

    static void unsupportedSchemaGoesToDlq() {
        DeadLetterQueue dlq = new DeadLetterQueue();
        BillingProjectionConsumer c = new BillingProjectionConsumer(new EventSchemaRegistry().register("OrderPlaced", 1), new IdempotencyStore(), dlq);
        EventEnvelope event = new EventFactory().orderPlaced("t", "ORD-4", "C4", 400, "EUR", 9);
        ok(c.handle(event) instanceof ProcessingResult.Rejected);
        eq(1, dlq.entries().size());
    }

    static void jmsQueueDistributesWork() {
        JmsQueue q = new JmsQueue();
        q.send(new CommandEnvelope("cmd-1", "trace", "inventory-service", "ReserveInventory", Instant.now(), Map.of(), new ReserveInventoryCommand("ORD-1", "P-1", 2)));
        q.send(new CommandEnvelope("cmd-2", "trace", "billing-service", "GenerateInvoice", Instant.now(), Map.of(), new GenerateInvoiceCommand("ORD-1")));
        ok(q.receive().isPresent());
        ok(q.receive().isPresent());
        ok(q.receive().isEmpty());
    }

    static void eq(Object expected, Object actual) {
        if (!java.util.Objects.equals(expected, actual)) throw new AssertionError(expected + " != " + actual);
    }
    static void ok(boolean condition) { if (!condition) throw new AssertionError(); }
}
Alle Projektdateien öffnen (36 Einträge)
⌂ Cockpit