メインコンテンツまでスキップ

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.propertiesmqtt.server.user / mqtt.server.password)。2026.05+ の install.sh がボックスごとにランダムパスワードを自動生成。HiveMQ broker の auth.properties は container-entrypoint が起動ごとに app.properties の mqtt.* と自動 sync — 運用者は app.properties のみ変更
2026.05+ コンテナモードのシークレット

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_idEdgeContext.site_id (例: SITE_00001)。sparkplug.group.id で override 可能
edge_node_idEdgeContext.id (例: EDGE_00303)。sparkplug.edge.node.id で override 可能
device_idopc_id (例: OPC_UA_Kepware) — sparkplug.device.strategy=per-opc
metric nametag_id

メッセージ種別

メッセージトピックペイロードの要点発行タイミング
NBIRTHspBv1.0/{g}/NBIRTH/{n}bdSeq metric、seq=0接続直後
DBIRTHspBv1.0/{g}/DBIRTH/{n}/{d}該当 OPC の全タグ metric (initial null)、seq をインクリメント接続直後に OPC 単位
DDATAspBv1.0/{g}/DDATA/{n}/{d}単一 metric (tag_id, value, datatype, ts)、seq をインクリメントキューの Point 1 件到着時
NDEATHspBv1.0/{g}/NDEATH/{n}bdSeq metric、seq なし (仕様)MQTT will (異常終了) / disconnect 時に明示発行
NCMDspBv1.0/{g}/NCMD/{n}host → edge ノードコマンド (例: Node Control/Rebirth=true)host のユーザー操作 / タイムアウト / state sync
DCMDspBv1.0/{g}/DCMD/{n}/{d}host → edge デバイス metric write (対象 metric name または alias + new value)host が OPC の特定タグ値を setpoint として変更するとき

ライフサイクル / シーケンス

  1. bdSeq — 0 ~ 255 wrap、ディスク永続化 (./data/sparkplug/bdseq がデフォルト)。セッション開始ごとに +1
  2. seq — 0 ~ 255 wrap。NBIRTH で 0 から開始し、DBIRTH/DDATA いずれも発行時点で 1 ずつインクリメント。
  3. NDEATH = MQTT Will — CONNECT 時に will として登録され、異常終了時に broker が自動発行。正常終了時は disconnect() が明示発行します。

データ型マッピング

APP_TAG.data_type (自由文字列) → MetricDataType (大文字小文字を無視、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
その他 / nullString (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.javaQueueTransfer 実装。init/trans/close のみ公開
SparkPlugSession.javaライフサイクル (CONNECT+will / NBIRTH / DBIRTH / DDATA / disconnect-NDEATH / seq / bdSeq)
SparkPlugConfig.javaプロパティ → 設定オブジェクト (ブローカーは mqtt.server.* を再利用)
SparkPlugTopicBuilder.javaトピック文字列ビルダー
SparkPlugDataTypeMapper.javadata_type → MetricDataType + value coerce
SparkPlugBdSeqStore.javabdSeq のディスク永続化 (0 ~ 255 wrap)
PointQueueProcessorFactory.javatransfer_list に SparkPlugQueueTransfer を追加

依存 jar

WebContent/WEB-INF/lib/ に追加:

jar役割
tahu-core-1.0.14.jarEclipse Tahu Sparkplug B ペイロードモデル + エンコーダ
protobuf-java-3.25.5.jartahu の protobuf 依存
commons-compress-1.27.1.jartahu の commons-compress 依存
lib の同期

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=truerebirthAsync() — NBIRTH 再発行 + 全 OPC の DBIRTH 再発行 + (alias 有効時) aliasReverseMap の更新 (refreshAliasReverseMap())
DCMD metric writemetric name または alias から tag_id を逆引き → 該当 OPC の driver write(tagId, value) を呼び出し → 結果を次の DDATA サイクルに反映

トピック検証 (parseCmdTopic)

入力結果
spBv1.0/Plant1/NCMD/EDGE_00303OK (NCMD, deviceId=null)
spBv1.0/Plant1/DCMD/EDGE_00303/OPC_LSOK (DCMD, deviceId=OPC_LS)
null / 空文字列 / 4 segment 未満reject (null)
spBv1.0/Plant1/NDATA/EDGE_00303reject (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/DEATHNBIRTH / 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_status metric の追加を検討します。