You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Docker运行Paho MQTT示例遇问题:接收容器启动后自动退出

MQTT客户端容器启动后退出&时序问题排查

问题描述

我花了好几个小时调试这个示例,本来是做给朋友学Python用的,结果自己卡壳了。我的Python知识不多,之前试过用time.sleep延迟执行(代码已经删掉了),但程序线程还是会直接结束。

  • 预期效果:sender容器要在receiver容器之后启动,保证receiver已经订阅MQTT Broker并在等消息。
  • 实际情况:receiver容器一启动就退出了。

提前谢过各位帮忙。

Docker Compose配置

services:
  mqtt_broker:
    image: eclipse-mosquitto
    volumes:
      - "./mosquitto.conf:/mosquitto/config/mosquitto.conf"
  client_send:
    build:
      context: ./client_send/
    environment:
      BROKER_HOST: mqtt_broker
    depends_on:
      - client_receive
  client_receive:
    build:
      context: ./client_receive/
    environment:
      BROKER_HOST: mqtt_broker
    depends_on:
      - mqtt_broker

Receiver代码(原版本)

import os
import paho.mqtt.client as mqtt

def on_connect(client, userdata, flags, rc):
    print("[receiver] Connected with result code " + str(rc))
    client.subscribe("sample_topic")

def on_message(client, userdata, msg):
    print("[receiver] got a message: " + str(msg.payload.decode()))
    client.loop_stop()

client = mqtt.Client()
client.on_connect = on_connect
client.on_message = on_message

client.connect(os.environ["BROKER_HOST"], 1883, 60)
client.loop_start()

Sender代码(原版本)

import os
import paho.mqtt.client as mqtt

def run():
    print("[sender] will send a message")
    client.publish("sample_topic", "message from sender")
    client.loop_stop()

def on_connect(client, userdata, flags, rc):
    print("[sender] Connected with result code " + str(rc))
    run()

client = mqtt.Client()
client.on_connect = on_connect

client.connect(os.environ["BROKER_HOST"], 1883, 60)
client.loop_start()

解决办法

1. 修复Receiver容器退出问题

原Receiver代码中,client.loop_start()启动后台线程处理MQTT消息,但主线程执行完所有代码后直接退出,导致整个容器进程终止。需要在主线程加阻塞逻辑,直到收到消息后再退出:

import os
import time
import paho.mqtt.client as mqtt

def on_connect(client, userdata, flags, rc):
    print("[receiver] Connected with result code " + str(rc))
    client.subscribe("sample_topic")

def on_message(client, userdata, msg):
    print("[receiver] got a message: " + str(msg.payload.decode()))
    client.loop_stop()
    # 标记主线程可以退出
    global keep_running
    keep_running = False

client = mqtt.Client()
client.on_connect = on_connect
client.on_message = on_message

client.connect(os.environ["BROKER_HOST"], 1883, 60)
client.loop_start()

keep_running = True
# 主线程阻塞等待消息
while keep_running:
    time.sleep(1)

2. 解决MQTT服务就绪问题

Docker Compose的depends_on只保证容器启动顺序,不保证Broker服务完全就绪。给Sender加重连逻辑,确保连接成功后再发消息,同时提高QoS保证消息投递:

import os
import time
import paho.mqtt.client as mqtt

def run():
    print("[sender] will send a message")
    # QoS=1保证消息至少被投递一次
    client.publish("sample_topic", "message from sender", qos=1)
    client.loop_stop()

def on_connect(client, userdata, flags, rc):
    print("[sender] Connected with result code " + str(rc))
    run()

def on_disconnect(client, userdata, rc):
    if rc != 0:
        print("[sender] Unexpected disconnection, reconnecting...")
        time.sleep(1)
        client.reconnect()

client = mqtt.Client()
client.on_connect = on_connect
client.on_disconnect = on_disconnect

client.connect(os.environ["BROKER_HOST"], 1883, 60)
client.loop_start()

# 阻塞主线程直到消息发布完成
time.sleep(2)

3. 检查MQTT Broker配置

确保mosquitto.conf允许匿名访问(如果没有设置认证的话),否则客户端无法连接:

allow_anonymous true
listener 1883

内容的提问来源于stack exchange,提问作者k1r0s

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.11 04:15:18