paho-mqttの使い方|PythonでMQTTを購読する方法

「PythonからMQTTブローカーに繋いで、センサーの値やデバイスのメッセージを受け取りたい」と思って調べ始めたものの、コールバックの種類が多くて何から書けばいいのか迷った経験はないでしょうか。特にpaho-mqttはバージョン2.0でコールバックの書き方が変わっており、古い記事のコードをそのまま写すと動かないことがあります。
この記事では、paho-mqttを使ってPythonでMQTTブローカーに接続し、メッセージを受信(サブスクライブ)して切断するまでの流れを、動くサンプルつきで順番に説明します。つまずきやすいポイントや、TLS(WSS)で接続したい場合の設定も最後にまとめています。
※paho-mqttはEclipse財団が公開している、MQTTプロトコル用のクライアントライブラリです。

この記事で解決できること

  • paho-mqttのインストールから接続、切断までの基本の流れ
  • 接続、受信、サブスクライブ、アンサブスクライブの各コールバックの役割
  • 「接続はできたのにメッセージが届かない」といった典型的なトラブルの切り分け方
  • TLS(WSS)で接続するときの設定と、開発時だけ使ってよい設定の注意点

MQTTとpaho-mqttの基本

MQTTは、IoT機器などでよく使われる軽量なメッセージング用プロトコルです。「ブローカー」と呼ばれる中継サーバーを間に置き、送る側(パブリッシャー)と受け取る側(サブスクライバー)がお互いを直接知らなくてもやり取りできる点が特徴です。
メッセージには「トピック」という宛先のような名前が付いています。受け取りたい側は、欲しいトピックを指定してブローカーに「購読」を申し込みます。これがサブスクライブです。申し込んだトピックにメッセージが届くと、ブローカーが自動的に転送してくれます。
paho-mqttは、このクライアント側の処理(接続、購読、受信、切断)をPythonで書くためのライブラリです。メッセージが届いたときや接続が完了したときに、あらかじめ登録しておいた関数(コールバック関数)が呼ばれる仕組みになっています。

インストール

pipでインストールします。

pip install paho-mqtt

なお、この記事のコードはpaho-mqtt 2.x系を前提にしています。1.x系が入っている環境では、コールバックの引数の数が異なるため、そのままでは動きません。pip show paho-mqttでバージョンを確認しておくと安心です。

MQTTクライアントの使用方法

MQTTクライアントの生成

最初にクライアントのオブジェクトを作ります。

import paho.mqtt.client as mqtt

client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2)

引数のCallbackAPIVersion.VERSION2は、コールバック関数の書き方(引数の形)をどのバージョンに合わせるかを指定するものです。paho-mqtt 2.0以降ではこの指定が必須になっており、省略するとエラーになります。「昔のサンプルをコピーしたらValueErrorが出た」という場合は、ほぼここが原因です。

MQTTブローカーへの接続

ブローカーのアドレス、ポート番号、キープアライブ(秒)を指定して接続します。

client.connect('192.168.1.10', 1883, 60)

ポート番号の1883は、暗号化なしのMQTTで標準的に使われる番号です。TLSで暗号化する場合は8883が使われることが多いですが、ブローカーの設定によって変わるので、接続先の管理者や設定ファイルで確認してください。
3つ目の60はキープアライブで、「このくらいの間隔で生存確認の通信をします」とブローカーに伝える値です。この間隔を過ぎても通信がないと、ブローカー側は接続が切れたものとして扱います。

MQTTブローカーから切断

使い終わったら切断します。

client.disconnect()

メッセージループ

connectを呼んだだけでは、実際の送受信は始まりません。メッセージループを開始し、プログラムが終了しないようにします。

client.loop_forever()

loop_forever()は処理をその場で待ち続ける関数で、disconnect()が呼ばれるまで次の行には進みません。受信専用のスクリプトであればこれで十分です。ほかの処理と並行して動かしたい場合は、別スレッドでループを回すloop_start()も用意されています。

各種コールバック関数

paho-mqttでは、「接続が完了した」「メッセージが届いた」といったイベントごとに、登録した関数が呼ばれます。ここでは使用頻度の高い4つを紹介します。

接続確認応答

on_connectには接続確認応答(CONNACK)のコールバック関数を指定します。ブローカーが接続を受け付けたかどうかは、ここで受け取るreason_codeで判断します。

def on_connect(client, userdata, flags, reason_code, properties):
    # 接続確認応答(CONNACK)を受信
    print(f"Connected with result code {reason_code}")

client.on_connect = on_connect

