paho-mqtt は 1.x系のサンプルをそのまま貼ると例外で止まります。2026年8月23日時点のPyPI最新は 2.1.0(2024年4月29日公開)で、Client() の第1引数が client_id から callback_api_version へ入れ替わったためです。この記事では、その移行で踏む地雷、4つのループ方式の選び分け、reconnect_delay_set の倍加待機と再購読の書き方、inflight 20件と送信キューの初期値、aiomqtt へ移すかの線引きを公式ソースに沿って整理しました。
まとめ|Python MQTT実装で先に決まるループ方式と再接続の設計
Pythonでクライアントを書くとき、最初に決まるのは2点です。ネットワーク処理をどのループ方式で回すか、切断からの復帰をどこまで自前で書くか。後から差し替えづらいのはこの2点です。
ループ方式は用途で決まります。収集専用の常駐プロセスなら loop_forever、Webアプリやバッチと同居させるなら loop_start。loop() 単体は公式docstringが「もはや推奨しない」と書いています。厄介なのは再接続で、paho-mqtt は接続の復帰こそ自動ですが購読の復元はやりません。subscribe() を on_connect に置く一手を怠ると、切断後は何も届かないプロセスが残ります。
QoS 1以上を使うなら、inflight 20件と送信キュー無制限という初期値の効き方まで見ておいてください。asyncio へ移すかは最後の判断でよく、収集が主目的なら loop_start で足ります。以下、移行・ループ・再接続・送達・asyncio・採用判断の順に根拠を示します。
paho-mqtt 2系の最小実装|1.x系のサンプルが例外で止まる三つの変更
日本語の解説は、2021年10月公開の 1.6.1 世代で書かれたものが今も多く残ります。2.0.0 が出たのは2024年2月10日。ここで入った破壊的変更を知らずに貼ると、実行前に例外で落ちるか、コールバックが引数の数違いで呼べません。
Client()の第1引数の入れ替わりと、文字列を渡すとValueErrorになる条件
1.x系では mqtt.Client("collector-01") と書けば、文字列はクライアントIDとして扱われました。2.0.0 で第1引数は callback_api_version になり、既定値のない必須引数に変わります。2.1.0 では VERSION1 という既定値が付いて省略は通りますが、文字列を第1引数へ渡すと「Unsupported callback API version」で始まる ValueError が飛びます。
# 1.x系の書き方(2.x系ではValueError)
client = mqtt.Client("collector-01")
# 2.x系
client = mqtt.Client(
mqtt.CallbackAPIVersion.VERSION2,
client_id="collector-01",
protocol=mqtt.MQTTv5,
)
VERSION1 を明示すれば旧来のコールバックのまま動きますが、DeprecationWarning が出ます。移行途中で一時的に置くのは構いません。新規に書くなら VERSION2 一択です。
VERSION2で並びが変わったon_connect・on_disconnectの引数
並びが変わるのは4つです。on_connect は (client, userdata, connect_flags, reason_code, properties)、on_disconnect は connect_flags が disconnect_flags に替わった同じ形。on_publish と on_subscribe は mid の次に reason_code か reason_code_list が入ります。on_message だけ全版共通で、受信処理の書き換えは不要です。
実務で効くのは、v3.1.1 でも reason_code が ReasonCode オブジェクトで渡る点。1.x系で if rc == 0: と整数比較していた箇所は if reason_code.is_failure: へ置き換えます。2.0.0 では loop() 系の max_packets 引数、loop_stop() の force 引数、message_retry_set() も削除されました。呼ぶ古いコードは TypeError か AttributeError で止まります。
publishとsubscribeを1ファイルで動かす最小コードと設定の呼び出し順
最小構成では、接続前に呼ぶものと接続後でよいものの区別が要ります。username_pw_set()・will_set()・tls_set()・reconnect_delay_set() は connect() より前。subscribe() は接続確立後、つまり on_connect の中です。
import paho.mqtt.client as mqtt
def on_connect(client, userdata, connect_flags, reason_code, properties):
if reason_code.is_failure:
print(f"connect failed: {reason_code}")
return
client.subscribe("factory/line1/+/temperature", qos=1)
def on_message(client, userdata, msg):
print(msg.topic, msg.payload.decode())
client = mqtt.Client(
mqtt.CallbackAPIVersion.VERSION2,
client_id="collector-01",
protocol=mqtt.MQTTv5,
)
client.on_connect = on_connect
client.on_message = on_message
client.username_pw_set("collector", "********")
client.connect("broker.example.jp", 1883, keepalive=60)
client.loop_forever()
keepalive の既定は60秒、connect() のタイムアウトは 5.0 秒。発行側は client.publish(topic, payload, qos=1) の1行で足ります。CONNECT から CONNACK までの往復やセッション状態はMQTTの仕組みとQoS・MQTT 5.0の実装解説を土台にしてください。1883番のまま本番へ出さないTLSと認証の設計も同記事の担当です。8883番のTLSと証明書認証、トピック単位の認可までを実装の粒度で決める段になったら、MQTTセキュリティの三層設計とブローカー別の権限粒度を扱った記事も併せて確認してください。
四つのループ方式の使い分け|再接続を内部に任せる境界とスレッド設計
paho-mqtt のネットワーク処理は呼び出し側が回します。回し方は4通りで、再接続の有無とスレッドの持ち方が違う。取り違えると、切断後に沈黙するプロセスができあがります。
loop_forever|収集専用プロセスで再接続まで内部に任せる標準の選び方
収集だけを行う常駐プロセスなら loop_forever() です。呼ぶと制御が返らない無限ループに入り、reconnect_on_failure が True(既定)である限り、切断の検知から再接続までを内部で面倒みます。disconnect() を呼べば戻ります。
初回接続だけは扱いが別で、connect() が失敗した時点で OSError が上がります。ブローカーより先にコレクタが起動するコンテナ構成では、これで即死する。loop_forever(retry_first_connection=True) なら初回も再試行の対象に入ります。起動順を制御できない環境では、この引数の有無が明暗を分けます。
loop_start|Webアプリと同居させるときのデーモンスレッドと終了処理
Django や FastAPI のプロセスから publish したい場合は loop_start() です。中身は loop_forever(retry_first_connection=True) を daemon スレッドで起動するだけで、再接続の挙動は同じ。制御はすぐ戻ります。
落とし穴は終了処理です。loop_stop() はスレッドの終了を待って join しますが、docstring が明記するとおりpublish パケットの送信までは保証しません。終了前に info.wait_for_publish(timeout=5) で送達を待ってから止めてください。daemon スレッドのため、loop_stop() を呼ばずにプロセスが終わると送信中のパケットは黙って捨てられます。
loop()とloop_read系|自前イベントループへ組み込むときの適用条件
loop() を自分で定期的に呼ぶ書き方は 1.x系のサンプルに多く残っています。2.1.0 の docstring は「loop() の単体使用はもはや推奨しない」と明言し、loop_start() か loop_forever()、外部イベントループを使うよう指示します。理由は単純で、loop() には再接続が入らないためです。
selectors や Tornado など既存のイベントループへ組み込む場合だけは、loop_read()・loop_write()・loop_misc() が正規ルートとして残ります。socket の読み書き可能通知を on_socket_register_write などで受け、キープアライブ送出のため loop_misc() を定期実行する形です。載せる必然性がないなら選びません。
四方式の比較|再接続・スレッド・向く構成を一枚の表で並べた判断材料
| 方式 | 再接続 | スレッド | 向く構成 |
|---|---|---|---|
| loop_forever() | 内部で実施 | 呼び出し元を占有 | 収集専用の常駐プロセス |
| loop_start() | 内部で実施 | daemonスレッド1本 | Web・バッチとの同居 |
| loop() | 自前で実装 | 呼び出し元 | 非推奨・移行対象 |
| loop_read系 | 自前で実装 | 既存ループに従う | 外部イベントループ組込 |
実務では上2つだけ押さえれば足ります。loop() が残る既存コードは、再接続が入っていないと考えて差し支えありません。
切断に耐える実装|指数バックオフ・再購読・LWTとQoS1の取りこぼし対策
工場や屋外の回線は落ちます。落ちた後にどう戻るかが品質差になる部分。paho-mqtt が自動でやる範囲と、書き手が埋める範囲の境界を先に押さえます。
reconnect_delay_setの1秒から120秒への倍加と、回数上限が無い挙動
再接続の待機は reconnect_delay_set(min_delay, max_delay) で決まり、既定は min_delay=1・max_delay=120 です。内部の _reconnect_wait() は待機を毎回2倍にし、max_delay で頭打ちにします。1秒、2秒、4秒と伸び、7回目以降は120秒間隔になる計算です。
ソースを読むと、再試行回数の上限はどこにも実装されていません。ブローカーが停止したままでもプロセスは生き続け、2分ごとに接続を試みる。常駐コレクタには都合のよい挙動ですが、短命なバッチでは終わらないジョブになります。経過時間による打ち切りを入れてください。多数の機器が同時に復帰する構成なら、min_delay に乱数を足して散らします。
再購読はon_connectに書く|pahoが購読を復元しない前提の実装手順
ここが最も事故になりやすい箇所です。paho-mqtt のソースに、再接続時へ購読を張り直す処理はありません。subscribe() を直線上に書くと、1回目は届き、復帰後は何も届かない状態になります。ログにエラーが出ないため発見も遅れます。
def on_connect(client, userdata, connect_flags, reason_code, properties):
if reason_code.is_failure:
return
# 再接続のたびに呼ばれるため、ここに購読をまとめる
client.subscribe([
("factory/line1/+/temperature", 1),
("factory/line1/+/status", 1),
])
client.reconnect_delay_set(min_delay=1, max_delay=30)
セッションを残す設定(v3.1.1 なら clean_session=False、v5 なら Session Expiry Interval を非ゼロ)でも購読は保持されますが、on_connect に置くほうが安全です。破棄されたかは connect_flags の session_present で判定できます。なお clean_session=False では client_id が必須で、空だと ValueError です。
will_setとClean Session|オフライン中のメッセージを残す組み合わせ
機器が落ちたことを他の購読者へ知らせるのが Will(LWT)。client.will_set(topic, payload, qos=1, retain=True) を connect() より前に呼びます。retain を True にすると、後から購読した側も最後の状態を受け取れる。復帰時に同じトピックを publish で上書きする対にして機能します。
受信の取りこぼしを防ぐ側はセッションの設定です。clean_session=False と固定 client_id なら、オフライン中の QoS 1以上のメッセージをブローカーが保持し、復帰時にまとめて配送します。保持は上限までで、長時間の停止では溢れる。製品ごとの差と構築手順はMQTTブローカーの仕組みと製品比較にまとめています。
inflight 20件と送信キューの初期値|publishの戻り値で取りこぼしを検知
publish() は MQTTMessageInfo を返し、info.rc と info.mid を持ちます。(rc, mid) = client.publish(...) というタプル展開も互換で残っています。未接続のまま呼ぶと rc は MQTT_ERR_NO_CONN。戻り値を捨てるコードは、送れていない事実を取りこぼします。
QoS 1以上では2つの初期値が効きます。_max_inflight_messages は 20 で、確認応答待ちのまま並行できる上限。_max_queued_messages は 0、つまり無制限です。細い回線へ高頻度で publish すると、20件の枠が埋まった分がキューへ積み上がりメモリを食い続ける。max_queued_messages_set(1000) と上限を置けば、溢れた publish の rc に MQTT_ERR_QUEUE_SIZE が入ります。QoSレベルごとの再送のふるまいはMQTT QoSのレベル0・1・2の違いと選び方を参照してください。
asyncio採用の判断軸|aiomqtt 2.5系とgmqttの保守状況と移行条件
既存の asyncio アプリへ組み込みたい場合の選択肢は2つ。paho-mqtt をラップする aiomqtt と、独自実装の gmqtt です。
aiomqtt 2.5.1の再接続は自前のwhileループ|3.0.0a1で変わる前提
PyPI上の aiomqtt の最新配布は 2.5.1(2026年3月5日)。この2系で押さえるべきは、paho-mqtt と違って再接続が組み込まれていない点です。公式ドキュメントが示す手順は、MqttError を捕まえて自分で待って入り直す形になります。
import asyncio
import aiomqtt
async def main():
client = aiomqtt.Client("broker.example.jp")
interval = 5
while True:
try:
async with client:
await client.subscribe("factory/line1/#")
async for message in client.messages:
print(message.payload)
except aiomqtt.MqttError:
await asyncio.sleep(interval)
asyncio.run(main())
Client のコンテキストは再利用できますが再入はできません。同じ client を並行タスクから同時に async with へ入れると壊れます。接続の共有は Client インスタンスを引数で配る形に。asyncio 側の前提が曖昧なら非同期処理と同期処理の違いと実装方式を先に確認してください。
開発中の 3.0.0 は前提が変わります。paho-mqtt 依存を外して mqtt5 という別ライブラリへ置き換え、MQTTv5 のみ対応。client.messages がメソッド呼び出しへ変わり、payload は bytes 必須、reconnect=True で自動再接続が入ります。2026年4月2日公開の 3.0.0a1 はアルファのため本番で採るのは早い。
三つのライブラリのリリース間隔と、未リリース修正が溜まる状況の実測
選定では保守のテンポも見ます。3つを並べます。
| ライブラリ | 最新配布 | 公開日 | 前版からの間隔 |
|---|---|---|---|
| paho-mqtt | 2.1.0 | 2024-04-29 | 2.0.0から2か月半 |
| aiomqtt | 2.5.1 | 2026-03-05 | 2.5.0から2か月 |
| gmqtt | 0.7.0 | 2024-11-22 | 0.6.9から約4年 |
paho-mqtt は 2.1.0 を最後にリリースが止まっています。ただしリポジトリが死んでいるわけではなく、master への直近コミットは2026年8月6日。MQTT v5 の Server Keep Alive を CONNACK から反映する修正などがマージ済で、PyPI配布版には入っていません。バグ報告の前に master を確認する手間が要るという話で、乗り換え理由にはなりません。gmqtt は 0.6.9 から 0.7.0 まで約4年空いており、新規採用は避ける判断になります。
asyncioへ移す条件と、paho-mqttのloop_startで足りる場面の線引き
判断は言い切ります。asyncio へ移すのは、すでにアプリ本体が asyncio で動いている場合だけ。FastAPI のバックグラウンドタスクから publish する、受信したメッセージを非同期のHTTPクライアントで転送するといった構成では、橋渡しコードが消える分だけ aiomqtt が有利になります。
逆に、同期のコードベースへ「並行性が上がりそうだから」という理由で持ち込むのは割に合いません。MQTTの受信は I/O 待ちが支配的で、loop_start() のスレッド1本で数千件毎秒は捌けます。ボトルネックはたいてい受信後のDB書き込みや変換側で、ループ方式を変えても解決しない。重い処理だけ queue.Queue で逃がすほうが見合います。
Pythonで書いてよい三条件と、MQTT実装を見送るべき三つの場面
ここまでを踏まえ、そもそもPythonで書くべきかを判断します。言語選定を後から覆すのは高くつきます。
Pythonでクライアントを書いてよい三条件|台数・処理内容・運用体制
第1に、1プロセスが扱う接続が数十本までで、受信を秒間数千件以下で処理する規模。第2に、受信後の処理がデータ変換・DB書き込み・機械学習の推論のいずれかであること。Pythonのライブラリ資産が効くのはこの部分で、C++やGoで書き直す動機が消えます。第3に、運用側にPythonを読める人員がいること。現地で障害対応する構成では復旧時間に直結します。
3条件のうち第2が最も効きます。単に転送するだけならPythonである必然性は薄い。AWS IoT Coreの機能・料金・接続方法で扱うルール転送で完結するなら、クライアントを書くこと自体を省けます。
見送るべき三つの場面|ブローカー代替・ミリ秒制御・GILで詰まる変換
見送る場面も条件付きで示します。第1に、Pythonでブローカーを代替しようとする構成。数万接続を捌く用途にクライアント実装を並べるのは設計として誤りで、Mosquitto や EMQX を立てるかマネージドを使ってください。
第2に、ミリ秒単位の応答保証が要る制御ループ。GC の停止時間とGILの切り替えが読めず、上限を約束できません。制御は PLC 側やC言語のタスクへ置きます。第3に、受信のたび重いCPU処理が入る構成。on_message で画像処理や大きなJSONの変換を回すと、その間ネットワークループが止まり、切断されます。キューへ積むだけにして別プロセスへ分けてください。産業用途でデータの並べ方まで規約化するなら、Sparkplug BによるMQTT上の産業データ規約を先に検討すると手戻りが減ります。
MQTT収集基盤の開発を外部へ委ねるとき見積書で確かめる四項目
外部へ出す場合、見積書で確認すべきは4点です。第1に、再接続と再購読の実装が含まれているか。第2に、QoS レベルと送信キューの上限が明記されているか。第3に、オフライン中のメッセージ保持を何時間分見込むか。第4に、台数が増えたときの同時再接続をどう散らすか。
この4点を書けるベンダーには、切断を前提にした実装経験があります。「MQTTで接続します」としか書かれていない見積は、正常系だけのコードになる公算が高い。設備データの収集から可視化・推論までを一体で相談したい場合は、AI/IoTソリューションで要件整理から対応しています。
よくある質問
移行と運用でつまずきやすい5点に答えます。
paho-mqttでClient()を作るとDeprecationWarningが出るのはなぜですか?
callback_api_version を省略したか、VERSION1 を明示したためです。2.1.0 では第1引数の既定値が VERSION1 で、この値だと「Callback API version 1 is deprecated」という警告を出します。消すには VERSION2 を渡し、4つのコールバックの引数を新しい並びへ書き換えてください。
loop_start()とloop_forever()はどちらを使うべきですか?
そのプロセスがMQTT以外の処理も抱えるかで決まります。収集だけの常駐プロセスなら loop_forever()、Webアプリやバッチと同居させるなら loop_start()。後者の中身は前者を daemon スレッドで回すもので、再接続の挙動は同じです。loop_start() なら終了時に loop_stop() を呼び、その前に wait_for_publish() で送達を確かめてください。
再接続したあとに購読が復活しないのはなぜですか?
paho-mqtt が購読を張り直さないためです。ソースに再購読の実装はなく、subscribe() を直線上に一度だけ書くと、再接続後は購読ゼロのまま接続だけが維持されます。subscribe() を on_connect の中へ移してください。セッションを残す設定でも復元されますが、on_connect に置くほうが確実です。
publish()したのにメッセージが届かないことがあるのはなぜですか?
疑うのは3点です。未接続のまま呼んだ場合、戻り値の rc に MQTT_ERR_NO_CONN が入ります。送信キューが溢れた場合は MQTT_ERR_QUEUE_SIZE。3点目は loop_stop() やプロセス終了で送信前に打ち切られた場合で、loop_stop() は送達を保証しません。戻り値を捨てていると気づけないため、rc の確認を入れてください。
paho-mqttとaiomqttのどちらを選べばよいですか?
アプリ本体が asyncio で動いているなら aiomqtt、そうでなければ paho-mqtt です。aiomqtt 2.5.1 は内部で paho-mqtt を使うラッパーで、再接続は自前の while ループで書く必要があります。開発中の 3.0.0 では paho-mqtt 依存が外れて MQTTv5 専用になるため、2系で書いたコードは移行対象になります。
関連記事
- MQTTとは?Pub/Subの仕組み・QoS・MQTT 5.0の実装からHTTPとの使い分けまで実装者向けに解説:プロトコル本体の仕組みとTLS・認証の設計
- MQTTブローカーとは?仕組み・比較・構築手順を解説:接続先となるブローカー製品の選定と構築
- MQTT QoSとは?レベル0・1・2の違い・仕組み・選び方を解説:QoSレベルごとの再送のふるまい
- 非同期処理とは?同期処理との違いから実装方式まで実装者目線で解説:asyncioへ移す前提となる考え方