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 ステップです。

  1. 作業フォルダに mosquitto\config フォルダを作る(mkdir mosquitto\config
  2. 下の docker-compose.yml を作業フォルダ直下に保存する
  3. 下の mosquitto.confmosquitto\config\ に保存する
  4. docker compose up -d で起動する
  5. 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.pyconnected: 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 は責任を負いません。

関連記事

参考文献・一次情報