メインコンテンツまでスキップ

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 を 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 は 2 種類の 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 (1 階層) — マッチするトピックの最後のメッセージ (最も新しく到着したもの)
factory/#wildcard (マルチ) — 最終 segment、すべての下位トピック
"" / nullempty を返す

wildcard の注意: Paho の callback は実際に到着した topic ごとに個別 entry をキャッシュするのではなく、 到着した topic 自体が key になります。つまり device/+/status で subscribe していても lastValueByTopic には device/sensor01/statusdevice/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 でもドライバは 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 なしで動作します。