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.*
| クラス | 責務 |
|---|---|
MQTTClientDriver | ProtocolDriver 実装 — 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_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 (1 階層) — マッチするトピックの最後のメッセージ (最も新しく到着したもの) |
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 でもドライバは 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 なしで動作します。