Zum Hauptinhalt springen

Sparkplug B (Phase 1 + 3 aktiv)

Handbuch zum Sparkplug-B-Publisher (outbound) + Befehlsempfänger (inbound) von PlantPulse Edge. NBIRTH / DBIRTH / DDATA / NDEATH werden über die Standard-SPB-1.0-Topics veröffentlicht; vom Host (SCADA) gesendete NCMD / DCMD werden verarbeitet, inklusive Tag-Schreibzugriff und Rebirth. Kompatibel mit SCADA / HiveMQ / Ignition usw.

PunktWert
StatusPhase 1 (BIRTH/DATA/DEATH) + Phase 3 (NCMD/DCMD inbound) aktiv. Phase 2 (alias) optional per Toggle
Controller-Paketplantpulse.app.edge.component.queue.transfer.sparkplug
Standardsparkplug.enable=true (im Betrieb aktiv)
BefehlsempfangBei sparkplug.cmd.enable=true Abonnement + Verarbeitung von NCMD/DCMD
BrokerHiveMQ der Edge (mqtt.server.* wiederverwendet)
AuthentifizierungMQTT username/password (mqtt.server.user / mqtt.server.password aus app.properties). Ab 2026.05+ erzeugt install.sh pro Box automatisch ein Zufallspasswort. auth.properties des HiveMQ-Brokers wird vom Container-Entrypoint bei jedem Start automatisch mit mqtt.* aus app.properties synchronisiert — der Betreiber ändert nur app.properties
Secrets im Containermodus ab 2026.05+

Die MQTT-Konten werden über MQTT_USER / MQTT_PASSWORD_FILE oder MQTT_PASSWORD injiziert. Wenn ein externes SCADA (Ignition / HiveMQ Edge o. Ä.) per Sparkplug auf diese Box zugreift, entnimmt der Betreiber das Passwort der von ihm verwalteten Secret-Datei.


Topic-Schema

spBv1.0/{group_id}/(N|D){BIRTH|DEATH|DATA|CMD}/{edge_node_id}[/{device_id}]
PunktMapping
group_idEdgeContext.site_id (z. B. SITE_00001). Über sparkplug.group.id überschreibbar
edge_node_idEdgeContext.id (z. B. EDGE_00303). Über sparkplug.edge.node.id überschreibbar
device_idopc_id (z. B. OPC_UA_Kepware) — sparkplug.device.strategy=per-opc
metric nametag_id

Nachrichtentypen

NachrichtTopicKern der PayloadZeitpunkt der Veröffentlichung
NBIRTHspBv1.0/{g}/NBIRTH/{n}bdSeq metric, seq=0Direkt nach der Verbindung
DBIRTHspBv1.0/{g}/DBIRTH/{n}/{d}Alle Tag-Metriken des jeweiligen OPC (initial null), seq erhöhtDirekt nach der Verbindung, je OPC
DDATAspBv1.0/{g}/DDATA/{n}/{d}Einzelne Metrik (tag_id, value, datatype, ts), seq erhöhtBei Eintreffen eines Point aus der Queue
NDEATHspBv1.0/{g}/NDEATH/{n}bdSeq metric, kein seq (Spezifikation)MQTT will (unerwartetes Beenden) / explizite Veröffentlichung beim Disconnect
NCMDspBv1.0/{g}/NCMD/{n}Host → Edge-Knotenbefehl (z. B. Node Control/Rebirth=true)Benutzeraktion am Host / Timeout / State-Sync
DCMDspBv1.0/{g}/DCMD/{n}/{d}Host → Edge-Device-Metrik-Write (Ziel-Metrikname oder alias + neuer Wert)Wenn der Host einen bestimmten Tag-Wert eines OPC als Sollwert ändert

Lebenszyklus / Sequenz

  1. bdSeq — 0 ~ 255 wrap, auf Disk persistiert (./data/sparkplug/bdseq Standard). Zu jedem Sessionstart +1.
  2. seq — 0 ~ 255 wrap. Beginnt bei NBIRTH mit 0, wird bei jeder DBIRTH-/DDATA-Veröffentlichung um 1 erhöht.
  3. NDEATH = MQTT Will — wird beim CONNECT als will registriert und vom Broker bei unerwartetem Beenden automatisch veröffentlicht. Bei normalem Beenden veröffentlicht disconnect() explizit.

