跳到主要内容

MQTT Client Driver — 技术参考

PlantPulse Edge 的 MQTT 客户端驱动是对 Eclipse Paho mqttv3 1.2.5 库进行封装的 Java 实现。 它订阅外部 IIoT broker(HiveMQ / Mosquitto / EMQX / AWS IoT Core / Azure IoT Hub 等)的 topic,将消息保存在内存缓存中, 并在标签的每个 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,涵盖所有子主题
"" / null返回 empty

wildcard 注意事项:Paho 的 callback 并不会按 wildcard 模式缓存条目,而是以实际到达的 topic 作为 key。 也就是说,即使以 device/+/status 进行 subscribe, lastValueByTopic 中也会分别写入 device/sensor01/statusdevice/sensor02/status 两个条目; 若用户以 wildcard 主题 本身执行 read,cache 中并不存在对应 entry(因为 wildcard 本身 不是 publish 的目标)。注册 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() → 批量 subscribe tag_map
read()cache miss 时 lazy subscribe → cache hit 则返回值
write()publish (retain=false)
close()disconnect() + close()(best-effort,忽略异常)→ 清空 cache
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 存储)。
  • JSON path 提取 —— 即使 payload 是 JSON,驱动也只返回 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
  • 静态辅助方法 —— clampQosdefaultClientId(长度/random/null)、nullIfEmptyparseIntparseBool
  • 常量校验(DEFAULT_PORTDEFAULT_PORT_TLSDEFAULT_KEEP_ALIVE_SECDEFAULT_QOS

依赖真实 broker 的集成测试另行提供 —— 本单元测试无需外部 broker 即可运行。