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

Python IoT脚本:如何在MQTT发布前用线程执行PID生成逻辑?

问题

我是Python新手,为IoT应用编写了这套MQTT发布脚本。我希望在while循环中使用线程,具体是在等待时间(9.980秒)之前,将代码中从PID = random.randint开始的34-38行代码交由线程执行。请问这是否可行?

发布端代码

import paho.mqtt.client as mqtt
import time
import random
import thread
brokerAddress = "foobar.s2.eu.hivemq.cloud"
userName = "username"
passWord = "password"
topic = "Sensor"
#客户端收到代理的connack响应时的回调函数
def on_connect(client, userdata, flags, rc):
    
    if rc == 0:
        print("Connection Established")
    else:
        print("Coneect returened result code: " + str(rc))
def on_message(client, userdata, msg):
        print("Received message: " +msg.topic + " -> " + msg.payload.decode("utf-8"))
client = mqtt.Client()
client.on_connect = on_connect
client.on_message = on_message
client.tls_set(tls_version=mqtt.ssl.PROTOCOL_TLS)
client.username_pw_set(userName, passWord)
print("Connecting to broker::",brokerAddress)
client.connect(brokerAddress, 8883)
client.loop_start()
time.sleep(2)

while True:
    try:
        input("press enter to put the box on line:")
    except SyntaxError:
        pass
    PID = random.randint(0 , 2000000)
    print(PID)
    print("wait for the box to arrive to sensor:9.980")
    time.sleep(9.980)
    client.publish(topic, PID)
    print("..............................................................")

client.loop_stop()

发布端输出

output first (enter):
Connecting to broker:: d615ed48421a497f892c98c1b04cfb82.s2.eu.hivemq.cloud
Connection Established
press enter to put the box on line:
168488

订阅端代码

import paho.mqtt.client as mqtt
import datetime
import time

brokerAddress = "foobar.s2.eu.hivemq.cloud"
userName = "user"
passWord = "password"
topic = "Sensor"
#客户端收到代理的connack响应时的回调函数
def on_connect(client, userdata, flags, rc ):
    if rc == 0:
        print("Connection Established")
    else:
        print("Connect retured result code: " +str(rc))
        
#收到代理发布消息时的回调函数
def on_message(client, userdata, msg):
        print("Received message: " +msg.topic + " -> " + msg.payload.decode("utf-8"))
        #等待执行器
        time.sleep(.02)
        Prduction_time = datetime.datetime.now()
        print("production Time",Prduction_time)    
#创建客户端
client = mqtt.Client()
client.on_connect = on_connect
client.on_message = on_message

client.tls_set(tls_version=mqtt.ssl.PROTOCOL_TLS)
client.username_pw_set(userName, passWord)
print("Connecting to broker::",brokerAddress)
client.connect(brokerAddress, 8883)

client.subscribe(topic)

client.loop_forever()
解答

完全可行,这种方式能让主线程在等待盒子到达传感器的9.980秒期间,继续响应新的输入操作,不会被阻塞。

修改后的发布端代码

推荐使用Python标准库的threading模块(比旧的thread模块更安全易用),将需要异步执行的逻辑封装成独立函数,每次按下回车后启动一个线程执行:

import paho.mqtt.client as mqtt
import time
import random
import threading  # 替换旧的thread模块
brokerAddress = "foobar.s2.eu.hivemq.cloud"
userName = "username"
passWord = "password"
topic = "Sensor"

# 客户端收到代理的connack响应时的回调函数
def on_connect(client, userdata, flags, rc):
    if rc == 0:
        print("Connection Established")
    else:
        print("Connect returned result code: " + str(rc))

def on_message(client, userdata, msg):
    print("Received message: " + msg.topic + " -> " + msg.payload.decode("utf-8"))

# 封装线程要执行的任务
def handle_box_publish(client, topic):
    PID = random.randint(0 , 2000000)
    print(PID)
    print("wait for the box to arrive to sensor:9.980")
    time.sleep(9.980)
    client.publish(topic, PID)
    print("..............................................................")

client = mqtt.Client()
client.on_connect = on_connect
client.on_message = on_message
client.tls_set(tls_version=mqtt.ssl.PROTOCOL_TLS)
client.username_pw_set(userName, passWord)
print("Connecting to broker::", brokerAddress)
client.connect(brokerAddress, 8883)
client.loop_start()
time.sleep(2)

while True:
    try:
        input("press enter to put the box on line:")
    except SyntaxError:
        pass
    # 启动新线程执行任务,daemon=True确保主线程退出时子线程自动终止
    threading.Thread(target=handle_box_publish, args=(client, topic), daemon=True).start()

client.loop_stop()

关键说明

  1. 线程模块选择:用threading替代thread,thread是Python早期低级API,threading提供更完善的线程管理功能,比如守护线程、同步机制等。
  2. 线程安全:paho-mqtt的Client对象是线程安全的,多个线程同时调用publish方法不会出现冲突。
  3. 异步执行:每次按下回车就启动独立线程,线程自行完成生成PID、等待、发布的流程,主线程可立即回到等待输入状态,支持同时处理多个盒子的流程。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 07:05:24