跳到主要内容

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.* 自动同步 —— 运维人员只需修改 app.properties
2026.05+ 容器模式的密钥

MQTT 账号通过 MQTT_USER / MQTT_PASSWORD_FILEMQTT_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 覆盖
edge_node_idEdgeContext.id (例: EDGE_00303)。可用 sparkplug.edge.node.id 覆盖
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(初始 null),seq 递增连接建立后按 OPC 逐个发布
DDATAspBv1.0/{g}/DDATA/{n}/{d}单个 metric(tag_id, value, datatype, ts),seq 递增队列中到达 1 条 Point 时
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 写入(目标 metric name 或 alias + 新值)host 将某 OPC 的特定标签值改为 setpoint 时

生命周期 / 时序

  1. bdSeq — 0 ~ 255 回绕,持久化到磁盘(默认 ./data/sparkplug/bdseq)。每次会话启动时 +1
  2. seq — 0 ~ 255 回绕。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属性 → 配置对象(broker 复用 mqtt.server.*)
SparkPlugTopicBuilder.java主题字符串构建器
SparkPlugDataTypeMapper.javadata_type → MetricDataType + value coerce
SparkPlugBdSeqStore.javabdSeq 磁盘持久化(0 ~ 255 回绕)
PointQueueProcessorFactory.java在 transfer_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 write通过 metric 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 段reject (null)
spBv1.0/Plant1/NDATA/EDGE_00303reject(出站类型)
其他 namespace (spBv2.0/...)reject
spBv1.0/Plant1/DCMD/EDGE_00303 (4 段,缺少 deviceId)reject
6 段及以上忽略多余段(仅使用第 4-5 段)

→ 单元测试 SparkPlugCommandTopicTest 覆盖上述 10 个用例。

rebirth 安全网

原先 rebirthAsync() 在重发 NBIRTH/DBIRTH 后遗漏刷新 aliasReverseMap,导致后续 DCMD 使用过期 alias 查找的缺陷已修复 —— 新增 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: 队列逐条传输,因此未做批量发布优化。若频率进一步提高, 可演进为 batched DDATA。
  • bdSeq path 默认是相对于工作目录的路径: 生产服务器建议显式使用绝对路径 (/opt/kopens/.../data/sparkplug/bdseq)。
  • 未实现 DCMD 写入失败响应: 当前即使 driver write() 失败,也不会通过 NDATA/DDATA 向 host 通报结果。 后续考虑新增 result_status metric。