Zum Hauptinhalt springen

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.*

KlasseVerantwortung
MQTTClientDriverProtocolDriver-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.

TreiberPayloadEinsatzfall
MQTT_CLIENT (diese Seite)beliebiger String/JSON/Binärinhaltallgemeine IIoT-Broker, benutzerdefinierte Topics, ASCII-/JSON-Payload
SPARKPLUG_BSparkplug 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üsselStandardwertBeschreibung
username(keiner)Benutzername für die Broker-Authentifizierung
password(keiner)Passwort für die Broker-Authentifizierung
tlsfalsetruessl:// (Standard 8883), falsetcp:// (Standard 1883)
qos0QoS für Subscribe/Publish (0 / 1 / 2). Werte außerhalb des Bereichs werden automatisch begrenzt.
keep-alive60Keep-Alive-Intervall (Sekunden) — MqttConnectOptions.setKeepAliveInterval
clean-sessiontrueMqttConnectOptions.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

addressVerhalten
factory/line1/tempeinzelnes Topic — genau die letzte Nachricht dieses Topics
device/+/statusWildcard (eine Ebene) — letzte Nachricht der passenden Topics (zuletzt eingetroffen)
factory/#Wildcard (mehrstufig) — letztes Segment, alle untergeordneten Topics
"" / nullliefert 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 value werden unverändert als Payload gesendet (UTF-8 als Standard-Charset).
  • retain=false ist 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

MethodeVerhalten
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.setSocketFactory noch 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.