認証エラーなどで接続が拒否された場合も、connect()自体は例外にならず、このreason_codeにエラーの内容が入ります。接続の成否を確かめたいときは、まずここの出力を見てください。

MQTTメッセージ

on_messageにはMQTTメッセージ受信のコールバック関数を指定します。購読したトピックにメッセージが届くたびに呼ばれます。

def on_message(client, userdata, msg):
    # MQTTメッセージを受信
    print(f"{msg.topic} " + str(msg.payload))

client.on_message = on_message

msg.topicが届いたトピック名、msg.payloadが本文です。payloadは文字列ではなくバイト列(bytes)で届くため、文字として扱いたい場合はmsg.payload.decode()などで変換する必要があります。

上記str(msg.payload)はb'...'が付いた文字列表現になります。

サブスクライブ確認応答

on_subscribeにはサブスクライブ確認応答(SUBACK)のコールバック関数を指定します。購読の申し込みをブローカーが受け付けたかどうかを確認できます。

def on_subscribe(client, userdata, mid, reason_code_list, properties):
    if reason_code_list[0].is_failure:
        print(f"Broker rejected you subscription: {reason_code_list[0]}")
    else:
        print(f"Broker granted the following QoS: {reason_code_list[0].value}")

client.on_subscribe = on_subscribe

購読を申し込んだだけでは、受け付けられたかどうかは分かりません。権限がないトピックを購読しようとした場合などは、ブローカーに拒否されます。メッセージが届かないときは、この関数で拒否されていないかを確認すると原因を絞り込めます。成功した場合に表示されるQoSは、メッセージ配信の確実さを表す設定値(0、1、2)です。

アンサブスクライブ確認応答

on_unsubscribeにはアンサブスクライブ確認応答(UNSUBACK)のコールバック関数を指定します。購読の解除をブローカーが受け付けたときに呼ばれます。

def on_unsubscribe(client, userdata, mid, reason_code_list, properties):
    # アンサブスクライブ確認応答(UNSUBACK)を受信
    if len(reason_code_list) == 0 or not reason_code_list[0].is_failure:
        print("unsubscribe succeeded (if SUBACK is received in MQTTv3 it success)")
    else:
        print(f"Broker replied with failure: {reason_code_list[0]}")

client.on_unsubscribe = on_unsubscribe

MQTTv3では、アンサブスクライブの応答に結果コードが含まれません。そのためreason_code_listが空になることがあり、コード側でlen(reason_code_list) == 0を成功として扱っています。

サブスクライブの実装サンプル

ここまでの内容を組み合わせた、受信用のサンプルです。「接続する、購読を申し込む、メッセージを1件受け取る、購読を解除する、切断する」という一連の流れを確認できます。

import paho.mqtt.client as mqtt

BROKER = "mqtt.sample.io"
PORT = 1883
KEEP_ALIVE = 60
MQTT_TOPIC="/mytopic"

def on_connect(client, userdata, flags, reason_code, properties):
    # 接続確認応答(CONNACK)を受信
    print(f"Connected with result code {reason_code}")
    # サブスクライブ要求
    client.subscribe(MQTT_TOPIC)
#;

def on_subscribe(client, userdata, mid, reason_code_list, properties):
    # サブスクライブ確認応答(SUBACK)を受信
    if reason_code_list[0].is_failure:
        print(f"Broker rejected you subscription: {reason_code_list[0]}")
    else:
        print(f"Broker granted the following QoS: {reason_code_list[0].value}")

def on_unsubscribe(client, userdata, mid, reason_code_list, properties):
    # アンサブスクライブ確認応答(UNSUBACK)を受信
    if len(reason_code_list) == 0 or not reason_code_list[0].is_failure:
        print("unsubscribe succeeded (if SUBACK is received in MQTTv3 it success)")
    else:
        print(f"Broker replied with failure: {reason_code_list[0]}")
    # ブローカーから切断
    client.disconnect()

def on_message(client, userdata, msg):
    # MQTTメッセージを受信
    print(f"{msg.topic} " + str(msg.payload))
    # アンサブスクライブ要求
    client.unsubscribe(MQTT_TOPIC)

client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2)
client.on_connect = on_connect
client.on_message = on_message
client.on_subscribe = on_subscribe
client.on_unsubscribe = on_unsubscribe

# ブローカーに接続
client.connect(BROKER, PORT, KEEP_ALIVE)

client.loop_forever()

処理の流れ

