Zum Inhalt springen
Deutsch

ConsumerController

ConsumerController ist die Empfänger-Seite zuverlässiger Zustellung. Er:

  • Empfängt Delivery<T>-Envelopes von einem oder mehreren ProducerControllern.
  • Dedupliziert Duplikate (gleiches producerId + gleiche Producer-Inkarnation
    • gleiche seq) und hält diese Buchführung dabei auf einem begrenzten Budget — höchstens maxProducers Producer, und Einträge, die producerIdleTtlMs lang untätig waren, werden weggeräumt.
  • Weist ein fehlerhaftes Envelope ab — ohne replyTo, mit einer seq, 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:

Auto-ACK bei Resolve

ProducerController

ConsumerController

handler(body)

type ConsumerControllerOptionsType<T> = {
handler: (body: T) => void | Promise<void>;
maxProducers?: number; // default 1024, Infinity = no cap
producerIdleTtlMs?: number; // default 300000, Infinity = no sweep
};
OptionWirkung
handlerWird für jeden neuen (nicht duplizierten) Message-Body aufgerufen. Pflichtfeld.
maxProducersWie viele Producer die Dedup-Map gleichzeitig hält. Darüber hinaus wird der Eintrag des am längsten nicht genutzten Producers verdrängt.
producerIdleTtlMsWie 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.

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.

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.

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.

Producer sendet 1, 2, 3
Netzwerk liefert: 1, 3, 2
ConsumerController 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.

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 erneut

at-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.

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.

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.