Zum Inhalt springen
Deutsch

ProducerController

ProducerController umwickelt die ausgehenden Nachrichten eines Senders, vergibt Sequence Numbers und hält sie, bis sie vom Consumer geACKt sind.

import { ProducerController, ProducerControllerOptions } from 'actor-ts/delivery';
const producerControllerOptions = ProducerControllerOptions.create<OrderEvent>()
.withProducerId('orders') // stabile Identität
.withConsumer(consumerRef)
.withWindowSize(16) // Flow-Control-Window
.withResendTimeout(500);
const producer = system.spawn(
() => new ProducerController<OrderEvent>(producerControllerOptions),
'producer',
);
producer.tell({ kind: 'reliable-delivery.send', body: { orderId: 'o-1' } });

Intern:

  • Vergibt Seq 1 für die erste Nachricht, 2 für die nächste, …
  • Sendet jede via Wrapped Tell an den Consumer.
  • Hält sie in einem Buffer, bis sie geACKt sind.
  • Retransmittiert, wenn binnen resendTimeout kein ACK kommt.
type ProducerControllerOptionsType<T> = {
consumer: ActorRef<Delivery<T>>;
resendTimeout?: number; // Default 500
windowSize?: number; // Default 16
producerId?: string; // wird zufällig generiert, wenn weggelassen
};
FeldZweck
consumerDie ConsumerController-Ref, an die zugestellt wird.
resendTimeoutWie lange auf ein ACK gewartet wird, bevor retransmittiert wird.
windowSizeFlow-Control-Window — pausiert das Queuing nach N unbestätigten In-Flight-Nachrichten.
producerIdStabiler Identifier, an dem der Consumer seinen Dedup-State festmacht. Wird zufällig generiert, wenn weggelassen; setze ihn fest, wenn Du über Restarts hinweg eine Identität in Logs, Metriken und der Map des Consumers willst. Er trägt kein Dedup-Fenster über einen Restart — siehe Restart und die Producer-Inkarnation.

Eine generierte producerId ist zufällig, nicht fortlaufend. Früher war sie ein Modul-Zähler — producer-1, producer-2, … — und das war gleich doppelt falsch. Ein Acknowledgment nennt eine producerId und eine seq, und die Seq ist bauartbedingt eine kleine ganze Zahl; ein Zähler ließ von diesem Paar also keine Hälfte mehr zu erraten übrig. Der Zähler war außerdem geteilt: zwei Prozesse desselben Service erzeugten beide producer-1, stempelten dieselbe Identität auf ihre Deliveries und setzten sich danach gegenseitig laufend das Dedup-Fenster bei einem gemeinsamen Consumer zurück. Das Framework zieht jetzt stattdessen 16 Hex-Zeichen.

Diese Id wird bei jeder Konstruktion neu gezogen — producerId wegzulassen heißt also, dass es gar keine Identität gibt, die stabil sein könnte. Setze sie selbst, sobald irgendetwas weiter unten diesen Producer nach einem Restart wiedererkennen soll.

windowSize: 16

Der Producer hält bis zu 16 unbestätigte Nachrichten in Flight. Darüber hinaus werden eingehende { kind: 'reliable-delivery.send' }-Nachrichten eingereiht — das tell des Callers ist erfolgreich, aber das eigentliche Senden pausiert, bis ACKs den Slot freigeben.

Es gibt kein Signal zurück an den Caller, wenn dieses Queuing passiert: ein reliable-delivery.send ist Fire-and-Forget — der Producer antwortet nie darauf, es gibt also kein ask-basiertes “Ich habe sie gesendet”-ACK, auf das man warten könnte.

Für Sender, die wissen müssen, ob eine Nachricht durchgekommen ist, übergibst du einen confirm-Callback beim Send. Er feuert, sobald der Consumer diese Nachricht per ACK bestätigt — mit null bei Erfolg oder einem Error, wenn der Producer stoppt, während die Nachricht noch in der Warteschlange liegt:

producer.tell({
kind: 'reliable-delivery.send',
body: { orderId: 'o-1' },
confirm: (err) => {
if (err) console.error('delivery failed', err);
else console.log('delivered and acked');
},
});

Um daraus echtes Backpressure zu machen, richte dein Sende-Tempo nach diesen Bestätigungen aus — halte die nächste Charge zurück, bis frühere Nachrichten bestätigt sind.

