Перейти к содержимому
Русский

ProducerController

Это содержимое пока не доступно на вашем языке.

ProducerController wraps a sender’s outgoing messages, assigning sequence numbers + holding them until acked by the consumer.

import { ProducerController, ProducerControllerOptions } from 'actor-ts/delivery';
const producerControllerOptions = ProducerControllerOptions.create<OrderEvent>()
.withProducerId('orders') // stable identity
.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' } });

Internally:

  • Assigns seq 1 to the first message, 2 to the next, …
  • Sends each via wrapped tell to the consumer.
  • Holds in a buffer until acked.
  • Retransmits if no ack within resendTimeout.
type ProducerControllerOptionsType<T> = {
consumer: ActorRef<Delivery<T>>;
resendTimeout?: number; // default 500
windowSize?: number; // default 16
producerId?: string; // randomly generated if omitted
};
FieldPurpose
consumerThe ConsumerController ref to deliver to.
resendTimeoutHow long to wait for an ack before retransmitting.
windowSizeFlow-control window — pauses queueing after N unacked messages in flight.
producerIdStable identifier the consumer keys its dedup state on. Randomly generated if omitted; pin it when you want one identity across restarts in logs, metrics and the consumer’s map. It does not carry a dedup window across a restart — see Restart and the producer incarnation.

A generated producerId is random, not sequential. It used to be a module counter — producer-1, producer-2, … — which was wrong twice over. An Acknowledgment names a producerId and a seq, and the seq is a small integer by construction, so a counter left neither half of that pair to guess. The counter was also shared: two processes running the same service each minted producer-1, stamped one identity onto their deliveries, and then kept resetting each other’s dedup window at a consumer they had in common. The framework now draws 16 hex characters instead.

That id is minted fresh on every construction, so leaving producerId unset means there is no identity to be stable about — set it yourself whenever something downstream is supposed to recognise this producer after a restart.

windowSize: 16

The producer holds up to 16 unacked messages in flight. Beyond that, incoming { kind: 'reliable-delivery.send' } messages are queued — the caller’s tell succeeds, but actual sending pauses until the consumer’s acks free up the buffer.

There is no signal back to the caller when this queueing happens: a reliable-delivery.send is fire-and-forget — the producer never replies to it, so there is no ask-based “I sent it” acknowledgment to await.

For senders that need to know a message got through, pass a confirm callback on the send. It fires once the consumer acks that message — with null on success, or an Error if the producer stops while the message is still queued:

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');
},
});

To turn that into real backpressure, pace your sends against these confirmations — hold off on the next batch until earlier messages confirm.

When an ack doesn’t arrive within resendTimeout:

seq 5 sent at t=0
seq 5 ack never arrives
at t=500 → retransmit seq 5
at t=1000 → retransmit seq 5
... until ack arrives or producer stops

Retransmits use the same seq as the original, and carry the same producer incarnation, so the consumer recognises them as the same message.

Every construction of a ProducerController mints a fresh, unguessable incarnation token and stamps it onto each Delivery. Nothing configures it; it exists because producerId alone cannot answer two questions.

Is this a retransmit or a restarted producer? The sequence counter is in-memory, so a restarted producer numbers from 1 again while producerId stays as configured. A consumer keyed on producerId alone would see the whole post-restart prefix as sequence numbers it had already handled, absorb them as duplicates, and — because an absorbed delivery is still acked — report each one back as successfully delivered. The consumer keys on (producerId, incarnation) instead, so a new incarnation starts a fresh dedup window and its messages reach the handler.

Did this acknowledgment come from a real delivery? producerId and seq are both guessable, so a three-field ack could be manufactured by anything able to address the producer, cancelling its retransmit and firing your confirm with success. The incarnation travels only on deliveries the producer actually sent, so echoing it back is the evidence the producer requires.

An ack that names the wrong incarnation — forged, or a straggler from the previous incarnation of the same producerId — is ignored.

// Without persistence:
producer crashes → buffer lost → unacked messages are gone
// On restart, send resumes from seq 1 under a new incarnation, so the
// consumer treats those sends as new rather than as duplicates
// With ProducerController + PersistentActor wrapper:
producer crashes → recovery rebuilds buffer + last-acked-seq
// On restart, send resumes from the right seq

For full durability, wrap the producer in a PersistentActor pattern — persist outgoing messages before sending; on restart, replay any unacked.

The framework’s plain ProducerController is in-memory — sufficient for ephemeral streams. For durable streams, layer persistence on top.

// One producer per consumer:
const p1 = system.spawnAnonymous(() => new ProducerController({ consumer: c1Ref, ... }));
const p2 = system.spawnAnonymous(() => new ProducerController({ consumer: c2Ref, ... }));

Each ProducerController is 1:1 with one ConsumerController. For fan-out, multiple producers send to different consumers.

For routing-based fan-out (one logical stream split across N consumers), you’d write a custom router on top.

ProducerController’s consumer ref can point at a remote consumer (different cluster node). The cluster transport serializes envelopes; retransmissions work the same.

The producer holds the buffer locally — on producer-host crash, those messages are lost unless persisted.

ref.tell(msg): ProducerController:
- No seq, no ack, no retransmit - Reliable delivery contract
- Lost on dead recipient - Survives transient failures
- Sub-microsecond cost - ~50µs per message overhead

Use the controller only for streams where loss is unacceptable; raw tell for everything else.