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()
关键说明
- 线程模块选择:用
threading替代thread,thread是Python早期低级API,threading提供更完善的线程管理功能,比如守护线程、同步机制等。 - 线程安全:paho-mqtt的
Client对象是线程安全的,多个线程同时调用publish方法不会出现冲突。 - 异步执行:每次按下回车就启动独立线程,线程自行完成生成PID、等待、发布的流程,主线程可立即回到等待输入状态,支持同时处理多个盒子的流程。
内容的提问来源于stack exchange,提问作者mohamad mohamady
相关产品推荐
相关产品推荐

