1. はじめに — なぜ工場 IoT に MQTT なのか
OPC UA や Modbus が「機器とコントローラの会話」のためのプロトコルであるのに対し、MQTT は「分散したデバイスから中央へデータを集約する」ための非同期メッセージングプロトコルです。1999 年に IBM の Andy Stanford-Clark 氏と Arcom(現 Eurotech)の Arlen Nipper 氏が考案した経緯から、長く石油・ガス業界で使われてきましたが、IoT の文脈で爆発的に普及しました。
Pub/Sub モデルのおかげで、Publisher と Subscriber が直接知り合う必要がなく、ブローカ(Broker)を経由して疎結合に通信します。「設備 A の温度を、MES と Grafana と AWS の 3 つに同時に送りたい」といった要件が、配線変更なしで実現します(ここで言う「配線変更なし」はデータの配信先を増やすソフト側の話です。センサー側の配線・電源・ノイズ対策は従来どおり必要です)。
2. MQTT の基本概念
2.1 ブローカ・パブリッシャ・サブスクライバ
- ブローカ(Broker): メッセージを中継するサーバー。Eclipse Mosquitto / EMQX / HiveMQ が代表的
- パブリッシャ(Publisher): メッセージを「トピック」に対して発行する側
- サブスクライバ(Subscriber): トピックを購読してメッセージを受け取る側
2.2 トピック
トピックは / 区切りの階層構造で、ワイルドカード +(1 階層)と #(任意階層)が使えます。
factory/lineA/equipmentA/temperature— 単一の値factory/lineA/+/temperature— ライン A の全設備の温度を購読factory/#— 工場全体のすべて
2.3 QoS(Quality of Service)
- QoS 0(At most once): 一度送って終わり。ロスする可能性あり
- QoS 1(At least once): 受信確認まで再送。重複の可能性あり
- QoS 2(Exactly once): 確実に 1 回だけ届く(4 段階のハンドシェイク)
製造データの計測値は QoS 1 が現実的です。QoS 2 は理論的には完璧ですが、レイテンシとブローカ負荷のコストが高くなります。
2.4 Retain と LWT
- Retain: トピックの「最新値」をブローカが保持。あとから購読したクライアントに即座に渡される。「現在の状態」を表す値に向く
- LWT(Last Will and Testament): 接続が異常切断したときにブローカが代理で発行する遺言メッセージ。「設備 A はオンライン」フラグを LWT で「オフライン」に切り替える設計が定番
3. 環境構築
3.1 ライブラリ
paho-mqtt は Eclipse Foundation がメンテナンスする公式クライアントです(Eclipse Paho 公式ページ)。
python -m pip install "paho-mqtt>=2.0"
社内プロキシで pip が通らない場合は 学習ロードマップ STEP 2 を参照してください。
3.2 ブローカ(Mosquitto)を Docker で立てる
Docker は「アプリを隔離された箱ごと動かす仕組み」で、ここではブローカを 1 コマンドで立てるために使います(Docker Desktop のインストールが必要です。本節は自宅 PC での再現を想定しています)。手順は次の 5 ステップです。
- 作業フォルダに
mosquitto\configフォルダを作る(mkdir mosquitto\config) - 下の
docker-compose.ymlを作業フォルダ直下に保存する - 下の
mosquitto.confをmosquitto\config\に保存する docker compose up -dで起動するdocker compose logsで「mosquitto version ... running」が出ていれば成功(5 章の subscriber.py を動かして確認しても OK)
# docker-compose.yml
services:
mosquitto:
image: eclipse-mosquitto:2
restart: unless-stopped # PC 再起動後も自動で立ち上がる(24 時間運用の基本)
ports:
- "1883:1883"
volumes:
- ./mosquitto/config:/mosquitto/config
- ./mosquitto/data:/mosquitto/data
- ./mosquitto/log:/mosquitto/log
# mosquitto/config/mosquitto.conf
listener 1883
allow_anonymous true
persistence true
persistence_location /mosquitto/data/
本番では匿名接続を許可せず、必ず認証を入れます(後述)。Docker がどうしても使えない場合の代替は、Mosquitto の Windows インストーラ(要管理者権限)か、実験専用の公開テストブローカ test.mosquitto.org(平文・誰でも購読可能なので業務データは絶対に流さないこと)です。
4. Publisher の最小実装
# publisher.py
import json
import time
import paho.mqtt.client as mqtt
BROKER = "localhost"
PORT = 1883
CLIENT_ID = "publisher-edge-01"
TOPIC = "factory/lineA/equipmentA/temperature"
def main() -> None:
client = mqtt.Client(
callback_api_version=mqtt.CallbackAPIVersion.VERSION2,
client_id=CLIENT_ID,
)
client.connect(BROKER, PORT, keepalive=60)
client.loop_start()
try:
for i in range(100):
payload = {
"value": 25.0 + i * 0.1,
"unit": "C",
"ts": time.time(),
}
client.publish(TOPIC, json.dumps(payload), qos=1)
time.sleep(1.0)
finally:
client.loop_stop()
client.disconnect()
if __name__ == "__main__":
main()
ポイント:
callback_api_version=...VERSION2: paho-mqtt 2.x の新 API。古いサンプルは V1 で書かれているので注意loop_start(): バックグラウンドスレッドでネットワーク I/O を回す。注意——loop_stop()は送信スレッドを止めるだけで未送信メッセージの送達は保証しません。確実に届けてから終了したい場合はclient.publish(...).wait_for_publish()で送達を待ちます(このサンプルは publish ごとに 1 秒 sleep するため実質問題になりません)- ペイロードは JSON で構造化。
valueだけ送るとあとから単位や時刻情報が必要になったときに困る
5. Subscriber の最小実装
# subscriber.py
import json
import paho.mqtt.client as mqtt
BROKER = "localhost"
PORT = 1883
CLIENT_ID = "subscriber-mes-01"
TOPIC = "factory/+/+/temperature"
def on_connect(client, userdata, flags, reason_code, properties) -> None:
print(f"connected: {reason_code}")
client.subscribe(TOPIC, qos=1)
def on_message(client, userdata, message) -> None:
try:
payload = json.loads(message.payload.decode("utf-8"))
except (UnicodeDecodeError, json.JSONDecodeError):
# 同じトピックに非 JSON が流れてきてもコールバックを死なせない
print(f"{message.topic}: (non-JSON) {message.payload!r}")
return
print(f"{message.topic}: {payload}")
def main() -> None:
client = mqtt.Client(
callback_api_version=mqtt.CallbackAPIVersion.VERSION2,
client_id=CLIENT_ID,
)
client.on_connect = on_connect
client.on_message = on_message
client.connect(BROKER, PORT, keepalive=60)
client.loop_forever()
if __name__ == "__main__":
main()
loop_forever() はメインスレッドでネットワーク I/O を回し続けます。終了は Ctrl+C もしくは client.disconnect() です。
動かし方: ターミナルを 2 つ開き、先に python subscriber.py(connected: Success と表示され待機)、次に別ターミナルで python publisher.py を実行します。subscriber 側に factory/lineA/equipmentA/temperature: {'value': 25.0, ...} が毎秒流れれば成功です。
6. トピック設計のベストプラクティス
トピック設計は MQTT 運用の成否を分けます。一度動き始めると後から変更しにくいので、最初の設計に時間をかける価値があります。
6.1 階層は「物理 → 論理 → 計測」の順
# 良い例
factory/lineA/equipmentA/temperature
factory/lineA/equipmentA/pressure
factory/lineA/equipmentA/status/running
factory/lineB/equipmentX/temperature
# 悪い例(計測値が先に来ると拡張時に困る)
temperature/factory/lineA/equipmentA
「ワイルドカードでまとめて購読する」シーンを想像して、よく一緒に購読する単位がワイルドカードで取れるように並べます。
6.2 トピックは「動詞」を含めない
get_temperature, set_threshold のような動詞を入れると、Pub/Sub のフラットなモデルが崩れます。動作の指示は別トピック(commands/...)に分け、結果は元の状態トピックの更新で表現するのが自然です。
6.3 ペイロードは JSON で構造化
{
"value": 25.3,
"unit": "C",
"ts": 1714723200,
"quality": "good",
"source": "sensor-A1"
}
ts(タイムスタンプ)と quality はあとから必ず欲しくなる情報です。最初から入れておくと拡張が楽になります。
7. Retain と LWT を使う
7.1 Retain で「現在状態」を保持
# 設備のオンライン状態を retain=True で発行
client.publish(
"factory/lineA/equipmentA/status/online",
payload="true",
qos=1,
retain=True,
)
Retain したメッセージは、後から接続したサブスクライバが購読開始した瞬間にブローカから配信されます。「ダッシュボードを開いた瞬間に最新値が見える」のはこの仕組みです。
7.2 LWT で異常切断を検知
# will_set は connect より「前」に呼ぶこと(後だと黙って LWT が効かない)
client.will_set(
"factory/lineA/equipmentA/status/online",
payload="false",
qos=1,
retain=True,
)
client.connect(BROKER, PORT, keepalive=60)
このクライアントが正常に切断せずに音信不通になると、ブローカが代理で "false" を発行します。検知タイミングの設計値は、仕様上 keepalive の 1.5 倍経過で切断判定です(keepalive=60 なら最悪 90 秒後)。「設備 A の通信が落ちた」が、消費側で検知できます。
8. TLS と認証
8.1 ユーザー認証
まず mosquitto.conf を次のように変更します(3.2 の persistence 2 行はそのまま残します——10 章の復元機能が消えてしまうため)。
listener 1883
allow_anonymous false
password_file /mosquitto/config/passwd
persistence true
persistence_location /mosquitto/data/
次にパスワードファイルを作ってから、設定反映のために再起動します(先に restart すると password_file が見つからず起動に失敗しうるので、この順序で)。
docker compose exec mosquitto mosquitto_passwd -c /mosquitto/config/passwd edge-01
docker compose restart
Python 側は、§4・§5 の両スクリプトで client.connect(...) の直前に次の 2 行を追加します(認証を有効化した時点で、追加しないと接続できなくなります)。パスワードはコードに直書きせず、環境変数から読むのが基本です。
import os
client.username_pw_set("edge-01", os.environ["MQTT_PASSWORD"])
環境変数は、PowerShell なら実行前に $env:MQTT_PASSWORD = "設定したパスワード"(cmd なら set MQTT_PASSWORD=...)で設定します。
8.2 TLS(暗号化)
# Python 側(TLS バージョンは指定せずライブラリの既定に任せるのが現行推奨)
client.tls_set(
ca_certs="ca.crt",
certfile="client.crt", # クライアント証明書が必須のサービス(AWS IoT Core 等)向け
keyfile="client.key", # サーバー認証だけなら ca_certs のみで可
)
client.connect(BROKER, 8883, keepalive=60)
なお、証明書ファイル(ca.crt / client.crt / client.key)の生成と、Mosquitto 側の 8883 リスナー設定(listener 8883 + cafile/certfile/keyfile)は本記事の範囲外です——閉域の工場 LAN 内なら認証(8.1)のみで運用を始め、TLS はクラウド接続時に導入する、という段階論が現実的です。クラウドの IoT サービス(AWS IoT Core 等)に接続する場合は相互 TLS 認証(クライアント証明書)が必須で、証明書はサービス側から発行されます。証明書の期限管理を運用設計に含める前提が必要です。
9. 実装パターン
パターン A: PLC からデータを取得して MQTT に流すゲートウェイ
Modbus / OPC UA で PLC から値を読み、MQTT にパブリッシュするブリッジ。Raspberry Pi クラスでも十分動く負荷で、エッジデバイスの定番ユースケースです。
パターン B: 複数センサーから集約 → ダッシュボード
IoT センサー(環境温湿度、振動、電力)から MQTT で集約し、Telegraf(データ収集役)→ InfluxDB(時系列データベース)→ Grafana(ブラウザのグラフ画面)のスタックで時系列ダッシュボード化。この構成は InfluxDB + Grafana の記事で詳しく扱います。
パターン C: コマンド系(Pub/Sub の双方向利用)
# 状態は state トピックに発行(retain)
factory/lineA/equipmentA/state/threshold
# 指示は commands トピックに発行
factory/lineA/equipmentA/commands/set_threshold
状態と指示を別トピックに分けると、「再起動した瞬間に最新の指示が誤って再実行される」事故を避けられます。
10. 24 時間運用のチェックリスト
- 自動再接続: paho-mqtt は
loop_start()中に切断されると自動で再接続を試みる。on_disconnectでログを残す - クライアント ID は固定: ランダム生成すると、再接続のたびに別クライアントとして見える。LWT も期待通りに動かない
- persistent session:
clean_session=False相当にすると、切断中に届いたメッセージをブローカが保持する。MQTT v5 ではclean_start=Falseに加えて Session Expiry Interval の指定が必要(既定 0 のままだと切断でセッションが即消える) - メッセージサイズ: 1 件あたり数 KB に抑える。大きい場合は別チャンネル(HTTP / S3)に置いて URL だけ MQTT で流す
- ブローカの監視: ブローカは単一障害点なので、自身の死活監視も忘れない。Mosquitto の
$SYS/#トピックでメトリクスが取れる。persistence true(3.2 の設定)にしておけば、ブローカ再起動後も retain メッセージや未配信キューがおおむね復元される(保存は一定間隔のため、突然死の直前分は失われうる) - QoS の選定: 計測値は QoS 1、ハートビート(生存確認の定期信号)や非クリティカルなテレメトリは QoS 0 で帯域節約
11. おわりに
MQTT は「軽量・非同期・疎結合」という性質から、工場 IoT のメッセージング基盤としてほぼ第一選択肢になっています。OPC UA で集めたデータを、MQTT 経由でクラウドや社内システムに渡す——という構成は、これからの製造業 IT で頻出するアーキテクチャです。
本記事のサンプルは Mosquitto + Python だけで自宅 PC で再現できます。クラウドの IoT サービスへの接続も、認証部分以外はほぼ同じコードで動きます。
⚠️ 実機接続時の注意: 本記事のサンプルは検証用(ローカル Mosquitto / ダミートピック想定)です。実際の産業機器・PLC・センサーを MQTT に載せる場合は、必ず読み出し(テレメトリ)から始めて、書き込み(commands/* トピックによる設備制御・設定変更)は安全インターロック・計画停止のうえで実施してください。とくに §9 パターン C の commands/set_threshold のような設備への指示メッセージを扱う場合は、①TLS 暗号化+クライアント証明書認証、②ブローカ側 ACL(読み書き権限分離)、③リトライ・重複配送(QoS 1 の副作用)を考慮した冪等な受信ハンドラ、④異常時のフェイルセーフ状態、⑤書き込み履歴の監査ログ、を必ず設計に含めてください。24 時間稼働機器の停止・書き込みは、自社の安全基準と保全側の承認を経ることが前提です。本記事のコードによる損害・事故(機器故障・ライン停止・情報漏洩を含む)について筆者・GenbaPy は責任を負いません。