Zum Inhalt springen
Deutsch

ACK-Semantik

Das ACK ist das Kern-Primitiv zuverlässiger Zustellung — der Producer gibt gebufferte Nachrichten frei, sobald sie geACKt sind. Der Controller verwaltet das ACK-Timing für Dich; wie der Handler resolved, steuert, ob das ACK feuert.

1. ConsumerController lässt das Delivery-Envelope zu
→ fehlerhaft (kein replyTo, ungültige seq, zu lange Ids) → Dead Letter,
kein Handler
2. Controller prüft den Dedup-State für (producerId, incarnation)
→ bekannte Seq → nur re-ACK, Handler überspringen
→ neue Seq → handler(body) aufrufen, auf Resolve warten
3. Handler resolved → Controller sendet ACK an msg.replyTo und
spiegelt die Inkarnation des Delivery zurück
Handler wirft → Controller ACKt NICHT
4. Producer prüft, ob das ACK seine eigene Inkarnation spiegelt
→ Abweichung → ignoriert, der Retransmit-Timer bleibt stehen
→ Treffer → entfernt Seq aus dem In-Flight-Buffer
→ gibt einen Window-Slot für den nächsten Send frei

Der Handler ist der einzige User-seitige Schalter. Was im handler(body) passiert, ist Dein Code; was darum herum passiert (Dedup, ACK, replyTo-Verdrahtung, Retransmit-Timing) ist Aufgabe des Controllers.

MusterHandler-FormGarantie
Resolve bei Erfolg (Default)Handler erledigt die Arbeit, gibt void oder ein resolved Promise zurückat-least-once mit In-Memory-Dedup
Wirft bei MisserfolgHandler wirft / rejected, wenn die Verarbeitung fehlschlug → Controller überspringt das ACK → Producer sendet erneutat-least-once mit Retry
Vor Resolve persistierenHandler awaitet einen Journal-Write, bevor er zurückkehrteffectively-once über Consumer-Restart hinweg
new ConsumerController<T>({
handler: async (body) => {
await this.processIdempotently(body); // ACK-bei-Resolve
},
});
1. Handler beginnt die Verarbeitung (body)
2. Handler abgeschlossen (resolved)
3. Controller serialisiert ein ACK-Envelope
4. Controller tellt die replyTo-Ref des Producers
5. Producer entfernt msg aus dem Unacked-Buffer

Wenn der Consumer zwischen 2 und 4 crasht, erreicht das ACK den Producer nie. Der Producer retransmittiert nach resendTimeoutMs. Der In-Memory-Dedup des Controllers ist nach dem Crash weg, also läuft der Handler für dieselbe Seq erneut — mach ihn idempotent oder persistiere die processed-Seq zusammen mit dem Business-State.

Für effectively-once-Verarbeitung von Seiteneffekten await den Journal-Write innerhalb des Handlers. Wenn handler(body) resolved, sind sowohl Dein Event ALS AUCH der implizite “diese Seq ist verarbeitet”-Marker durable; das folgende ACK des Controllers teilt dem Producer nur mit, dass er seinen Slot freigeben kann.

class PersistentEventLog extends PersistentActor<Command, Event, State> {
// …onCommand / onEvent ausgelassen…
}
const log = system.spawnAnonymous(PersistentEventLog);
const consumer = system.spawn(
() => new ConsumerController<OrderEvent>({
handler: async (order) => {
// ask() resolved, sobald der PersistentActor das Event
// journalt hat. Crasht der Consumer, bevor das resolved,
// retransmittiert der Producer und der Handler läuft
// erneut — aber der Dedup des PersistentActor (Event-ID
// oder Seq im Event-Payload) fängt den Re-Run ab.
await log.ask({ kind: 'append', order }, 30_000);
},
}),
);

Die Sequenz bei sauberem Lauf:

  1. Delivery kommt an.
  2. Controller ruft handler(body) auf.
  3. Handler fragt das Journal, das Journal schreibt, das ask resolved.
  4. Controller sendet ACK → Producer gibt seinen Slot frei.

Bei Crash zwischen Schritt 2 und Schritt 3’s Resolve:

  • Bevor der Journal-Write committed: Producer retransmittiert; der Handler läuft erneut; der Dedup auf Journal-Seite fängt ihn ab. Netto: keine Duplizierung.
  • Nachdem der Journal-Write committed hat, aber bevor das ACK des Controllers den Producer erreicht: Producer retransmittiert; der Handler läuft erneut; der Dedup auf Journal-Seite springt sofort raus und der Controller ACKt sofort. Netto: keine Duplizierung.

Für End-to-End-Durability paare das mit persistentem Producer-State — speichere den In-Flight-Buffer des Producers in einem Journal, damit ein Producer-seitiger Crash keine unbestätigten Nachrichten verliert.

Producer sendet msg #5 → Crash, bevor das Transmit abgeschlossen ist

Recovery: Producer startet neu. Ohne Persistenz ist msg #5 verloren. Mit Persistenz ist die Nachricht im Unacked-Buffer-Journal des Producers; wird beim Recovery gesendet.

