Ir al contenido
Español

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 more ProducerControllers.
  • Dedups duplicates (same producerId + producer incarnation + seq), keeping that bookkeeping on a bounded budget — at most maxProducers producers, and entries idle for producerIdleTtlMs are swept.
  • Refuses a malformed envelope — one with no replyTo, a seq that 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:

auto-Ack on 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
};
OptionWhat it does
handlerInvoked for every new (un-duplicated) message body. Required.
maxProducersMost producers the dedup map holds at once. Past it, the least-recently-used producer’s entry is evicted.
producerIdleTtlMsHow 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.

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.

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.

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.

Producer sends 1, 2, 3
Network delivers them as: 1, 3, 2
ConsumerController tracks { contiguous: 1, above: {} } at start
After 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.

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 again

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

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.

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.