KafkaActor
이 콘텐츠는 아직 번역되지 않았습니다.
Defined in: src/io/broker/KafkaActor.ts:151
Kafka producer + consumer in one actor, backed by kafkajs. When
consumer.groupId is set, a consumer is started after connectImplementation
and consumed records are delivered to target. When a producer is
the only goal, leave consumer and topics empty.
Offset-commit semantics.
commitMode: 'auto'(default) — kafkajs commits after each handler returns successfully → at-least-once. Cheap; OK for idempotent handlers.commitMode: 'manual'— pump pauses on each message and waits for an explicitcommitcommand from the handler (#2). The handler is responsible for sending exactly onecommit(ornack) per delivered record. If neither arrives withincommitTimeoutMsthe pump rejects internally, kafkajs treats the partition as failed, and re-delivery happens on rebalance. Produces exactly-once-with-processing: a message that successfully passed throughcommitis committed; a crash ornackre-delivers.
const kafka = system.spawnAnonymous(() => new KafkaActor( KafkaOptions.create() .withBrokers([‘kafka:9092’]) .withConsumer({ groupId: ‘orders’, commitMode: ‘manual’ }) .withTopics([‘orders’]) .withTarget(orderProcessor), )));
class OrderProcessor extends Actor
Extends
Section titled “Extends”BrokerActor<KafkaOptionsType,KafkaCommand,KafkaPublish,KafkaSubscription>
Constructors
Section titled “Constructors”Constructor
Section titled “Constructor”new KafkaActor(
options?):KafkaActor
Defined in: src/io/broker/KafkaActor.ts:164
Parameters
Section titled “Parameters”options?
Section titled “options?”KafkaOptions = {}
Returns
Section titled “Returns”KafkaActor
Overrides
Section titled “Overrides”BrokerActor<KafkaOptionsType, KafkaCommand, KafkaPublish, KafkaSubscription>.constructor
Methods
Section titled “Methods”displayName()
Section titled “displayName()”displayName():
string
Defined in: src/Actor.ts:192
Human-readable name for this actor in log lines and in the DevTools actor tree (#891). Defaults to the full path — which is already the log source, so an actor that doesn’t override this logs exactly as it did before.
override displayName(): string { return `User(${this.entityId})`; }Purely cosmetic. The path stays the identity everywhere that routes,
correlates or aggregates — metric labels, tracing attributes, dead
letters, ActorRef.toString(), every wire identifier — so a display
name is free to be ambiguous, unstable, or shared between actors.
Resolved on every record, not captured once. Two consequences:
keep it cheap and side-effect free, and expect it to be called before
preStart (hence the optional chain — the context is attached after
construction). In exchange a name may be derived from state, and it
updates when that state does. Throwing, or returning anything but a
non-empty string, falls back to the path and warns once: a naming
hook must not be able to take a log line down with it.
ActorOptions.withDisplayName(...) outranks this, for the same reason
withSupervisorStrategy(...) outranks supervisorStrategy
— the spawn site is the more specific statement. It has to: every
Behaviors actor is a TypedActor that inherits this default, so a
method that won would silently swallow the spawn-site value for exactly
the actors that have no subclass to override. For a name that only
becomes known at runtime, this.context.setDisplayName(...) outranks
both.
Returns
Section titled “Returns”string
Inherited from
Section titled “Inherited from”onReceive()
Section titled “onReceive()”onReceive(
command):void
Defined in: src/io/broker/KafkaActor.ts:364
Main message handler. Receives each envelope dequeued from the mailbox. A thrown error (sync or async) is caught by the supervisor.
Parameters
Section titled “Parameters”command
Section titled “command”Returns
Section titled “Returns”void
Overrides
Section titled “Overrides”postRestart()
Section titled “postRestart()”postRestart(
_reason):void|Promise<void>
Defined in: src/Actor.ts:152
Called on the fresh instance after a restart. Default: call preStart().
Parameters
Section titled “Parameters”_reason
Section titled “_reason”Error
Returns
Section titled “Returns”void | Promise<void>
Inherited from
Section titled “Inherited from”postStop()
Section titled “postStop()”postStop():
Promise<void>
Defined in: src/io/broker/BrokerActor.ts:486
Called after the actor has been terminated. Children are already stopped.
Returns
Section titled “Returns”Promise<void>
Inherited from
Section titled “Inherited from”preRestart()
Section titled “preRestart()”preRestart(
_reason,_message?):void|Promise<void>
Defined in: src/Actor.ts:119
Called before a restart, on the instance about to be thrown away.
The default calls postStop() and nothing else.
Override to release what the instance holds outside itself — a file handle, an open socket, a broker connection — or to do something other than drop the message that failed.
Stopping this actor’s children is not done here: the framework tears
them down after this hook returns and waits for them before building the
replacement, because postRestart re-runs preStart and a named child
needs its name back. To keep the children instead, see
Actor.stopChildrenOnRestart.
Parameters
Section titled “Parameters”_reason
Section titled “_reason”Error
_message?
Section titled “_message?”Returns
Section titled “Returns”void | Promise<void>
Inherited from
Section titled “Inherited from”preStart()
Section titled “preStart()”preStart():
Promise<void>
Defined in: src/io/broker/BrokerActor.ts:477
Called after construction and before the first message is processed.
Returns
Section titled “Returns”Promise<void>
Inherited from
Section titled “Inherited from”stopChildrenOnRestart()
Section titled “stopChildrenOnRestart()”stopChildrenOnRestart():
boolean
Defined in: src/Actor.ts:149
Whether a restart tears this actor’s children down before rebuilding it.
Default: true.
A restart replaces the Actor instance while the cell — and therefore
the child map — survives. Keeping the children was the old behaviour and
it made an ordinary pattern impossible: postRestart re-runs preStart,
so an actor that spawns a named child there hit Child name … is not unique on its first restart and never recovered (#634).
Override to false when the children are expensive to rebuild, hold
state the parent cannot restore, or are supervised independently — a
connection pool, say. They then outlive the restart exactly as before,
and it is on you to make preStart idempotent — by adopting the survivor
from the cell, this.child = this.context.child('name').toNullable() ?? this.context.spawn(Child, 'name'), or with context.spawnAnonymous.
An instance field cannot do it: preStart runs on a fresh instance
after every restart, so this.child ??= … is always unset and re-spawns
into the name the surviving child still holds, which fails the spawn and
restarts the actor again.
This is a separate hook rather than a preRestart override because the
teardown has to be awaited: the new instance cannot be built until the
old children are actually gone, and preRestart has no way to tell the
cell that it started something worth waiting for.
Returns
Section titled “Returns”boolean
Inherited from
Section titled “Inherited from”BrokerActor.stopChildrenOnRestart
supervisorStrategy()
Section titled “supervisorStrategy()”supervisorStrategy():
SupervisorStrategy
Defined in: src/Actor.ts:160
Supervisor strategy for this actor’s children. Defaults to restart, up to 10 times per minute, then stop.