Alles, was der neu gestartete Producer danach sendet, kommt normal an. Der Restart erzeugt eine neue Producer-Inkarnation, der Consumer eröffnet also ein frisches Dedup-Fenster, statt die wiederverwendeten Sequenznummern als schon behandelte Duplikate zu werten — siehe Restart und die Producer-Inkarnation auf der Seite ProducerController.

Consumer empfängt msg → crasht, bevor der Handler läuft

Kein ACK gesendet; Producer retransmittiert. Idempotenter Handler dedupliziert (oder verarbeitet zum ersten Mal).

Producer sendet msg #5
Consumer verarbeitet msg #5
Consumer sendet ACK #5
Netzwerk verliert das ACK-Envelope

Producer retransmittiert msg #5 nach Timeout. Consumer dedupliziert (highWatermark schon bei 5); ackt erneut.

Das Dedup-und-Re-ACK-Muster heißt, transiente ACKs sind wiederherstellbar — irgendwann kommt eines an.

Producer.windowSize = 16
Consumer braucht 10 s pro Nachricht

Backpressure: sobald 16 unbestätigte Nachrichten in Flight sind, sendet der Producer nicht weiter, bis ein ACK einen Slot freigibt. Eingehende { kind: 'reliable-delivery.send' }-Tells reihen sich innerhalb des Producers ein, bis das Window aufgeht.

Das ist gut — das System kann nicht weiter vor sich selbst davonlaufen als das Flow-Control-Window.

Consumer auf node-A verarbeitet msg #5; Controller queued das ACK
node-A crasht, bevor das ACK-Tell abgeschlossen ist

Producer sieht im resendTimeoutMs-Fenster kein ACK → retransmittiert.

Wenn der Consumer sharded ist (die Entity wandert beim Failover), wird die neue Instanz:

  • Den Retransmit empfangen.
  • Einen leeren In-Memory-Dedup-State haben (frischer ConsumerController).
  • Den Handler aufrufen — der für Effectively-once seine eigene processed-Seq persistieren muss.

Für sharded Consumer immer den Handler mit Persistenz paaren — sonst setzt der Dedup des Controllers bei jedem Rebalance zurück, und Du bekommst volle Wiederverarbeitung.

Das ACK ist Fire-and-Forget vom Controller zum replyTo des Producers. Der Controller wartet nicht auf die Zustellung — wenn das Netzwerk das ACK verliert, retransmittiert der Producer nach resendTimeoutMs, und der Dedup des Controllers re-ACKt sofort.

Das heißt: Dein Handler zahlt keine Latenz-Kosten für das ACK. Der Controller emittiert das ACK, nachdem handler(body) resolved, und macht weiter.

Der Consumer-Handler entscheidet, ob das ACK feuert; mit dem confirm-Callback des Producers beobachtet Dein eigener Code, dass es eintrifft. Greif zur ergonomischen ReliableDelivery-Fassade: ReliableDelivery.producer gibt einen ProducerHandle zurück, dessen tell(body, confirm?) einen optionalen ConfirmationCallback entgegennimmt — (err: Error | null) => void.

const consumer = ReliableDelivery.consumer<OrderEvent>(system, {
handler: async (order) => {
await save(order);
},
});
const producerOptions = ProducerControllerOptions.create<OrderEvent>().withConsumer(consumer.ref as never);
const producer = ReliableDelivery.producer<OrderEvent>(system, producerOptions);
producer.tell(order, (err) => {
// err === null → der Consumer hat diese Seq geACKt; der Slot wird frei
// err !== null → der Producer wurde vor Eintreffen des ACK gestoppt
});

confirm(null) feuert genau einmal — nachdem das ACK des Consumers den Producer erreicht hat und der In-Flight-Slot freigegeben ist. Ein erneut zugestelltes (doppeltes) ACK für dieselbe Seq wird ignoriert, sodass der Callback nie doppelt feuert. Das ist das einzige Producer-seitige Signal, dass eine Nachricht angekommen ist; ohne ihn ist der Send aus Sicht Deines Codes Fire-and-Forget.

Das Signal ist nur so gut wie das ACK dahinter — deshalb verlangt der Producer, dass das ACK sein eigenes Inkarnations-Token zurückspiegelt, bevor er darauf reagiert. seq zählt ab 1 hoch und producerId ist meist ein Name, den Dein eigener Code gewählt hat — beides ist also nicht schwer zu erraten. Das Inkarnations-Token schon, und es verlässt den Producer nur auf Deliveries, die er tatsächlich gesendet hat. Ohne diese Prüfung könnte alles, das den Producer adressieren kann, den Retransmit abbrechen und confirm(null) für eine Nachricht auslösen, die der Consumer nie gesehen hat — der Stream würde stillschweigend auf At-most-once herabgestuft und dabei Erfolg melden. Beachte, was das Signal weiterhin nicht verspricht: es sagt, dass der Consumer die Verantwortung für die Seq übernommen hat — bei einer erneuten Zustellung heißt das, der Handler lief bei einem früheren Versuch, nicht bei diesem.