Datentyp-Mapping

APP_TAG.data_type (freie Zeichenkette) → MetricDataType (Groß-/Kleinschreibung ignoriert, trim).

edgeSPB
Float / SingleFloat
DoubleDouble
Byte / Int8Int8
Short / Int16Int16
Int / Integer / Int32Int32
Long / Int64Int64
UInt8/16/32/64UInt8/16/32/64
Boolean / BoolBoolean
String / TextString
DateTime / DateDateTime
Sonstige / nullString (fallback, log.warn)

Point.value ist als String gespeichert und muss daher vor der Übernahme als SPB-Metrikwert entsprechend obiger Typen gecastet werden. Schlägt das Casting fehl, wird der Metrikwert als null veröffentlicht.


Konfiguration (src/app.properties)

# --------------------------------------------------------------------------------------
# SPARKPLUG B
# 브로커 접속은 mqtt.server.* 를 그대로 재사용한다.
# --------------------------------------------------------------------------------------
sparkplug.enable=false
sparkplug.group.id= # 비우면 site_id
sparkplug.edge.node.id= # 비우면 EdgeContext.id
sparkplug.device.strategy=per-opc # per-opc | single (현재 per-opc 만 구현)
sparkplug.qos=0 # DDATA 빈번 → 0 권장
sparkplug.bdseq.path= # 비우면 ./data/sparkplug/bdseq
sparkplug.alias.enable=true # Phase 2 토글 (현 무영향)
sparkplug.cmd.enable=false # Phase 3 토글 (현 무영향)

Komponenten

DateiFunktion
SparkPlugQueueTransfer.javaQueueTransfer-Implementierung. Nur init/trans/close wird exponiert
SparkPlugSession.javaLebenszyklus (CONNECT+will / NBIRTH / DBIRTH / DDATA / disconnect-NDEATH / seq / bdSeq)
SparkPlugConfig.javaProperties → Konfigurationsobjekt (Broker: mqtt.server.* wiederverwendet)
SparkPlugTopicBuilder.javaBuilder für Topic-Strings
SparkPlugDataTypeMapper.javadata_type → MetricDataType + value coerce
SparkPlugBdSeqStore.javabdSeq-Persistierung auf Disk (0 ~ 255 wrap)
PointQueueProcessorFactory.javaSparkPlugQueueTransfer zur transfer_list hinzugefügt

Abhängige jars

Ergänzung in WebContent/WEB-INF/lib/:

jarFunktion
tahu-core-1.0.14.jarEclipse Tahu Sparkplug B Payload-Modell + Encoder
protobuf-java-3.25.5.jarprotobuf-Abhängigkeit von tahu
commons-compress-1.27.1.jarcommons-compress-Abhängigkeit von tahu
lib-Synchronisierung

Die jars unter WEB-INF/lib werden beim Build aus Nexus materialisiert und daher nicht committet. Für neue Abhängigkeiten müssen die Koordinaten in build.gradle ergänzt werden.


Verifikation (dev-Server / 2026-05-06)

Capture von spBv1.0/# mit dem paho-mqtt-Subscriber:

[SUMMARY] total=511
NDEATH: 1 ← 재시작 시 will 자동 발행
NBIRTH: 1 ← seq=0
DBIRTH: 10 ← OPC 10개 모두, seq=1..10
DDATA: 499 ← 13초간 ≈ 38 msgs/s

Die Anzahl der DBIRTH-Metriken stimmt exakt mit der Anzahl der APP_TAG-Tags des jeweiligen OPC überein (OPC_UA_Kepware=11, OPC_00303=56, OPC_LS_XBM_0001=28 usw.).

Beispiel für ein DDATA-Topic:

spBv1.0/SITE_00001/DDATA/EDGE_00303/OPC_UA_Kepware bytes=44 seq=125 ts=...

Werkzeug zur Prüfung der Payload-Dekodierung:

MQTT_USER="${MQTT_USER:-edge}" \
MQTT_PASSWORD="$(tr -d '\r\n' < /run/secrets/mqtt-password)" \
python3 /tmp/spb_sub.py 100.106.92.21 1883 "$MQTT_USER" "$MQTT_PASSWORD" 60

(Dieses Skript ist eine temporäre Datei und nicht Bestandteil der Build-Artefakte.)


