MQTT
MqttActor integriert mit MQTT-Brokern (Mosquitto, EMQ X, HiveMQ,
AWS IoT Core). Unterstützt sowohl MQTT 3.1.1 als auch 5; Topic-
Wildcards (+, #); QoS 0/1/2; Retained Messages.
Es ist das MQTT-Gegenstück zu
WebsocketClientActor: eine abstrakte
Basisklasse, von der du ableitest. Deklariere Subscriptions im
Konstruktor, verarbeite eingehenden Traffic in onMessage und publishe
mit this.publish(...). Von außen bleibt er per ref.tell(...)
steuerbar.
import { ActorSystem, ActorSystemOptions, MqttActor, MqttOptions, type MqttMessage } from 'actor-ts';
type Reading = { sensor: string; celsius: number };
const actorSystemOptions = ActorSystemOptions.create().withConfig({ 'actor-ts': { io: { broker: { mqtt: { brokerUrl: 'mqtt://localhost:1883' } } } }, });class TemperatureHub extends MqttActor<Reading> { constructor(options: MqttOptions) { super(options.withQos(1).withClientId('temperature-hub')); this.subscribe('sensors/+/temp'); this.subscribe('alerts/#', { qos: 2 }); }
override onMessage(message: MqttMessage<Reading>): void { const { sensor, celsius } = message.payload.entity(); // via den Codec dekodiert this.log.info(`${sensor}: ${celsius}°C`); this.publish(`ack/${sensor}`, 'ok'); // roher String this.publish('rollup', { sensor, celsius }); // kodierte Entity }}
const system = ActorSystem.create('demo', actorSystemOptions);
system.spawn(() => new TemperatureHub(MqttOptions.create()), 'hub');T (hier Reading) typisiert das eingehende Payload —
message.payload.entity() liefert ein T zurück. Ein zweiter Generic,
TSelf, typisiert Anwendungsnachrichten, die andere Aktoren dem Ref
tellen können (siehe Externe Steuerung unten); Default ist never.
Diese überschreibst du in deiner Subklasse. onReceive ist
versiegelt — die Basisklasse dispatcht an die Hooks unten;
überschreibe es nicht.
| Hook | Wann |
|---|---|
onMessage(message) | Eine Nachricht traf auf einer eigenen Subscription dieses Aktors ein. Erforderlich. |
onConnected() | Die Verbindung wurde (neu) geöffnet; die Subscription-Registry ist beim Broker erneut angewendet. |
onDisconnected(cause?) | Die Verbindung brach ab; ein Reconnect-Zyklus kann folgen (gemäß Settings). |
onInvalidMessage(err, message) | onMessage warf einen MqttDecodeError (ein lazy entity() auf einem fehlerhaften Payload). Default: loggen + verwerfen. Erneut werfen zum Eskalieren. |
onSelfMessage(message) | Eine Anwendungsnachricht (TSelf) wurde diesem Ref getellt. |
Eingehende und Lifecycle-Ereignisse werden über die Mailbox zugestellt,
sodass onMessage und die Hooks stets auf dem Aktor-Thread laufen
(single-threaded, Reihenfolge pro Verbindung erhalten: connected →
Nachrichten → disconnected).
Konfiguration
Abschnitt betitelt „Konfiguration“Settings lösen mit der üblichen Präzedenz auf: Konstruktor / Builder >
HOCON (actor-ts.io.broker.mqtt) > eingebaute Defaults. Nutze den
Fluent-Builder MqttOptions oder übergib ein einfaches
Partial<MqttOptionsType>.
const options = MqttOptions.create() .withBrokerUrl('mqtts://mqtt.example.com:8883') .withClientId('my-app') .withCredentials(process.env.MQTT_USER, process.env.MQTT_PASS) .withQos(1) // Default-QoS für publish/subscribe .withProtocolVersion(5) // MQTT 5.0 aktivieren .withCleanSession(false) .withKeepAlive(30);interface MqttOptionsType extends BrokerCommonOptionsType { brokerUrl?: string; // mqtt:// mqtts:// ws:// wss:// clientId?: string; credentials?: { username?: string; password?: string }; qos?: 0 | 1 | 2; // Default-QoS für publish/subscribe cleanSession?: boolean; // Default true keepAlive?: number; // Sekunden; Default 60 protocolVersion?: 4 | 5; // Default 4 (= MQTT 3.1.1) codec?: MqttCodec<unknown>; // Default mqttJsonCodec() will?: { topic: string; payload: string | Uint8Array; qos?: 0 | 1 | 2; retain?: boolean };}Häufige Muster:
withCleanSession(false)— die Session bleibt bestehen; der Broker liefert Nachrichten, die während der Trennung verpasst wurden (abhängig von der Broker-Konfiguration).withWill({ ... })— der Broker publisht das, wenn der Client unsauber die Verbindung verliert. Nützlich für Präsenz (“device-42-offline”).withProtocolVersion(5)— aktiviert MQTT-5-Features (User- Properties, Reason-Codes aufMqttMessage). Braucht einen v5-fähigen Broker.
Typisierte Payloads
Abschnitt betitelt „Typisierte Payloads“Eingehende MqttMessage<T> tragen ein lazy dekodierendes
MqttPayload<T>:
override onMessage(message: MqttMessage<Reading>): void { message.payload.bytes; // rohes Uint8Array message.payload.text(); // UTF-8-String (gecacht) message.payload.entity(); // dekodiertes T via den Codec (gecacht) message.payload.entity<Acknowledgment>(); // auf einem bestimmten Topic als anderen Typ dekodieren}Das Dekodieren ist lazy — es passiert erst, wenn du text() /
entity() aufrufst. Ein fehlerhaftes Payload lässt entity() einen
MqttDecodeError werfen; weil das innerhalb von onMessage geschieht,
fängt die Basisklasse ihn und leitet ihn an onInvalidMessage weiter
(Default: loggen + verwerfen, kein Neustart):
protected override onInvalidMessage(err: MqttDecodeError, message: MqttMessage<Reading>): void { this.log.warn(`bad payload on ${err.topic}: ${err.message}`);}Publishen
Abschnitt betitelt „Publishen“this.publish(topic, payload, options?) gibt false zurück, wenn die
Nachricht verworfen wurde (Encode-Fehler oder Überlauf des Outbound-
Puffers):
- ein
string- oderUint8Array-Payload wird roh gesendet; - jeder andere Wert wird via den Codec kodiert (Default JSON).
this.publish('ack/1', 'ok'); // rohe Bytes: okthis.publish('rollup', { sensor, celsius }); // JSON: {"sensor":...,"celsius":...}this.publish('cfg', data, { qos: 1, retain: true });Um einen nackten String als JSON-Entity zu publishen (die Wire-Bytes
"pong" statt pong), kodiere ihn explizit:
this.publish('topic', this.codec().encode('pong'));Der Codec wandelt Entities in Bytes und zurück. Default ist
mqttJsonCodec() (einfaches JSON über UTF-8). Übergib deinen eigenen
via withCodec(...) — z. B. einen JSON-Codec mit Laufzeit-Validierung:
import { mqttJsonCodec, MqttOptions } from 'actor-ts';
const codec = mqttJsonCodec<Reading>({ validate: (v) => ReadingSchema.parse(v), // zod usw.; wirft → MqttDecodeError});
MqttOptions.create().withBrokerUrl('mqtt://localhost:1883').withCodec(codec);Externe Steuerung
Abschnitt betitelt „Externe Steuerung“Über die Subklassen-API hinaus kann jeder Aktor einen MqttActor
steuern, indem er ihm ein MqttCommand tellt:
ref.tell({ kind: 'publish', publish: { topic: 't', payload: 'hi', qos: 1 } });
// Subscribe und an das eigene onMessage dieses Aktors routen (kein target):ref.tell({ kind: 'subscribe', topic: 'x/#', qos: 1 });
// Subscribe und an einen anderen Aktor fan-outen (externes target):ref.tell({ kind: 'subscribe', topic: 'y/#', target: someHandler });
ref.tell({ kind: 'unsubscribe', topic: 'y/#', target: someHandler });- Ein
subscribeohnetargetliefert passende Nachrichten an das eigeneonMessagedes Aktors; mittargetfan-outet es an jenen Aktor. Überlappende Patterns liefern an jedes Ref höchstens einmal. - Ein
unsubscribemittargetentfernt jenes Target. Ohnetargetentfernt es alle fremden Targets, lässt aber die eigene Subscription des Aktors bestehen — ein externer Controller kann die im Konstruktor deklarierte Subscription der Subklasse nicht stummschalten. - Fan-out-Targets werden per Deathwatch beobachtet: stoppt ein Target- Aktor, wird er automatisch entfernt (und ein Broker-UNSUBSCRIBE feuert, sobald ein Pattern keine Konsumenten mehr hat).
Da diese Kommandos einfache Objekte sind, halte ein selbst definiertes
TSelf von den kind-Werten publish / subscribe / unsubscribe
fern — sonst würde es als Kommando dispatcht, statt onSelfMessage zu
erreichen.
QoS-Level
Abschnitt betitelt „QoS-Level“| QoS | Zustellung | Kosten |
|---|---|---|
| 0 | At-most-once. Fire-and-Forget. | Am günstigsten; Nachrichten können verloren gehen. |
| 1 | At-least-once. Broker speichert bis ACK. | Default für die meisten Anwendungsfälle. Duplikate möglich. |
| 2 | Exactly-once. Voller Handshake. | Am langsamsten; selten nötig. |
Für IoT-Telemetrie ist QoS 1 die typische Balance. Befehle an Geräte: ebenfalls QoS 1 (oder QoS 2, falls Duplikate schädlich wären). Statusupdates, die schnell überholt werden (Sensorwerte): QoS 0 ist okay.
Topic-Wildcards
Abschnitt betitelt „Topic-Wildcards“this.subscribe('sensors/+/temp'); // passt zu sensors/dev1/temp, sensors/dev2/temp, ...this.subscribe('devices/#'); // passt zu devices/irgendwas/irgendwo/...| Wildcard | Trifft |
|---|---|
+ | Genau eine Ebene. |
# | Mehrere Ebenen (muss das letzte Segment sein). |
Retained Messages
Abschnitt betitelt „Retained Messages“this.publish('config/device-42', { mode: 'eco' }, { retain: true });Eine Retained Message wird vom Broker gespeichert und an jeden neuen
Subscriber sofort ausgeliefert. Verwende sie für Konfiguration (Geräte,
die beim Verbinden die aktuellste Config lesen) oder den letzten
bekannten Status. Lass retain bei normalem Pub/Sub weg — sonst hält
der Broker jede alte Nachricht für immer sichtbar.
Peer-Dependency
Abschnitt betitelt „Peer-Dependency“npm install mqtt# oder: bun add mqttWann MQTT
Abschnitt betitelt „Wann MQTT“- IoT / Telemetrie — Geräte publishen Sensordaten und subscribieren Befehle.
- Brücken aus bestehender MQTT-Infrastruktur — wenn der Broker schon da ist.
- Leichtgewichtiges Pub/Sub — wenn Kafka Overkill ist (kleine Payloads, keine Replay-Historie nötig).
Für clusterinternes Pub/Sub ist DistributedPubSub einfacher (kein externer Broker).
Wohin als Nächstes
Abschnitt betitelt „Wohin als Nächstes“- I/O-Übersicht — das große Bild.
- BrokerActor-Basis — der gemeinsame Lifecycle.
- WebSocket-Client — dieselbe subclass-first Form.
- Kafka — für durables Streaming.
