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

MQTT订阅者无法接收同一主题消息问题排查求助

MQTT订阅端无法接收消息,持续打印连接成功的问题排查

我用Python的paho-mqtt库开发了一个「中转航班时刻表变更」程序,逻辑如下:

  • 发布端:输入航班号(作为Topic)、目的地国家和新航班时间,从指定列表随机选一个和目的地不同的中转地点,然后把航班信息发布到对应Topic。
  • 订阅端:输入相同航班号作为Topic,尝试接收发布端消息,但运行时一直打印Connected successfully,触发不了on_message回调,也收不到消息。

发布端代码

import paho.mqtt.client as mqtt
import time
from datetime import datetime, date
import random

def on_connect(client, userdata, flags, rc):
    if (rc==0):
        global connected
        connected = True
        #print("Successfully Connected.")
        client.on_publish = on_publish
    else:
        print("Failed to connect.")
        
def on_publish(client, userdata, mid):
    print("Published successfully. MID: "+str(mid))
    

listTransit = ["Singapura", "Qatar", "Korea Selatan", "Turki", "Republik Tiongkok",
                "Amerika Serikat", "Jepang", "Uni Emirat Arab", "Oman", "Islandia"]
                
broker_address="broker.emqx.io"

client = mqtt.Client("Publisher")
client.on_connect = on_connect
client.connect(broker_address, port=1883)

client.loop_start()

topic = input("Masukkan nomor penerbangan: ")
negaraTujuan = input("negara tujuan: ")
print("Masukkan waktu penerbangan baru (Format: [jam::menit::detik])")
str_time = input()

Date = date.today()
Time = datetime.strptime(str_time, '%H::%M::%S').time()
Location = random.randrange(0,len(listTransit))

if (listTransit[Location] != negaraTujuan):
    message = Date.strftime("%Y/%m/%d")+"\nTujuan: "+negaraTujuan+"\nLokasi Transit  : "+listTransit[Location]+"\nJam terbang    : "+Time.strftime("%H:%M:%S")
    client.publish(topic, message)
    print("Topic: ",topic)
    print(message)
else:
    while listTransit[Location] == negaraTujuan:
        Location = random.randrange(0,len(listTransit))
    message = Date.strftime("%Y/%m/%d")+"\nTujuan: "+negaraTujuan+"\nLokasi Transit  : "+listTransit[Location]+"\nJam terbang    : "+Time.strftime("%H:%M:%S")
    client.publish(topic, message)
    print(message)
    
client.loop_stop()

订阅端代码

import paho.mqtt.client as mqtt
import time
from datetime import datetime, datetime
import re

def on_connect(client, userdata, flags, rc):
    if (rc == 0):
        print("Connected successfully.")
        #global topic
        #topic = input("Masukkan nomor penerbangan anda: ")
        #client.subscribe(topic)
    else:
        print("Connection failed.")
        
def on_message(client, userdata, msg):
    print("on_message callback function activated.")
    sched = str(msg.payload.decode("utf-8"))
    print(sched)
    
def on_subscribe(client, userdata, mid, granted_qos):
    print("Subscribed to "+topic+" successfully")
    
    
broker_address="broker.emqx.io"
topic = input("Masukkan nomor penerbangan anda: ")
negaraTujuan = input("negara tujuan: ")

client = mqtt.Client("Subscriber")
client.subscribe(topic)
client.on_connect = on_connect
client.on_message = on_message
client.on_subscribe = on_subscribe
client.connect(broker_address, port=1883)

client.loop_forever()

运行输出

发布端输出

Masukkan nomor penerbangan: YT05TA
negara tujuan: Australia
Masukkan waktu penerbangan baru (Format: [jam::menit::detik])
12::50::00
Topic:  YT05TA
2023/01/03
Tujuan: Australia
Lokasi Transit  : Amerika Serikat
Jam terbang    : 12:50:00
Published successfully. MID: 1

订阅端输出

Masukkan nomor penerbangan anda: YT05TA
negara tujuan: Australia
Connected successfully.
Connected successfully.
Connected successfully.
Connected successfully.
问题排查与解决

核心问题分析

  1. 订阅时机错误:订阅端在调用client.connect()之前就执行了client.subscribe(topic),此时客户端未与Broker建立有效连接,订阅请求无法被Broker处理,导致后续收不到消息。
  2. 客户端ID冲突风险:固定的客户端ID("Subscriber")如果被其他重复连接的客户端使用,会触发Broker的踢线机制,导致客户端反复断开重连,出现持续打印Connected successfully的情况。
  3. 冗余代码:订阅端中from datetime import datetime, datetime属于重复导入,虽不影响功能,但需修正。

修正后的订阅端代码

import paho.mqtt.client as mqtt
import time
from datetime import datetime
import re
import random

def on_connect(client, userdata, flags, rc):
    if (rc == 0):
        print("Connected successfully.")
        # 连接成功后执行订阅操作
        client.subscribe(userdata["topic"])
    else:
        print("Connection failed.")
        
def on_message(client, userdata, msg):
    print("on_message callback function activated.")
    sched = str(msg.payload.decode("utf-8"))
    print(sched)
    
def on_subscribe(client, userdata, mid, granted_qos):
    print(f"Subscribed to {userdata['topic']} successfully")
    
    
broker_address="broker.emqx.io"
topic = input("Masukkan nomor penerbangan anda: ")
negaraTujuan = input("negara tujuan: ")

# 使用随机后缀生成唯一客户端ID,避免冲突
client_id = f"Subscriber_{random.randint(1000, 9999)}"
client = mqtt.Client(client_id)
# 将topic存入userdata,在回调中使用
client.user_data_set({"topic": topic})
client.on_connect = on_connect
client.on_message = on_message
client.on_subscribe = on_subscribe
client.connect(broker_address, port=1883)

client.loop_forever()

修正说明

  • 将订阅操作移至on_connect回调中,确保只有连接成功后才发送订阅请求。
  • 生成带随机后缀的客户端ID,避免重复连接时的ID冲突问题。
  • 使用user_data_set传递topic参数,避免依赖全局变量,代码更健壮。
  • 修正了datetime重复导入的问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 10:10:34