Zum Inhalt springen
Deutsch

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.

HookWann
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).

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 auf MqttMessage). Braucht einen v5-fähigen Broker.

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}`);
}

this.publish(topic, payload, options?) gibt false zurück, wenn die Nachricht verworfen wurde (Encode-Fehler oder Überlauf des Outbound- Puffers):

  • ein string- oder Uint8Array-Payload wird roh gesendet;
  • jeder andere Wert wird via den Codec kodiert (Default JSON).
this.publish('ack/1', 'ok'); // rohe Bytes: ok
this.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);

Ü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 subscribe ohne target liefert passende Nachrichten an das eigene onMessage des Aktors; mit target fan-outet es an jenen Aktor. Überlappende Patterns liefern an jedes Ref höchstens einmal.
  • Ein unsubscribe mit target entfernt jenes Target. Ohne target entfernt 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.

QoSZustellungKosten
0At-most-once. Fire-and-Forget.Am günstigsten; Nachrichten können verloren gehen.
1At-least-once. Broker speichert bis ACK.Default für die meisten Anwendungsfälle. Duplikate möglich.
2Exactly-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.

this.subscribe('sensors/+/temp'); // passt zu sensors/dev1/temp, sensors/dev2/temp, ...
this.subscribe('devices/#'); // passt zu devices/irgendwas/irgendwo/...
WildcardTrifft
+Genau eine Ebene.
#Mehrere Ebenen (muss das letzte Segment sein).
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.

Terminal-Fenster
npm install mqtt
# oder: bun add mqtt
  1. IoT / Telemetrie — Geräte publishen Sensordaten und subscribieren Befehle.
  2. Brücken aus bestehender MQTT-Infrastruktur — wenn der Broker schon da ist.
  3. 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).