【问题标题】:MQTT Paho Python Client subscriber, how subscribe forever?MQTT Paho Python 客户端订阅者,如何永久订阅?
【发布时间】:2016-10-05 09:34:43
【问题描述】:

尝试在不断开连接到 Mosquitto 代理的情况下进行简单订阅,以从在特定主题中发布数据的设备获取所有消息,将它们保存在 BD 中并将它们发布到执行“工作人员”的 php。

这是我的 subscribe.py

import paho.mqtt.client as mqtt
from mqtt_myapp import *

topic="topic/#"                 # MQTT broker topic
myclient="my-paho-client"           # MQTT broker My Client
user="user"                     # MQTT broker user
pw="pass"                       # MQTT broker password
host="localhost"                # MQTT broker host
port=1883                       # MQTT broker port
value="123"                     # somethin i need for myapp

def on_connect(mqttc, userdata, rc):
    print('connected...rc=' + str(rc))
    mqttc.subscribe(topic, qos=0)

def on_disconnect(mqttc, userdata, rc):
    print('disconnected...rc=' + str(rc))

def on_message(mqttc, userdata, msg):
    print('message received...')
    print('topic: ' + msg.topic + ', qos: ' + 
          str(msg.qos) + ', message: ' + str(msg.payload))
    save_to_db(msg)
    post_data(msg.payload,value)

def on_subscribe(mqttc, userdata, mid, granted_qos):
    print('subscribed (qos=' + str(granted_qos) + ')')

def on_unsubscribe(mqttc, userdata, mid, granted_qos):
    print('unsubscribed (qos=' + str(granted_qos) + ')')

mqttc = mqtt.Client(myclient)
mqttc.on_connect = on_connect
mqttc.on_disconnect = on_disconnect
mqttc.on_message = on_message
mqttc.on_subscribe = on_subscribe
mqttc.on_unsubscribe = on_unsubscribe
mqttc.username_pw_set(user,pw)
mqttc.connect(host, port, 60)
mqttc.loop_forever()

这是我的 mqtt_myapp.py:

import MySQLdb
import requests # pip install requests

url = "mydomain/data_from_broker.php"

def save_to_db(msg):
    with db:
        cursor = db.cursor()
        try:
            cursor.execute("INSERT INTO MQTT_LOGS (topic, payload) VALUES (%s,%s)", (msg.topic, msg.payload))
        except (MySQLdb.Error, MySQLdb.Warning) as e:
            print('excepttion BD ' + e)
            return None

def post_data(payload,value):
    datos = {'VALUE': value,'data-from-broker': payload}
    r = requests.post(url, datos)
    r.status_code
    print('response POST' + str(r.status_code))

db = MySQLdb.connect("localhost","user_db","pass_db","db" )

当我使用 python -t mqtt_subscribe.py & 在后台运行我的 python 脚本时,我会收到为其他客户端发布的消息,但 在我的 subscribe.py 脚本运行几个小时后,会发生套接字错误。

Mosquito.log:

 ...
    1475614815: Received PINGREQ from my-paho-client
    1475614815: Sending PINGRESP to my-paho-client
    1475614872: New connection from xxx.xxx.xxx.xxx on port 1883.
    1475614872: Client device1 disconnected.
    1475614872: New client connected from xxx.xxx.xxx.xxx as device1(c0, k0, u'user1').
    1475614872: Sending CONNACK to device1(0, 0)
    1475614873: Received PUBLISH from device1(d0, q1, r0, m1, 'topic/data', ... (33 bytes))
    1475614873: Sending PUBACK to device1 (Mid: 1)
    1475614873: Sending PUBLISH to my-paho-client (d0, q0, r0, m0, 'topic/data', ... (33 bytes))
    1475614874: Received DISCONNECT from device1
    1475614874: Client device1 disconnected.
...
    1475625566: Received PINGREQ from my-paho-client
    1475625566: Sending PINGRESP to my-paho-client
    1475625626: Received PINGREQ from my-paho-client
    1475625626: Sending PINGRESP to my-paho-client
    1475625675: New connection from xxx.xxx.xxx.xxx on port 1883.
    1475625675: Client device1 disconnected.
    1475625675: New client connected from xxx.xxx.xxx.xxx as device1 (c0, k0, u'user1').
    1475625675: Sending CONNACK to device1 (0, 0)
    1475625677: Received PUBLISH from device1 (d0, q1, r0, m1, 'topic/data', ... (33 bytes))
    1475625677: Sending PUBACK to device1 (Mid: 1)
    1475625677: Sending PUBLISH to my-paho-client (d0, q0, r0, m0, 'topic/data', ... (33 bytes))
    1475625677: Socket error on client my-paho-client, disconnecting.
    1475625677: Received DISCONNECT from device1
...

可能是什么问题?有什么想法或建议吗?

提前致谢

【问题讨论】:

    标签: python mqtt mosquitto paho


    【解决方案1】:

    如果您在方法“on_message”中的代码抛出异常并且您没有捕获它,您将被断开连接。 尝试取消注释除打印语句之外的所有语句。可能以下语句之一正在引发异常。

     save_to_db(msg)
     post_data(msg.payload,value)
    

    【讨论】:

    • 如果我想捕捉“on_message”事件,我应该在哪里编写“save_to_db”代码?
    • 您可以在那里编码,只要确保您在此处捕获所有异常,如果您不想断开连接。
    • 试着理解,所以就因为我没有捕捉到异常,脚本就断开了?即使我捕捉到异常并且什么都不做脚本也不会停止? SebastianK,你能给我一个参考来研究这个问题吗?非常感谢。
    • 我在 java 实现中遇到了这个问题,我猜 python 的行为是一样的。查看 javadoc link --> MqttCallback --> messageArrived。那里描述了这种行为。
    • "脚本已断开" - 我认为未捕获的异常会导致处理 MQTT 连接的线程死亡,从而中断 TCP 连接,即您不再订阅,因为您是'不再连接了。如果您(在通过不可靠/不可预测的网络连接时始终应该这样做)明智地使用 try/except(即不要简单地忽略所有异常,IMO 这比不使用 try/except 更糟糕)您将能够确保异常发生时线程不会死。
    【解决方案2】:

    在插入数据库时​​,使用? 而不是%slink for reference

    cursor.execute('''INSERT INTO users(name, phone, email, password)
                      VALUES(?,?,?,?)''', (name1,phone1, email1, password1))
    

    【讨论】:

    • 这是针对 sqlite 的,不是针对 MySQLdb 的。
    猜你喜欢
    • 1970-01-01
    • 2018-05-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-10-17
    相关资源
    最近更新 更多