ConsumerController
ConsumerController ist die Empfänger-Seite zuverlässiger
Zustellung. Er:
- Empfängt
Delivery<T>-Envelopes von einem oder mehrerenProducerControllern. - Dedupliziert Duplikate (gleiches
producerId+ gleiche Producer-Inkarnation- gleiche
seq) und hält diese Buchführung dabei auf einem begrenzten Budget — höchstensmaxProducersProducer, und Einträge, dieproducerIdleTtlMslang untätig waren, werden weggeräumt.
- gleiche
- Weist ein fehlerhaftes Envelope ab — ohne
replyTo, mit einerseq, die keine positive Safe-Integer ist, oder mit einem leeren bzw. zu langen Identifier — als Dead Letter, ohne den Handler zu rufen und ohne zu ACKen. - Ruft eine User-Handler-Funktion mit dem Message-Body auf.
- ACKt automatisch an den Producer zurück, sobald der Handler resolved; wirft der Handler, kein ACK — der Producer retransmittiert.
import { ConsumerController, ConsumerControllerOptions } from 'actor-ts/delivery';
const consumer = system.spawn( () => new ConsumerController<OrderEvent>({ handler: async (order) => { await processOrder(order); }, }),);Die Nachrichten des Producers fließen:
Konfiguration
Abschnitt betitelt „Konfiguration“type ConsumerControllerOptionsType<T> = { handler: (body: T) => void | Promise<void>; maxProducers?: number; // default 1024, Infinity = no cap producerIdleTtlMs?: number; // default 300000, Infinity = no sweep};| Option | Wirkung |
|---|---|
handler | Wird für jeden neuen (nicht duplizierten) Message-Body aufgerufen. Pflichtfeld. |
maxProducers | Wie viele Producer die Dedup-Map gleichzeitig hält. Darüber hinaus wird der Eintrag des am längsten nicht genutzten Producers verdrängt. |
producerIdleTtlMs | Wie lange der Eintrag eines Producers ohne Delivery von ihm überlebt. Ein Hintergrund-Sweep räumt die abgelaufenen weg. |
Du konfigurierst weiterhin weder eine Buffer-Größe noch eine Consumer-ID, und die Grenzen dafür, was ein eingehendes Envelope enthalten darf, bleiben nicht konfigurierbar: sie existieren, um Wire-Input zu begrenzen, und ein Limit, das der Sender anheben könnte, ist kein Limit. Die beiden Optionen oben sind eine andere Art von Grenze — ein Budget, das dieser Consumer auf seinen eigenen Heap anwendet, und genau so etwas sollte ein Betreiber einstellen können.
const consumerOptions = ConsumerControllerOptions.create<OrderEvent>() .withHandler(async (order) => { await processOrder(order); }) .withMaxProducers(64) .withProducerIdleTtlMs(60_000);
const consumer = system.spawn( () => new ConsumerController<OrderEvent>(consumerOptions), 'orders-consumer',);Ein einfaches Objekt funktioniert überall dort, wo auch der Builder
funktioniert — { handler, maxProducers } — und ist die Kurzform, die der Rest
dieser Seite verwendet.
Warum der Dedup-State ein Budget hat
Abschnitt betitelt „Warum der Dedup-State ein Budget hat“Der Schlüssel der Map ist die producerId aus dem Envelope, und die wählt
der Sender. Früher kostete jede einzelne davon einen permanenten Eintrag,
die Größe der Map war also eine Funktion davon, wie viele IDs jemals angekommen
waren, und nicht davon, wie viele Producer es gibt. Das ist ein Leck, ganz ohne
Angreifer: ein ProducerController, dessen Aufrufer producerId nicht setzt,
zieht pro Konstruktion eine frische Zufalls-ID — ein langlebiger Consumer, den
kurzlebige Producer bedienen, sammelt also einen Eintrag pro Producer an und
gibt keinen wieder her. Über einen Cluster-Port ist derselbe Pfad ein
Verstärker: ein behaltener Eintrag pro Delivery.
Verdrängung ist nicht gratis — deshalb warnt der Controller darüber, statt Fenster stillschweigend zu verlieren (getaktet auf höchstens eine Zeile pro Minute, damit eine Flut aus einem begrenzten Heap kein unbegrenztes Log macht). Einen Eintrag zu verwerfen verwirft die Duplikat-Unterdrückung dieses Producers, ein danach eintreffender Retransmit lässt Deinen Handler also ein zweites Mal laufen. Dieses Protokoll ist At-least-once und erlaubt das ohnehin; was es nicht erlaubt hat, war unbegrenztes Wachstum.
maxProducers gibt nur dann etwas frei, wenn ein neuer Producer einen Slot
braucht — dafür gibt es zusätzlich producerIdleTtlMs: ein Consumer, der einen
Schwall von Producern gesehen hat und danach still wurde, hat nichts, was eine
Verdrängung auslöst, und nur das Alter gibt diese Einträge wieder frei. Halte
die TTL deutlich über dem resendTimeout Deiner Producer (Default 500 ms) — ein
Eintrag, der verworfen wird, während sein Producer dieselbe Seq noch
retransmittiert, kostet einen doppelten Handler-Aufruf. Der Sweep läuft im
Takt der TTL selbst, ein untätiger Eintrag verschwindet also nach ein bis zwei
davon.
consumer.trackedProducers liefert die aktuelle Anzahl der Einträge. Sie
sollte nahe an der Zahl der Producer liegen, die tatsächlich mit diesem
Consumer sprechen; steigt sie mit der Nachrichtenzahl, ist genau das das
Symptom, das dieses Budget verhindern soll.
Handler-Kontrakt
Abschnitt betitelt „Handler-Kontrakt“Der Controller ruft handler(body) einmal pro neuem
(producerId, incarnation, seq)-Tripel auf. Sobald der Handler resolved,
schickt der Controller ein ACK zurück an den Producer (via
Delivery.replyTo), und der Producer gibt seinen In-Flight-Slot
frei.
Die Inkarnation ist es, die einen neu gestarteten Producer aus dem Dedup-Fenster heraushält: er verwendet seine Sequenznummern wieder, seine Inkarnation aber nicht — seine Nachrichten sind also neu und keine Duplikate. Siehe Restart und die Producer-Inkarnation auf der Seite ProducerController.
new ConsumerController<OrderEvent>({ handler: async (order) => { await processOrder(order); // ACK wird automatisch gesendet, sobald das hier zurückkehrt. },});Handler wirft oder rejected → kein ACK → Producer retransmittiert
Abschnitt betitelt „Handler wirft oder rejected → kein ACK → Producer retransmittiert“new ConsumerController<OrderEvent>({ handler: async (order) => { if (!isValid(order)) { throw new Error('invalid order'); } await processOrder(order); },});Werfen (oder ein gerejectetes Promise zurückgeben) verhindert
das ACK. Der resendTimeoutMs-Timer des Producers feuert und
schickt dieselbe Seq erneut. Das ist der einzige
“NACK”-Mechanismus — es gibt keine explizite NACK-Nachricht.
Duplikat-Behandlung läuft automatisch
Abschnitt betitelt „Duplikat-Behandlung läuft automatisch“Wenn dasselbe (producerId, incarnation, seq) erneut ankommt (z. B. weil der
Producer retransmittiert hat, bevor unser ACK ihn erreicht hat),
überspringt der Controller den Handler und schickt einfach
nochmal das ACK. Kein Applikations-Code nötig — der Controller
hält contiguous und Out-of-order-Seqs pro Producer-Inkarnation intern fest.
Für fachliche Idempotenz (z. B. “eine Kreditkarte nicht doppelt belasten, auch nicht über Consumer-Neustarts hinweg”) persistiere Deine eigene processed-Seq zusammen mit dem Business-State — der In-Memory-Dedup des Controllers ist nach einem Restart weg.
Out-of-order-Behandlung
Abschnitt betitelt „Out-of-order-Behandlung“Producer sendet 1, 2, 3Netzwerk liefert: 1, 3, 2ConsumerController hält am Start { contiguous: 1, above: {} }Nach Ankunft von 3 → { contiguous: 1, above: {3} } (Handler lief für seq 3)Nach Ankunft von 2 → { contiguous: 3, above: {} } (Handler lief für seq 2, dann gleitet contiguous bis 3)Der Controller pflegt einen contiguous-High-Watermark plus
eine Menge von Out-of-order-Seqs, die darüber bereits ausgeliefert
wurden. Jede neu ankommende Seq wird gegen beides geprüft —
Duplikate überspringen den Handler und re-ACKen; neue Seqs
führen den Handler aus und updaten danach den State.
Es gibt keine feste Buffer-Größe für Nachrichten und keinen stillen Drop einer solchen — der Producer retransmittiert ungeACKte Nachrichten weiter, bis der Handler resolved und das ACK ankommt. Begrenzt ist die Dedup-Buchführung dahinter: siehe Warum der Dedup-State ein Budget hat. Ein verlorener Eintrag verwirft nie eine Nachricht, er kostet ein Duplikat.
Failure-Semantik
Abschnitt betitelt „Failure-Semantik“Consumer crasht mitten im Handler: - Handler resolved nie → kein ACK gesendet - Producer retransmittiert nach `resendTimeoutMs` - Neuer Consumer (oder restarteter alter) sieht die Seq erneut - Der In-Memory-Dedup-State des Controllers ist weg, Handler läuft erneutat-least-once-Zustellung; Idempotenz liegt in Deiner Verantwortung. Für Consumer-Crashes, die nicht erneut verarbeiten dürfen, persistiere die processed-Seq zusammen mit dem Business-State im Handler.
Mehrere Producer
Abschnitt betitelt „Mehrere Producer“const consumer = system.spawn( () => new ConsumerController<OrderEvent>({ handler: async (order) => { await processOrder(order); }, }),);
// Mehrere Producer senden an denselben Consumer:const p1 = system.spawnAnonymous(() => new ProducerController({ consumer, producerId: 'p1' }));const p2 = system.spawnAnonymous(() => new ProducerController({ consumer, producerId: 'p2' }));Jeder Producer pflegt seine eigene Seq. Der Consumer
dedupliziert pro Producer automatisch — er hält einen
separaten (contiguous, above)-Dedup-State pro producerId,
sodass die Seq-Räume zweier Producer nicht kollidieren. Der Eintrag hält
außerdem fest, welche Inkarnation dieser producerId die Zähler beschreiben —
ein Delivery einer neuen Inkarnation ersetzt ihn, statt einen zweiten Eintrag
anzulegen. Damit bleibt ein neu startender Producer bei einem Eintrag statt
einem pro Neustart, zum Preis eines möglichen weiteren Duplikats: ein
nachlaufendes Delivery der ausscheidenden Inkarnation setzt das Fenster erneut
zurück, ein paar schon behandelte Seqs können um den Wechsel herum also zweimal
laufen.
Wie viele Einträge es insgesamt gibt, ist Sache von maxProducers und nicht der
Inkarnation — bei mehr verschiedenen Producern als das werden die am längsten
nicht gehörten verdrängt. Hebe den Wert an, wenn ein Consumer legitim mehr
Producer bedient als die voreingestellten 1 024.
Dein Handler sieht nur den Body — producerId, incarnation und seq sind
die Sache des Controllers. Falls Du sie für Business-Logik
brauchst, persistiere einen Wrapper, der sie durch Deine
Verarbeitungs-Pipeline trägt.
Cluster-aware
Abschnitt betitelt „Cluster-aware“Die consumer-Ref des Producers kann auf einen Remote-Consumer
zeigen (anderer Cluster-Node). Der Cluster-Transport
serialisiert die Envelopes; der ACK-Pfad nutzt denselben
Transport in Gegenrichtung.
In der Praxis: platziere den Consumer nah an den Daten seines Handlers — wenn der Handler eine lokale Datenbank trifft, lass den Consumer auf demselben Node laufen wie die Datenbank. Cross-cluster-Delivery ist für durchsatzgebundene Workloads in Ordnung, fügt aber pro ACK eine Round-Trip hinzu.
Wohin als Nächstes
Abschnitt betitelt „Wohin als Nächstes“- Delivery im Überblick — das Gesamtbild.
- ProducerController — die Sender-Seite.
- ACK-Semantik — wann ACKs feuern.
- PersistentActor — zum Persistieren der processed-Seq, wenn Du Effectively-once-Verarbeitung über Consumer-Restarts hinweg brauchst.
