Sparkplug B (Phase 1 + 3 有効)
PlantPulse Edge の Sparkplug B パブリッシャ(outbound) + コマンド受信(inbound) マニュアルです。標準 SPB 1.0 トピックで NBIRTH / DBIRTH / DDATA / NDEATH を発行し、ホスト(SCADA) から送られる NCMD / DCMD を処理してタグ書き込みおよび rebirth までサポートします。SCADA / HiveMQ / Ignition などと互換です。
| 項目 | 値 |
|---|---|
| 状態 | Phase 1 (BIRTH/DATA/DEATH) + Phase 3 (NCMD/DCMD インバウンド) 有効。Phase 2 (alias) はトグルによるオプション |
| コントローラパッケージ | plantpulse.app.edge.component.queue.transfer.sparkplug |
| デフォルト | sparkplug.enable=true (運用時有効) |
| コマンド受信 | sparkplug.cmd.enable=true の場合に NCMD/DCMD を購読 + 処理 |
| ブローカー | Edge の HiveMQ (mqtt.server.* を再利用) |
| 認証 | MQTT username/password (app.properties の mqtt.server.user / mqtt.server.password)。2026.05+ の install.sh がボックスごとにランダムパスワードを自動生成。HiveMQ broker の auth.properties は container-entrypoint が起動ごとに app.properties の mqtt.* と自動 sync — 運用者は app.properties のみ変更 |
MQTT アカウントは MQTT_USER / MQTT_PASSWORD_FILE または MQTT_PASSWORD で注入します。
外部 SCADA (Ignition / HiveMQ Edge など) が Sparkplug でこのボックスに接続する場合は、運用者が管理する
secret ファイルでパスワードを確認してから使用します。
トピックスキーマ
spBv1.0/{group_id}/(N|D){BIRTH|DEATH|DATA|CMD}/{edge_node_id}[/{device_id}]
| 項目 | マッピング |
|---|---|
group_id | EdgeContext.site_id (例: SITE_00001)。sparkplug.group.id で override 可能 |
edge_node_id | EdgeContext.id (例: EDGE_00303)。sparkplug.edge.node.id で override 可能 |
device_id | opc_id (例: OPC_UA_Kepware) — sparkplug.device.strategy=per-opc |
metric name | tag_id |
メッセージ種別
| メッセージ | トピック | ペイロードの要点 | 発行タイミング |
|---|---|---|---|
| NBIRTH | spBv1.0/{g}/NBIRTH/{n} | bdSeq metric、seq=0 | 接続直後 |
| DBIRTH | spBv1.0/{g}/DBIRTH/{n}/{d} | 該当 OPC の全タグ metric (initial null)、seq をインクリメント | 接続直後に OPC 単位 |
| DDATA | spBv1.0/{g}/DDATA/{n}/{d} | 単一 metric (tag_id, value, datatype, ts)、seq をインクリメント | キューの Point 1 件到着時 |
| NDEATH | spBv1.0/{g}/NDEATH/{n} | bdSeq metric、seq なし (仕様) | MQTT will (異常終了) / disconnect 時に明示発行 |
| NCMD | spBv1.0/{g}/NCMD/{n} | host → edge ノードコマンド (例: Node Control/Rebirth=true) | host のユーザー操作 / タイムアウト / state sync |
| DCMD | spBv1.0/{g}/DCMD/{n}/{d} | host → edge デバイス metric write (対象 metric name または alias + new value) | host が OPC の特定タグ値を setpoint として変更するとき |
ライフサイクル / シーケンス
bdSeq— 0 ~ 255 wrap、ディスク永続化 (./data/sparkplug/bdseqがデフォルト)。セッション開始ごとに+1。seq— 0 ~ 255 wrap。NBIRTH で 0 から開始し、DBIRTH/DDATA いずれも発行時点で 1 ずつインクリメント。- NDEATH = MQTT Will — CONNECT 時に will として登録され、異常終了時に broker が自動発行。正常終了時は
disconnect()が明示発行します。
データ型マッピング
APP_TAG.data_type (自由文字列) → MetricDataType (大文字小文字を無視、trim)。
| edge | SPB |
|---|---|
| Float / Single | Float |
| Double | Double |
| Byte / Int8 | Int8 |
| Short / Int16 | Int16 |
| Int / Integer / Int32 | Int32 |
| Long / Int64 | Int64 |
| UInt8/16/32/64 | UInt8/16/32/64 |
| Boolean / Bool | Boolean |
| String / Text | String |
| DateTime / Date | DateTime |
| その他 / null | String (fallback、log.warn) |
Point.value は String として保存されているため、SPB metric value に入れる前に上記の型へキャストします。
キャスト失敗時は metric 値を null として発行します。
設定 (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 토글 (현 무영향)
構成要素
| ファイル | 役割 |
|---|---|
SparkPlugQueueTransfer.java | QueueTransfer 実装。init/trans/close のみ公開 |
SparkPlugSession.java | ライフサイクル (CONNECT+will / NBIRTH / DBIRTH / DDATA / disconnect-NDEATH / seq / bdSeq) |
SparkPlugConfig.java | プロパティ → 設定オブジェクト (ブローカーは mqtt.server.* を再利用) |
SparkPlugTopicBuilder.java | トピック文字列ビルダー |
SparkPlugDataTypeMapper.java | data_type → MetricDataType + value coerce |
SparkPlugBdSeqStore.java | bdSeq のディスク永続化 (0 ~ 255 wrap) |
PointQueueProcessorFactory.java | transfer_list に SparkPlugQueueTransfer を追加 |
依存 jar
WebContent/WEB-INF/lib/ に追加:
| jar | 役割 |
|---|---|
tahu-core-1.0.14.jar | Eclipse Tahu Sparkplug B ペイロードモデル + エンコーダ |
protobuf-java-3.25.5.jar | tahu の protobuf 依存 |
commons-compress-1.27.1.jar | tahu の commons-compress 依存 |
WEB-INF/lib の jar はビルド時に Nexus から materialize されるためコミットしません。新規依存は build.gradle に座標を追加
する必要があります。
検証 (dev サーバー / 2026-05-06)
paho-mqtt subscriber で spBv1.0/# をキャプチャ:
[SUMMARY] total=511
NDEATH: 1 ← 재시작 시 will 자동 발행
NBIRTH: 1 ← seq=0
DBIRTH: 10 ← OPC 10개 모두, seq=1..10
DDATA: 499 ← 13초간 ≈ 38 msgs/s
DBIRTH のメトリック数 = APP_TAG の該当 OPC タグ数と完全に一致 (OPC_UA_Kepware=11、OPC_00303=56、
OPC_LS_XBM_0001=28 など)。
DDATA トピック例:
spBv1.0/SITE_00001/DDATA/EDGE_00303/OPC_UA_Kepware bytes=44 seq=125 ts=...
ペイロードのデコード確認ツール:
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
(このスクリプトは一時ファイルであり、ビルド成果物には含まれません。)
Phase 3 — コマンド(NCMD/DCMD)受信動作
sparkplug.cmd.enable=true の場合、SparkPlugSession が次のトピックを自動購読します。
spBv1.0/{group}/NCMD/{edge_node}
spBv1.0/{group}/DCMD/{edge_node}/+ # 모든 device 와일드카드
受信したメッセージは静的メソッド parseCmdTopic(String) → CmdTopic で検証した後、分岐処理されます。
| コマンド | 処理 |
|---|---|
NCMD Node Control/Rebirth=true | rebirthAsync() — NBIRTH 再発行 + 全 OPC の DBIRTH 再発行 + (alias 有効時) aliasReverseMap の更新 (refreshAliasReverseMap()) |
DCMD metric write | metric name または alias から tag_id を逆引き → 該当 OPC の driver write(tagId, value) を呼び出し → 結果を次の DDATA サイクルに反映 |
トピック検証 (parseCmdTopic)
| 入力 | 結果 |
|---|---|
spBv1.0/Plant1/NCMD/EDGE_00303 | OK (NCMD, deviceId=null) |
spBv1.0/Plant1/DCMD/EDGE_00303/OPC_LS | OK (DCMD, deviceId=OPC_LS) |
null / 空文字列 / 4 segment 未満 | reject (null) |
spBv1.0/Plant1/NDATA/EDGE_00303 | reject (outbound タイプ) |
別 namespace (spBv2.0/...) | reject |
spBv1.0/Plant1/DCMD/EDGE_00303 (4 seg、deviceId 欠落) | reject |
| 6+ segment | 追加 segment を無視 (4-5 のみ使用) |
→ 単体テスト SparkPlugCommandTopicTest が上記 10 ケースを保証します。
rebirth のセーフティネット
rebirthAsync() が NBIRTH/DBIRTH 再発行後に aliasReverseMap の更新を漏らし、次の DCMD が stale alias で lookup していた不具合は修正済みです — refreshAliasReverseMap() 呼び出しを追加 (subscribeCommands() も同じヘルパーを使用、DRY)。
ロードマップ
| Phase | 内容 | 状態 |
|---|---|---|
| 1 — BIRTH/DATA/DEATH | NBIRTH / DBIRTH / DDATA / NDEATH(will) | ✅ 完了 |
| 2 — alias 最適化 | NBIRTH/DBIRTH の時点で tag_id → alias(Long) テーブルを生成、DDATA は alias のみ発行 (ペイロード 30 ~ 40% 削減)。トグル: sparkplug.alias.enable | オプション (toggle) |
| 3 — NCMD/DCMD インバウンド | host が送る metric write コマンドの処理 + Node Control/Rebirth | ✅ 完了 (sparkplug.cmd.enable=true) |
| 4 — STATE / 運用可視性 | STATE/{primary_host_id} を購読 → primary が OFFLINE の場合 DDATA を保留。/api/v1/edge のレスポンスに sparkplug セクションを追加 (connected / bdSeq / seq / devices) | TBD |
既知の制限
- DDATA value=null 時に isNull フラグが未設定: 現在は
Metric.setValue(null)のみを呼び出しています。仕様上はisNull=trueを明示するほうが安全 → 今後の補強候補。 - DDATA あたり単一 metric: キューが 1 件ずつ trans するため、まとめ発行の最適化は行っていません。頻度がさらに高くなれば batched DDATA へ発展可能です。
- bdSeq path のデフォルトが作業ディレクトリ相対パス: 実サーバーでは明示的な絶対パス
(
/opt/kopens/.../data/sparkplug/bdseq) を推奨します。 - DCMD write 失敗応答が未実装: 現在は driver
write()が失敗しても host に NDATA/DDATA で結果を通知 しません。今後result_statusmetric の追加を検討します。