본문으로 건너뛰기

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.*

클래스책임
MQTTClientDriverProtocolDriver 구현체 — 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_BSparkplug 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 인증 비밀번호
tlsfalsetruessl:// (default 8883), falsetcp:// (default 1883)
qos0subscribe / publish QoS (0 / 1 / 2). 범위 외 값은 자동 clamp.
keep-alive60keep-alive 주기 (초) — MqttConnectOptions.setKeepAliveInterval
clean-sessiontrueMqttConnectOptions.setCleanSession
client-id(자동)명시 client id. 미지정 시 PP-<opc_id>-<random6> (≤23자) 자동 생성

MqttConnectOptions.setAutomaticReconnect(true) 가 항상 켜져 있어 일시적 broker 단절 후 자동 재연결됩니다 (Paho 기본 backoff 적용).


Address 매핑

address동작
factory/line1/temp단일 토픽 — 정확히 이 topic 의 마지막 메시지
device/+/statuswildcard (한 단계) — 매칭되는 토픽들의 마지막 메시지 (가장 최근 도착)
factory/#wildcard (멀티) — 마지막 segment, 모든 하위 토픽
"" / nullempty 반환

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 없이 동작합니다.