Spring Kafka Order Flow
Ein vollständiger Java-21-Pfad für Transactional Outbox, Kafka/Redpanda, versionierte Verträge, idempotente Konsumenten, Retry/DLT und Betriebsdiagnose. Nutze die Seite bei einer Integration über Systemgrenzen: Kläre Vertrag, Zustellung, Fehlerpfad, Wiederholung und den Nachweis für kompatibles Verhalten.
Order + Outbox
Ein lokaler PostgreSQL-Commit verhindert den Dual-Write-Fehler zwischen Geschäftsdaten und Ereignis.
Relay + Broker
Das Relay skaliert mit SKIP LOCKED; ein
Wiederholungsfall wird bewusst nicht verschleiert.
Idempotenter Consumer
Verarbeitungsmarke und Projektion werden in derselben Transaktion gespeichert.
1. Architektur und Verantwortlichkeiten
Die Anwendung trennt fachliche Entscheidung, lokale Transaktionsgrenze, Broker-Auslieferung und Consumer-Projektion. Genau diese Grenzen machen Fehlerzustände sichtbar und testbar.
spring-kafka-order-flow/
├── src/main/java/.../domain
├── src/main/java/.../application
├── src/main/java/.../infrastructure
│ ├── persistence
│ ├── messaging
│ ├── observability
│ ├── security
│ └── web
├── src/main/resources/db/migration
├── src/main/resources/schemas
├── src/test/java
├── docs
├── compose-kafka.yaml
├── compose-redpanda.yaml
└── pom.xml
| Grenze | Garantie | bewusst verbleibendes Risiko |
|---|---|---|
| Order + Outbox | atomarer DB-Commit | keine Broker-Verfügbarkeit nötig |
| Relay → Broker | at-least-once | Duplikat nach Absturz möglich |
| Consumer | einmalige Fachwirkung | Marker muss dieselbe DB-Transaktion teilen |
| Schema | V1/V2 tolerant lesbar | semantische Änderungen bleiben Breaking Changes |
2. Transactional Outbox und skalierbares Relay
Mehrere Relay-Instanzen beanspruchen disjunkte Zeilen. Ein
Lock-Timeout macht verwaiste PUBLISHING-Einträge
erneut verarbeitbar. Exponentieller Backoff verhindert eine
enge Fehlerschleife.
@Transactional
public List<OutboxMessage> claim(int limit, int lockTimeoutSeconds) {
return jdbc.sql("""
with selected as (
select event_id from outbox_event
where status in ('PENDING','RETRY') and next_attempt_at <= now()
order by occurred_at
for update skip locked
limit :limit
)
update outbox_event o
set status='PUBLISHING', locked_at=now(), attempts=attempts+1
from selected s
where o.event_id=s.event_id
returning o.*
""").param("limit", limit).query(mapper).list();
}
PUBLISHED, ein normaler und getesteter
Zustand.
3. Idempotenz und Contract Evolution
Der Consumer setzt zuerst eine eindeutige Marke und aktualisiert anschließend die Projektion. Beide Operationen laufen in derselben Transaktion. V2 ergänzt Felder additiv; V1 bleibt lesbar.
// Pattern: Idempotent Consumer + Tolerant Reader
@Transactional
public void consume(String json) {
EventEnvelope event = codec.decode(json); // schemaVersion 1 oder 2
if (!processed.markIfFirst(event.eventId(), CONSUMER)) {
metrics.duplicate();
return;
}
projections.upsert(event); // fehlender salesChannel => UNKNOWN
metrics.consumed();
}
- stabiler Envelope mit Event-ID, Typ, Version, Aggregate-ID und Zeitpunkt
- Partition Key ist die Bestell-ID
- unbekannte Zusatzfelder werden toleriert
-
fehlender
salesChannelaus V1 wird alsUNKNOWNinterpretiert
4. Retry, DLT und Replay
Temporäre Laufzeitfehler werden zweimal wiederholt. Ein ungültiger oder nicht unterstützter Vertrag ist nicht retrybar und geht direkt in das Dead-Letter-Topic. Replay ist ein kontrollierter Betriebsprozess, kein pauschales Zurückkopieren.
Retrybar
Timeout, temporäre Datenbankstörung, kurzzeitige Downstream-Überlastung.
Nicht retrybar
Unbekannte Schemaversion, fehlende Pflichtmetadaten, nicht unterstützter Eventtyp.
Replay-Gate
Ursache behoben, Stichprobe erfolgreich, Idempotenz und Metriken überprüft.
5. Observability und Security
Micrometer erfasst Pending-Outbox, Publish-Fehler, konsumierte Ereignisse und erkannte Duplikate. HTTP-Korrelations-IDs werden in Response und MDC übernommen. Lokale Basic-Auth-Benutzer trennen Reader, Writer und Operations.
| Signal | Aussage | erste Reaktion |
|---|---|---|
| outbox.pending steigt | Relay oder Broker kommt nicht nach | Publish-Fehler, Broker und DB-Locks prüfen |
| consumer lag steigt | Consumer-Durchsatz kleiner als Input | Partition, Rebalance, DB-Latenz prüfen |
| duplicates steigt stark | Relay-Retries oder Replay aktiv | Ursache und Marker-Transaktion prüfen |
| DLT wächst | dauerhafte technische/vertragliche Fehler | Fehlerklassen und Beispielpayloads analysieren |
6. Lokal ausführen und prüfen
Apache Kafka und Redpanda sind alternative lokale Laufzeiten. Beide nutzen dieselbe Kafka-Client-Schnittstelle; der Code enthält keine Laufzeitweiche.
# Apache Kafka 4.3.1
./mvnw test
docker compose -f compose-kafka.yaml up -d
./mvnw spring-boot:run
# Alternative Redpanda 26.1
# docker compose -f compose-redpanda.yaml up -d
# echte PostgreSQL- und Kafka-Integration
./mvnw -Pcontainer-it verify
Der Offline-Check kompiliert den frameworkfreien Kern, führt die Demo aus und prüft POM, YAML, JSON-Schemas, SQL, Pattern-Markierungen sowie Architekturgrenzen.