Phase 3 — Verhalten beim Befehlsempfang (NCMD/DCMD)

Bei sparkplug.cmd.enable=true abonniert die SparkPlugSession automatisch die folgenden Topics.

spBv1.0/{group}/NCMD/{edge_node}
spBv1.0/{group}/DCMD/{edge_node}/+ # 모든 device 와일드카드

Eingehende Nachrichten werden mit der statischen Methode parseCmdTopic(String) → CmdTopic validiert und anschließend verzweigt verarbeitet.

BefehlVerarbeitung
NCMD Node Control/Rebirth=truerebirthAsync() — erneute NBIRTH-Veröffentlichung + erneute DBIRTH-Veröffentlichung aller OPC + (bei aktiviertem alias) Aktualisierung von aliasReverseMap (refreshAliasReverseMap())
DCMD metric writeRückwärtssuche der tag_id über metric name oder alias → Aufruf von write(tagId, value) des Treibers des jeweiligen OPC → Ergebnis wird im nächsten DDATA-Zyklus übernommen

Topic-Validierung (parseCmdTopic)

EingabeErgebnis
spBv1.0/Plant1/NCMD/EDGE_00303OK (NCMD, deviceId=null)
spBv1.0/Plant1/DCMD/EDGE_00303/OPC_LSOK (DCMD, deviceId=OPC_LS)
null / leere Zeichenkette / weniger als 4 Segmentereject (null)
spBv1.0/Plant1/NDATA/EDGE_00303reject (outbound-Typ)
Anderer Namespace (spBv2.0/...)reject
spBv1.0/Plant1/DCMD/EDGE_00303 (4 seg, deviceId fehlt)reject
6+ SegmenteZusätzliche Segmente werden ignoriert (nur 4-5 werden genutzt)

→ Der Unit-Test SparkPlugCommandTopicTest deckt die obigen 10 Fälle ab.

Rebirth-Sicherung

Der Fehler, bei dem rebirthAsync() nach der erneuten NBIRTH-/DBIRTH-Veröffentlichung die Aktualisierung von aliasReverseMap ausließ und das nächste DCMD daher mit einem veralteten alias nachschlug, ist behoben — Aufruf von refreshAliasReverseMap() ergänzt (subscribeCommands() nutzt denselben Helper, DRY).


Roadmap

PhaseInhaltStatus
1 — BIRTH/DATA/DEATHNBIRTH / DBIRTH / DDATA / NDEATH(will)✅ abgeschlossen
2 — alias-OptimierungErzeugung der tag_id → alias(Long)-Tabelle zum Zeitpunkt von NBIRTH/DBIRTH, DDATA veröffentlicht nur noch alias (Payload-Reduktion um 30 ~ 40 %). Toggle: sparkplug.alias.enableOptional (Toggle)
3 — NCMD/DCMD inboundVerarbeitung der vom Host gesendeten Metrik-Write-Befehle + Node Control/Rebirth✅ abgeschlossen (sparkplug.cmd.enable=true)
4 — STATE / BetriebstransparenzAbonnement von STATE/{primary_host_id} → DDATA wird zurückgehalten, wenn primary OFFLINE ist. Ergänzung des Abschnitts sparkplug in der /api/v1/edge-Antwort (connected / bdSeq / seq / devices)TBD

Bekannte Einschränkungen

  • isNull-Flag wird bei DDATA value=null nicht gesetzt: Derzeit wird nur Metric.setValue(null) aufgerufen. Laut Spezifikation wäre die explizite Angabe von isNull=true sicherer → Kandidat für eine spätere Nachbesserung.
  • Eine einzelne Metrik pro DDATA: Da die Queue jeweils einen Eintrag überträgt, findet keine Optimierung durch gebündelte Veröffentlichung statt. Bei höherer Frequenz ist eine Weiterentwicklung zu batched DDATA möglich.
  • bdSeq-Pfad standardmäßig relativ zum Arbeitsverzeichnis: Auf Produktivservern wird ein expliziter absoluter Pfad (/opt/kopens/.../data/sparkplug/bdseq) empfohlen.
  • Keine Fehlerantwort bei fehlgeschlagenem DCMD-Write: Derzeit wird dem Host auch bei Fehlschlagen von write() des Treibers kein Ergebnis per NDATA/DDATA gemeldet. Die Ergänzung einer result_status-Metrik wird künftig geprüft.