MQTT Client Driver — Technische Referenz
Der MQTT-Client-Treiber von PlantPulse Edge ist eine Java-Implementierung, die die Bibliothek Eclipse Paho mqttv3 1.2.5 kapselt. Er abonniert Topics eines externen IIoT-Brokers (HiveMQ / Mosquitto / EMQX / AWS IoT Core / Azure IoT Hub usw.), hält die Nachrichten in einem In-Memory-Cache und liefert bei jedem Polling-Zyklus eines Tags den Cache-Wert zurück. Schreibvorgänge werden sofort publiziert.
Quelle: plantpulse.driver.protocol.mqtt_client.*
| Klasse | Verantwortung |
|---|---|
MQTTClientDriver | ProtocolDriver-Implementierung — connect/read/write/close, Topic-Cache, Lazy Subscribe |
Bibliothek: lib/org.eclipse.paho.client.mqttv3-1.2.5.jar
Unterschied zu Sparkplug B
PlantPulse stellt zwei getrennte MQTT-basierte Treiber bereit.
| Treiber | Payload | Einsatzfall |
|---|---|---|
MQTT_CLIENT (diese Seite) | beliebiger String/JSON/Binärinhalt | allgemeine IIoT-Broker, benutzerdefinierte Topics, ASCII-/JSON-Payload |
SPARKPLUG_B | Sparkplug B Payload (Protobuf) | Cirrus Link / Tahu / Ignition-Standard, NBIRTH/DBIRTH-Alias |
Sparkplug B bringt einen eigenen Wire-Codec und eine Alias-Lernlogik mit; in Umgebungen, die nur allgemeines MQTT benötigen, ist dieser Treiber schlanker und kommt mit der reinen Paho-Standard-API aus.
Ablauf
broker ──PUBLISH──> Paho MqttClient ──MqttCallback.messageArrived(topic, payload)──> Driver
│
└─ lastValueByTopic.put(topic, new String(payload))
PlantPulse 폴러 ──read(addr)──> lastValueByTopic[addr] (in-memory hit)
write(addr, value) ──Paho.publish(topic, value, qos, retain=false)──> broker
Ergibt der erste read-Aufruf einen Cache-Miss, erfolgt ein Lazy Subscribe (mit dem QoS-Optionswert)
und es wird ein leerer String zurückgegeben — ab dem nächsten Zyklus kommen die tatsächlichen Werte.
In tag_map registrierte Topics werden unmittelbar nach connect gesammelt über preSubscribeFromTagMap() abonniert, sodass
bereits im ersten Zyklus ein Wert eintreffen kann (abhängig von der Publish-Frequenz).
Optionen
| Schlüssel | Standardwert | Beschreibung |
|---|---|---|
username | (keiner) | Benutzername für die Broker-Authentifizierung |
password | (keiner) | Passwort für die Broker-Authentifizierung |
tls | false | true → ssl:// (Standard 8883), false → tcp:// (Standard 1883) |
qos | 0 | QoS für Subscribe/Publish (0 / 1 / 2). Werte außerhalb des Bereichs werden automatisch begrenzt. |
keep-alive | 60 | Keep-Alive-Intervall (Sekunden) — MqttConnectOptions.setKeepAliveInterval |
clean-session | true | MqttConnectOptions.setCleanSession |
client-id | (automatisch) | explizite Client-ID. Ohne Angabe wird PP-<opc_id>-<random6> (≤23 Zeichen) automatisch erzeugt |
MqttConnectOptions.setAutomaticReconnect(true) ist immer aktiv, sodass nach einer vorübergehenden Broker-Unterbrechung automatisch neu
verbunden wird (mit dem Paho-Standard-Backoff).
Adressabbildung
| address | Verhalten |
|---|---|
factory/line1/temp | einzelnes Topic — genau die letzte Nachricht dieses Topics |
device/+/status | Wildcard (eine Ebene) — letzte Nachricht der passenden Topics (zuletzt eingetroffen) |
factory/# | Wildcard (mehrstufig) — letztes Segment, alle untergeordneten Topics |
"" / null | liefert empty |
Wildcard-Hinweis: Der Paho-Callback legt keinen separaten Cache-Eintrag je Wildcard-Muster an;
Schlüssel ist stets das tatsächlich eingetroffene Topic. Auch wenn also mit device/+/status abonniert wurde,
liegen in lastValueByTopic getrennt device/sensor01/status und device/sensor02/status — liest man mit dem Wildcard-Topic selbst, existiert
kein Cache-Eintrag (da das Wildcard selbst kein Publish-Ziel ist). Wird ein Wildcard-Muster wie
device/+/status registriert, muss der Wert des zuletzt eingetroffenen Topics separat nachverfolgt werden —
im Normalbetrieb wird die explizite Registrierung je Topic empfohlen.
write (publish)
client.publish(topic, value.getBytes(), qos, /*retain=*/ false);
- Die Bytes von
valuewerden unverändert als Payload gesendet (UTF-8 als Standard-Charset). retain=falseist fest gesetzt und nicht änderbar — da der Broker die letzte Nachricht nicht aufbewahrt, erhalten später hinzukommende Clients erst beim nächsten Schreibvorgang einen Wert. Für Verbraucher, die Retained Messages benötigen, ist eine Broker-seitige Bridge oder ein separater Publisher vorzusehen.- QoS wird direkt aus der Option übernommen. Bei QoS 1/2 überträgt Paho intern bis zum ACK erneut.
Verbindungs-Lebenszyklus
| Methode | Verhalten |
|---|---|
connect() | Optionen anwenden → Broker-URL zusammensetzen → MqttClient.connect() → tag_map gesammelt abonnieren |
read() | bei Cache-Miss Lazy Subscribe → bei Cache-Hit Wert zurückgeben |
write() | publish (retain=false) |
close() | disconnect() + close() (Best-Effort, Ausnahmen werden ignoriert) → Cache leeren |
connectionLost(cause) | im Callback setConnected(false) — Paho versucht automatisch, neu zu verbinden |
Einschränkungen und geplante Arbeiten
- mTLS / Client-Zertifikate — derzeit nur username/password. Umgebungen mit X.509-mTLS wie AWS IoT
erfordern eine separate Truststore-/Keystore-Konfiguration (Paho
MqttConnectOptions.setSocketFactorynoch nicht nach außen geführt). - Retained / Will Message — derzeit fest
retain=false, kein LWT konfiguriert. - getrennter Wildcard-Cache — ein read auf das Wildcard-Muster selbst liefert ein leeres Ergebnis (Einträge je eingetroffenem Topic).
- JSON-Path-Extraktion — auch bei JSON-Payload gibt der Treiber nur den rohen String zurück.
Die JSONPath-Extraktion muss in der Stufe
format/ Formel erfolgen (oder über einen separaten Parser-Node). - Payload-Encoding — derzeit wird das Standard-Charset der Plattform verwendet. Bei Binär-Payloads
kann es zu UTF-8-Verfälschungen kommen → künftig wird die Freigabe einer Option
payload-encoding(utf8/iso-8859-1/base64) geprüft.
Ort der Tests
test/java/plantpulse/driver/protocol/mqtt_client/MQTTClientDriverTest.java — 31 Tests.
- Ausgangszustand / Metadaten (
isConnected/isWriteSupported/isExternal/getDriverSource) - Sicherheit bei read/write/null/empty ohne Verbindung
getDebugBaseUrl— tcp/ssl, explizite Ports, Standardports (1883/8883)- Optionsparsing — username/password/tls/keep-alive/clean-session/qos/client-id, ungültige Werte → Standard
- statische Helper —
clampQos,defaultClientId(Länge/random/null),nullIfEmpty,parseInt,parseBool - Konstantenprüfung (
DEFAULT_PORT,DEFAULT_PORT_TLS,DEFAULT_KEEP_ALIVE_SEC,DEFAULT_QOS)
Integrationstests gegen einen realen Broker sind separat — diese Unit-Tests laufen ohne externen Broker.