このサンプルは、次の順番で動きます。

  1. connect()でブローカーに接続し、loop_forever()でメッセージループを開始
  2. 接続が完了するとon_connectが呼ばれ、その中でsubscribe()を実行
  3. ブローカーによる購読の受け付けと、on_subscribeの呼び出し
  4. 対象トピックにメッセージが届くとon_messageが呼ばれ、その中でunsubscribe()を実行
  5. ブローカーが解除を受け付けるとon_unsubscribeが呼ばれ、その中でdisconnect()を実行
  6. 切断されるとループが終わり、プログラムが終了

購読の申し込みをconnect()の直後ではなくon_connectの中に書いているのがポイントです。接続が完了する前にsubscribe()を呼ぶと取りこぼす可能性がありますし、通信が途切れて自動再接続したときにもon_connectが呼ばれるため、購読が自動的に張り直されます。
また、サンプルではメッセージを1件受け取った時点で購読を解除して終了しています。これは一連のコールバックを順に確認するための動きです。常時受信し続けたい場合は、on_message内のunsubscribe()を外してください。

TLS(WSS)で接続する場合の設定

ブローカーがTLSによる暗号化や、WebSocket over TLS(WSS)での接続を求めている場合は、クライアント側にも設定が必要です。ここでは、接続方式をprotocol変数で切り替える想定で、wssのときだけTLSを設定する例を紹介します。まず、標準ライブラリのsslモジュールをインポートし、接続方式を決めておきます。

import ssl

protocol = "wss"

そのうえで、接続前に次のような設定を入れます。

    # TLSが必要なら設定
    if protocol == "wss":
        #client.tls_set(ca_certs="/home/app/src/192.0.2.100.crt")  # 必要に応じて証明書ファイルパスを指定
        client.tls_set(cert_reqs=ssl.CERT_NONE)
        client.tls_insecure_set(True)  # 証明書検証無効化(開発時のみ推奨)

それぞれの意味は次のとおりです。

  • client.tls_set(ca_certs=...): サーバー証明書を検証するための証明書ファイルを指定する設定
  • client.tls_set(cert_reqs=ssl.CERT_NONE): サーバー証明書の検証を行わない設定
  • client.tls_insecure_set(True): 証明書のホスト名の検証を行わない設定

WebSocketで接続する場合は、クライアントの生成時にtransport="websockets"を指定する必要がある点にも注意してください。

本番環境では証明書の検証を有効にする

cert_reqs=ssl.CERT_NONEとtls_insecure_set(True)は、自己署名証明書を使った開発用のブローカーなど、検証用の環境で動作確認を手早く進めるための設定です。通信は暗号化されますが、接続先が本物のサーバーかどうかを確認しないため、なりすましや中間者攻撃を防げません。
本番環境では、この2行を使わないでください。コメントアウトされているca_certsにブローカーの証明書(またはそれを発行した認証局の証明書)のパスを指定し、証明書の検証を有効にした状態で運用するのが基本です。開発時に検証を無効にしていたコードが、そのまま本番に混ざってしまう事故は意外とよくあります。リリース前には、検証を無効にする設定が残っていないかを必ず確認してください。

うまく動かないときの確認ポイント

実際に動かしてみて困りやすい症状と、確認する場所をまとめます。

  • ValueErrorが出てクライアントを生成できない:mqtt.Client()へのCallbackAPIVersion.VERSION2の指定の有無
  • 接続できない、タイムアウトする: ブローカーのアドレスとポート番号、ファイアウォールの設定
  • on_connectのreason_codeが失敗を示している: ユーザー名やパスワードなどの認証設定
  • 接続はできるのにメッセージが届かない:on_subscribeでの拒否の有無、トピック名の綴りや先頭のスラッシュの有無
  • コールバックの引数の数が合わないというエラーが出る: paho-mqttのバージョンと、コールバックの書き方の整合性
  • 受信した文字がb'...'のように表示される:payloadがバイト列のため、decode()による文字列への変換

トピック名は大文字と小文字、先頭のスラッシュの有無も区別されます。/mytopicとmytopicは別のトピックとして扱われるため、受信できないときは真っ先に疑ってみてください。

まとめ

paho-mqttでメッセージを受信するには、クライアントを生成し、コールバック関数を登録して、ブローカーに接続したうえでloop_forever()で待ち続ける、という流れになります。購読の申し込みはon_connectの中で行うと、再接続のときにも自動で購読し直されるので安全です。
動かないときは、on_connectとon_subscribeの出力を見れば、接続の問題なのか購読の問題なのかを切り分けられます。TLSの設定では、検証を無効にする設定を開発時だけにとどめることを忘れないようにしてください。

このエントリーをはてなブックマークに追加
にほんブログ村 IT技術ブログへ

コメント

メールアドレスが公開されることはありません。 ※ が付いている欄は必須項目です