本文共 1577 字,大约阅读时间需要 5 分钟。
可以通过以下命令安装paho-mqtt库:
pip install paho-mqtt -i http://pypi.douban.com/simple --trusted-host pypi.douban.com
首先导入所需的库:
import paho.mqtt.client as mqttfrom multiprocessing import Process
定义MQTT服务器地址和端口:
MQTTHOST = "172.19.4.4" # 替换为实际MQTT服务器地址MQTTPORT = 1883 # 替换为实际MQTT服务器端口
创建MQTT客户端:
mqttClient = mqtt.Client()
定义连接回调函数:
def on_mqtt_connect(): mqttClient.connect(MQTTHOST, MQTTPORT, 60) mqttClient.loop_start() # 启动循环,处理消息
定义消息到来回调函数:
def on_message_come(client, userdata, msg): print(f"主题 {msg.topic} 收到消息:{msg.payload.decode('utf-8')}") # 启动多进程处理消息 p = Process(target=talk, args=("/camera/person/num/result", msg.payload.decode("utf-8"))) p.start() 订阅指定主题:
def on_subscribe(): mqttClient.subscribe("test", 1) # 主题为"test" mqttClient.on_message = on_message_come 定义发布回调函数:
def on_publish(topic, msg, qos): mqttClient.publish(topic, msg, qos)
在多进程中发布消息时重新初始化MQTT客户端:
def talk(topic, msg): global mqttClient camera_person_num = CameraPsersonNum(msg) t_max, t_mean = camera_person_num.personNum() # 初始化MQTT客户端 mqttClient = mqtt.Client() mqttClient.connect(MQTTHOST, MQTTPORT, 60) mqttClient.loop_start() # 发布处理结果 mqttClient.publish(topic, f'{{"max":{t_max),"mean":{t_mean}}}', 1) 定义主程序函数:
def main(): on_mqtt_connect() on_subscribe() while True: pass
在终端中运行:
python mqtt_example.py
以上代码实现了一个简单的MQTT消息收发应用,支持多进程处理消息并重新初始化MQTT客户端以发布处理结果。可以根据实际需求调整主题、端口和服务器地址。
转载地址:http://joafk.baihongyu.com/