Wenn ein ACK nicht binnen resendTimeout eintrifft:

Seq 5 zum Zeitpunkt t=0 gesendet
ACK für Seq 5 kommt nie
bei t=500 → Seq 5 retransmitten
bei t=1000 → Seq 5 retransmitten
... bis das ACK kommt oder der Producer stoppt

Retransmits nutzen dieselbe Seq wie das Original — und dieselbe Producer-Inkarnation, damit der Consumer sie als dieselbe Nachricht erkennt.

Jede Konstruktion eines ProducerController erzeugt ein frisches, nicht erratbares Inkarnations-Token und stempelt es auf jedes Delivery. Es ist nicht konfigurierbar; es existiert, weil producerId allein zwei Fragen nicht beantworten kann.

Ist das ein Retransmit oder ein neu gestarteter Producer? Der Sequenz-Zähler liegt im Speicher, also nummeriert ein neu gestarteter Producer wieder ab 1, während producerId unverändert bleibt. Ein Consumer, der nur an producerId festmacht, würde das ganze Präfix nach dem Restart als schon behandelte Sequenznummern sehen, es als Duplikate absorbieren und — weil ein absorbiertes Delivery trotzdem bestätigt wird — jede einzelne Nachricht als erfolgreich zugestellt zurückmelden. Der Consumer macht stattdessen an (producerId, incarnation) fest, also eröffnet eine neue Inkarnation ein frisches Dedup-Fenster und ihre Nachrichten erreichen den Handler.

Kam dieses Acknowledgment von einem echten Delivery? producerId und seq sind beide erratbar, ein dreifeldriges Ack ließe sich also von allem fabrizieren, das den Producer adressieren kann — es würde den Retransmit abbrechen und Dein confirm mit Erfolg auslösen. Das Inkarnations-Token verlässt den Producer nur auf Deliveries, die er tatsächlich gesendet hat; es zurückzuspiegeln ist damit der Nachweis, den der Producer verlangt.

Ein Ack mit der falschen Inkarnation — gefälscht oder ein Nachläufer der vorherigen Inkarnation derselben producerId — wird ignoriert.

// Ohne Persistenz:
Producer crasht → Buffer weg → unbestätigte Nachrichten sind verloren
// Beim Restart setzt das Senden bei Seq 1 wieder ein, unter einer neuen
// Inkarnation — der Consumer behandelt diese Sends also als neu und
// nicht als Duplikate
// Mit ProducerController + PersistentActor-Wrapper:
Producer crasht → Recovery baut Buffer + last-acked-Seq wieder auf
// Beim Restart setzt das Senden bei der richtigen Seq fort

Für volle Durability den Producer in ein PersistentActor-Muster einwickeln — ausgehende Nachrichten persistieren, bevor gesendet wird; beim Restart unbestätigte neu abspielen.

Das einfache ProducerController des Frameworks ist in-memory — ausreichend für ephemere Streams. Für durable Streams legst du Persistenz darauf.

// Ein Producer pro Consumer:
const p1 = system.spawnAnonymous(() => new ProducerController({ consumer: c1Ref, ... }));
const p2 = system.spawnAnonymous(() => new ProducerController({ consumer: c2Ref, ... }));

Jeder ProducerController ist 1:1 mit einem ConsumerController. Für Fan-out senden mehrere Producer an verschiedene Consumer.

Für routing-basiertes Fan-out (ein logischer Stream auf N Consumer aufgeteilt) würdest du einen eigenen Router darüber schreiben.

Die consumer-Ref des ProducerController kann auf einen Remote-Consumer zeigen (anderer Cluster-Node). Der Cluster-Transport serialisiert Envelopes; Retransmissions funktionieren gleich.

Der Producer hält den Buffer lokal — bei einem Crash des Producer-Hosts sind diese Nachrichten verloren, sofern nicht persistiert.

ref.tell(msg): ProducerController:
- Keine Seq, kein ACK, kein Retransmit - Vertrag zuverlässige Zustellung
- Verloren bei totem Empfänger - Übersteht transiente Fehler
- Sub-Mikrosekunden-Kosten - ~50µs Overhead pro Nachricht

Nimm den Controller nur für Streams, in denen Verlust inakzeptabel ist; rohes tell für alles andere.