MQTT Client Driver — 기술 레퍼런스
PlantPulse Edge 의 MQTT 클라이언트 드라이버는 Eclipse Paho mqttv3 1.2.5 라이브러리를 래핑한 자바 구현입니다. 외부 IIoT broker (HiveMQ / Mosquitto / EMQX / AWS IoT Core / Azure IoT Hub 등) 의 topic 을 subscribe 하여 메시지를 in-memory 캐시에 보관하고, 태그의 polling cycle 마다 캐시값을 반환합니다. write 는 즉시 publish.
소스: plantpulse.driver.protocol.mqtt_client.*
| 클래스 | 책임 |
|---|---|
MQTTClientDriver | ProtocolDriver 구현체 — connect/read/write/close, topic 캐시, lazy subscribe |
라이브러리: lib/org.eclipse.paho.client.mqttv3-1.2.5.jar
Sparkplug B 와의 차이
PlantPulse 는 두 가지 MQTT 기반 드라이버 를 별도로 제공합니다.
| 드라이버 | payload | 사용 시점 |
|---|---|---|
MQTT_CLIENT (이 페이지) | 임의 String/JSON/binary | 일반 IIoT broker, 사용자 정의 topic, ASCII / JSON payload |
SPARKPLUG_B | Sparkplug B Payload (Protobuf) | Cirrus Link / Tahu / Ignition 표준, NBIRTH/DBIRTH alias |
Sparkplug B 는 자체 wire codec + alias 학습 로직이 들어 있어, 일반 MQTT 만 필요한 환경에서는 이 드라이버 쪽이 가볍고 표준 Paho API 만으로 충분합니다.
동작 흐름
broker ──PUBLISH──> Paho MqttClient ──MqttCallback.messageArrived(topic, payload)──> Driver
│
└─ lastValueByTopic.put(topic, new String(payload))
PlantPulse 폴러 ──read(addr)──> lastValueByTopic[addr] (in-memory hit)
write(addr, value) ──Paho.publish(topic, value, qos, retain=false)──> broker
read 의 첫 호출에서 cache miss 면 lazy subscribe (QoS 옵션값) 후 빈 문자열 반환 — 다음 cycle 부터
실제 값이 들어옵니다. tag_map 에 등록된 topic 들은 connect 직후 preSubscribeFromTagMap() 으로
일괄 subscribe 되어 첫 cycle 에 바로 값이 도착할 수도 있습니다 (메시지가 publish 빈도에 따름).
옵션
| 키 | 기본값 | 설명 |
|---|---|---|
username | (없음) | broker 인증 사용자명 |
password | (없음) | broker 인증 비밀번호 |
tls | false | true → ssl:// (default 8883), false → tcp:// (default 1883) |
qos | 0 | subscribe / publish QoS (0 / 1 / 2). 범위 외 값은 자동 clamp. |
keep-alive | 60 | keep-alive 주기 (초) — MqttConnectOptions.setKeepAliveInterval |
clean-session | true | MqttConnectOptions.setCleanSession |
client-id | (자동) | 명시 client id. 미지정 시 PP-<opc_id>-<random6> (≤23자) 자동 생성 |
MqttConnectOptions.setAutomaticReconnect(true) 가 항상 켜져 있어 일시적 broker 단절 후
자동 재연결됩니다 (Paho 기본 backoff 적용).
Address 매핑
| address | 동작 |
|---|---|
factory/line1/temp | 단일 토픽 — 정확히 이 topic 의 마지막 메시지 |
device/+/status | wildcard (한 단계) — 매칭되는 토픽들의 마지막 메시지 (가장 최근 도착) |
factory/# | wildcard (멀티) — 마지막 segment, 모든 하위 토픽 |
"" / null | empty 반환 |
wildcard 주의: Paho 의 callback 은 실제 도착 topic 별로 별도 entry 를 캐시하지 않고,
도착한 topic 자체가 key 가 됩니다. 즉 device/+/status 로 subscribe 했더라도
lastValueByTopic 에는 device/sensor01/status, device/sensor02/status 가 따로 들어가고,
사용자가 wildcard 토픽 자체로 read 하면 cache 에는 entry 가 없습니다 (wildcard 자체는
publish 의 destination 이 아니므로). wildcard 패턴을 등록할 때는 device/+/status 처럼 등록한 뒤
read 의 결과로는 마지막 도착 토픽의 값을 별도로 추적해야 합니다 — 일반 사용에서는 토픽별로
명시 등록을 권장합니다.
write (publish)
client.publish(topic, value.getBytes(), qos, /*retain=*/ false);
value의 바이트가 그대로 payload 로 전송됩니다 (UTF-8 기본 charset).retain=false로 고정되어 있고 바꿀 수 없습니다 — broker 가 마지막 메시지를 보존하지 않으므로, 나중에 구독을 시작한 클라이언트는 다음 쓰기가 일어날 때까지 값을 받지 못합니다. retained 메시지가 필요한 소비자에게는 broker 쪽 bridge 나 별도 publisher 를 두세요.- QoS 는 옵션값 그대로 사용. QoS 1/2 면 Paho 가 내부적으로 ack 까지 재전송.
연결 라이프사이클
| 메서드 | 동작 |
|---|---|
connect() | options 적용 → broker URL 조립 → MqttClient.connect() → tag_map 일괄 subscribe |
read() | cache miss 시 lazy subscribe → cache hit 면 값 반환 |
write() | publish (retain=false) |
close() | disconnect() + close() (best-effort, 예외 무시) → cache clear |
connectionLost(cause) | callback 에서 setConnected(false) — Paho 가 자동 재연결 시도 |
한계 및 향후 작업
- mTLS / 클라이언트 인증서 — 현재 username/password 만. AWS IoT 같은 X.509 mTLS 환경은 별도
truststore / keystore 설정 필요 (Paho
MqttConnectOptions.setSocketFactory노출 미구현). - retained / will message — 현재
retain=false고정, LWT 미설정. - wildcard cache 분리 — wildcard 패턴 자체로 read 하면 결과가 비어 있음 (도착 topic 별 entry).
- JSON path 추출 — payload 가 JSON 이라도 driver 는 raw String 만 반환. JSONPath 추출은
format/ fomula 단계에서 처리 필요 (또는 별도 파서 노드). - payload encoding — 현재 platform default charset 사용. binary payload 는 UTF-8 깨짐 가능 →
추후
payload-encoding옵션 (utf8/iso-8859-1/base64) 노출 검토.
테스트 위치
test/java/plantpulse/driver/protocol/mqtt_client/MQTTClientDriverTest.java — 31 테스트.
- 초기 상태 / 메타 (
isConnected/isWriteSupported/isExternal/getDriverSource) - 미연결 read/write/null/empty 안전
getDebugBaseUrl— tcp/ssl, 명시 포트, default 포트 (1883/8883)- 옵션 파싱 — username/password/tls/keep-alive/clean-session/qos/client-id, 잘못된 값 → default
- 정적 헬퍼 —
clampQos,defaultClientId(길이/random/null),nullIfEmpty,parseInt,parseBool - 상수 검증 (
DEFAULT_PORT,DEFAULT_PORT_TLS,DEFAULT_KEEP_ALIVE_SEC,DEFAULT_QOS)
실 broker 의존 통합 테스트는 별도 — 본 단위 테스트는 외부 broker 없이 동작합니다.