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 구독 + 처리 |
| 브로커 | 엣지의 HiveMQ (mqtt.server.* 재사용) |
| 인증 | MQTT username/password (app.properties 의 mqtt.server.user / mqtt.server.password). 2026.05+ install.sh 가 박스마다 random 비번 자동 생성. 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 (target 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 tag 수와 정확히 일치 (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명시가 더 안전 → 후속 보강 후보. - 단일 metric per DDATA: 큐가 1건씩 trans 하므로 묶음 발행 최적화는 안 함. 빈도가 더 높아지면 batched DDATA 로 발전 가능.
- bdSeq path 디폴트가 작업디렉토리 상대경로: 실서버에서는 명시 절대경로
(
/opt/kopens/.../data/sparkplug/bdseq) 권장. - DCMD write 실패 응답 미구현: 현재 driver
write()가 실패해도 host 에 NDATA/DDATA 로 결과 통보를 하지 않습니다. 향후result_statusmetric 추가 검토.