ConsumerController
Esta página aún no está disponible en tu idioma.
ConsumerController is the receiver side of reliable delivery.
It:
- Receives
Delivery<T>envelopes from one or moreProducerControllers. - Dedups duplicates (same
producerId+ producer incarnation +seq), keeping that bookkeeping on a bounded budget — at mostmaxProducersproducers, and entries idle forproducerIdleTtlMsare swept. - Refuses a malformed envelope — one with no
replyTo, aseqthat is not a positive safe integer, or an empty / over-long identifier — as a dead letter, without running the handler and without acking. - Invokes a user handler function with each message body.
- Acks back to the producer automatically after the handler resolves; if the handler throws, no ack — the producer retransmits.
import { ConsumerController, ConsumerControllerOptions } from 'actor-ts/delivery';
const consumer = system.spawn( () => new ConsumerController<OrderEvent>({ handler: async (order) => { await processOrder(order); }, }),);The producer’s messages flow:
Configuration
Section titled “Configuration”type ConsumerControllerOptionsType<T> = { handler: (body: T) => void | Promise<void>; maxProducers?: number; // default 1024, Infinity = no cap producerIdleTtlMs?: number; // default 300000, Infinity = no sweep};| Option | What it does |
|---|---|
handler | Invoked for every new (un-duplicated) message body. Required. |
maxProducers | Most producers the dedup map holds at once. Past it, the least-recently-used producer’s entry is evicted. |
producerIdleTtlMs | How long a producer’s entry survives with no delivery from it. A background sweep drops the ones past it. |
You still don’t configure a buffer size or a consumer ID, and the bounds on what an incoming envelope may contain stay non-configurable: those exist to bound wire input, and a limit the sender could raise is not a limit. The two options above are a different kind of bound — a budget this consumer applies to its own heap, which is exactly the sort of thing an operator should be able to set.
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',);A plain object works everywhere the builder does — { handler, maxProducers }
— and is the shorthand the rest of this page uses.
Why the dedup state has a budget
Section titled “Why the dedup state has a budget”The map is keyed by the producerId on the envelope, which the sender
chooses. Every distinct one used to cost a permanent entry, so the map’s size
was a function of how many ids had ever arrived rather than of how many
producers exist. That is a leak with nobody attacking it: a ProducerController
whose caller leaves producerId unset mints a fresh random one per
construction, so a long-lived consumer served by short-lived producers
accumulates one entry per producer and never sheds one. Over a cluster port,
the same path is an amplifier — one retained entry per delivery.
Eviction is not free, which is why the controller warns about it instead of losing windows quietly (paced to at most one line a minute, so a flood cannot turn a bounded heap into an unbounded log). Dropping an entry drops that producer’s duplicate suppression, so a retransmit arriving afterwards runs your handler a second time. This protocol is at-least-once and already permits that; unbounded growth is what it did not.
maxProducers only reclaims when a new producer needs a slot, which is why
producerIdleTtlMs exists as well: a consumer that saw a burst of producers
and then went quiet has nothing to trigger an eviction, and age is the only
thing that releases those entries. Keep the TTL well above your producers’
resendTimeout (default 500 ms) — an entry dropped while its producer is
still retransmitting the same seq costs one duplicate invocation. The sweep
runs on the TTL’s own interval, so an idle entry goes between one and two of
them.
Read consumer.trackedProducers for the current entry count. It should sit
near the number of producers actually talking to this consumer; climbing with
the message count is the symptom the budget exists to prevent.
Handler contract
Section titled “Handler contract”The controller calls handler(body) once per new (producerId, incarnation, seq) triple. After the handler resolves, the controller sends an
Ack back to the producer (via Delivery.replyTo) and the producer
releases its in-flight slot.
The incarnation is what keeps a restarted producer out of the dedup window: it re-uses its sequence numbers but not its incarnation, so its messages are new rather than duplicates. See Restart and the producer incarnation on the ProducerController page.
new ConsumerController<OrderEvent>({ handler: async (order) => { await processOrder(order); // Ack is sent automatically after this returns. },});Handler throws or rejects → no Ack → producer retransmits
Section titled “Handler throws or rejects → no Ack → producer retransmits”new ConsumerController<OrderEvent>({ handler: async (order) => { if (!isValid(order)) { throw new Error('invalid order'); } await processOrder(order); },});Throwing (or returning a rejected promise) prevents the Ack.
The producer’s resendTimeoutMs timer fires and re-sends the
same seq. This is the only “nack” mechanism — there is no
explicit nack message.
Duplicate handling is automatic
Section titled “Duplicate handling is automatic”When the same (producerId, incarnation, seq) arrives again (e.g. the producer
retransmitted before our Ack reached it), the controller skips the
handler and re-sends the Ack. No application code needed — the
controller tracks contiguous + out-of-order seqs per producer
incarnation internally.
For business-level idempotency (e.g. “don’t charge a credit card twice even across consumer restarts”), persist your own processed-seq alongside the business state — the controller’s in-memory dedup is lost on restart.
Out-of-order handling
Section titled “Out-of-order handling”Producer sends 1, 2, 3Network delivers them as: 1, 3, 2ConsumerController tracks { contiguous: 1, above: {} } at startAfter 3 arrives → { contiguous: 1, above: {3} } (handler ran for seq 3)After 2 arrives → { contiguous: 3, above: {} } (handler ran for seq 2, then contiguous slides to 3)The controller keeps a contiguous high-watermark plus a set of
out-of-order seqs already delivered above it. Every newly
arrived seq is checked against both — duplicates skip the
handler and re-Ack; new seqs run the handler then update state.
There is no fixed buffer size for messages and no silent drop of one — the producer keeps retransmitting unacked messages until the consumer’s handler resolves and the Ack lands. What is bounded is the dedup bookkeeping behind it: see Why the dedup state has a budget. Losing an entry never drops a message; it costs a duplicate.
Failure semantics
Section titled “Failure semantics”Consumer crashes during the handler: - Handler never resolved → no Ack was sent - Producer retransmits after `resendTimeoutMs` - New consumer (or restarted same consumer) sees the seq again - Controller's in-memory dedup state is gone, so the handler runs againAt-least-once delivery; idempotency is your responsibility. For consumer crashes that must not re-process work, persist the processed-seq alongside business state in your handler.
Multiple producers
Section titled “Multiple producers”const consumer = system.spawn( () => new ConsumerController<OrderEvent>({ handler: async (order) => { await processOrder(order); }, }),);
// Multiple producers send to the same consumer:const p1 = system.spawnAnonymous(() => new ProducerController({ consumer, producerId: 'p1' }));const p2 = system.spawnAnonymous(() => new ProducerController({ consumer, producerId: 'p2' }));Each producer maintains its own seq. The consumer dedups
per producer automatically — it keeps a separate
(contiguous, above) dedup state per producerId, so two
producers’ seq spaces don’t collide. The entry also records which
incarnation of that producerId the counters describe, and a delivery from a
new incarnation replaces it rather than adding a second entry. That keeps a
restarting producer at one entry rather than one per restart, at the cost of
one more possible duplicate: a straggling delivery from the outgoing
incarnation resets the window again, so a handful of already-handled seqs may
run twice around the changeover.
How many entries there are in total is maxProducers’ business, not the
incarnation’s — with more distinct producers than that, the least recently
heard from are evicted. Raise it when a consumer legitimately serves more
producers than the default 1 024.
Your handler sees only the body — producerId, incarnation and seq are
the controller’s concern. If you need them for business
logic, persist a wrapper that carries them through your
processing pipeline.
Cluster-aware
Section titled “Cluster-aware”The producer’s consumer ref can point at a remote consumer
(different cluster node). The cluster transport serialises the
envelopes; the Ack path uses the same transport in reverse.
In practice: place the consumer near its handler’s data — if the handler hits a local database, run the consumer on the same node as the database. Cross-cluster delivery is fine for throughput-bounded workloads but adds a round-trip per ack.
Where to next
Section titled “Where to next”- Delivery overview — the bigger picture.
- ProducerController — the sender side.
- Ack semantics — when acks fire.
- PersistentActor — for persisting processed-seq state if you need effectively-once handling across consumer restarts.
