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.* 自动同步 —— 运维人员只需修改 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 覆盖 |
edge_node_id | EdgeContext.id (例: EDGE_00303)。可用 sparkplug.edge.node.id 覆盖 |
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(初始 null),seq 递增 | 连接建立后按 OPC 逐个发布 |
| DDATA | spBv1.0/{g}/DDATA/{n}/{d} | 单个 metric(tag_id, value, datatype, ts),seq 递增 | 队列中到达 1 条 Point 时 |
| 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 写入(目标 metric name 或 alias + 新值) | host 将某 OPC 的特定标签值改为 setpoint 时 |
生命周期 / 时序
bdSeq— 0 ~ 255 回绕,持久化到磁盘(默认./data/sparkplug/bdseq)。每次会话启动时+1。seq— 0 ~ 255 回绕。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 | 属性 → 配置对象(broker 复用 mqtt.server.*) |
SparkPlugTopicBuilder.java | 主题字符串构建器 |
SparkPlugDataTypeMapper.java | data_type → MetricDataType + value coerce |
SparkPlugBdSeqStore.java | bdSeq 磁盘持久化(0 ~ 255 回绕) |
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 段 | reject (null) |
spBv1.0/Plant1/NDATA/EDGE_00303 | reject(出站类型) |
其他 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/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: 队列逐条传输,因此未做批量发布优化。若频率进一步提高, 可演进为 batched DDATA。
- bdSeq path 默认是相对于工作目录的路径: 生产服务器建议显式使用绝对路径
(
/opt/kopens/.../data/sparkplug/bdseq)。 - 未实现 DCMD 写入失败响应: 当前即使 driver
write()失败,也不会通过 NDATA/DDATA 向 host 通报结果。 后续考虑新增result_statusmetric。