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

如何在MQTT循环外导出并打印message变量以供给其他Python程序使用

解决MQTT消息外部访问与导出问题

一、当前代码的问题

你的代码里message是onMessage回调的局部变量,外部无法直接访问;而且client.loop_start()启动后台线程后,主线程立刻执行print(message),此时要么变量未定义报错,要么还没收到消息拿不到有效内容。


二、让消息能在MQTT循环外部打印

方案1:使用全局变量

把消息存储为全局变量,在回调中更新,主线程循环读取并打印:

import paho.mqtt.client as paho 
import time 

# 全局变量存储收到的消息
received_message = ""

def onMessage(client, userdata, msg):
    global received_message
    received_message = str(msg.payload.decode())
    # 可选:回调内打印确认收到消息
    # print(f"回调内收到消息: {received_message}")

client = paho.Client()                   
client.on_message = onMessage       
client.connect("broker.mqtt-dashboard.com", 1883)
client.subscribe("AGV1/posisi") 
client.loop_start()                     

# 主线程循环打印最新消息
try:
    while True:
        if received_message:
            print(f"外部打印消息: {received_message}")
            # 可选:打印后清空,避免重复输出
            # received_message = ""
        time.sleep(1)
except KeyboardInterrupt:
    client.loop_stop()
    client.disconnect()

方案2:使用线程安全队列(更推荐)

用queue.Queue存储消息,避免多线程竞态问题:

import paho.mqtt.client as paho 
import time 
from queue import Queue

# 创建线程安全队列,用于跨线程传递消息
message_queue = Queue()

def onMessage(client, userdata, msg):
    message = str(msg.payload.decode())
    message_queue.put(message)

client = paho.Client()                   
client.on_message = onMessage       
client.connect("broker.mqtt-dashboard.com", 1883)
client.subscribe("AGV1/posisi") 
client.loop_start()                     

# 主线程从队列取消息打印
try:
    while True:
        if not message_queue.empty():
            msg = message_queue.get()
            print(f"外部打印消息: {msg}")
        time.sleep(0.1)
except KeyboardInterrupt:
    client.loop_stop()
    client.disconnect()

三、将消息导出给其他Python程序使用

方案1:本地Socket通信(跨进程实时传递)

让MQTT程序作为Socket服务端,其他Python程序连接后获取消息:

MQT程序(服务端)代码:

import paho.mqtt.client as paho 
import time 
import socket
from threading import Thread
from queue import Queue

message_queue = Queue()

def onMessage(client, userdata, msg):
    message = str(msg.payload.decode())
    message_queue.put(message)

def socket_server():
    # 创建本地TCP Socket
    server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
    server_socket.bind(('localhost', 12345))
    server_socket.listen(5)
    print("Socket服务已启动,等待客户端连接...")
    
    while True:
        conn, addr = server_socket.accept()
        print(f"客户端 {addr} 已连接")
        try:
            # 循环发送最新消息
            while True:
                if not message_queue.empty():
                    msg = message_queue.get()
                    conn.sendall(f"{msg}\n".encode('utf-8'))
                time.sleep(0.1)
        except Exception:
            conn.close()
            print(f"客户端 {addr} 已断开连接")

# 初始化MQTT客户端
client = paho.Client()                   
client.on_message = onMessage       
client.connect("broker.mqtt-dashboard.com", 1883)
client.subscribe("AGV1/posisi") 
client.loop_start()                     

# 启动Socket服务线程
Thread(target=socket_server, daemon=True).start()

# 主线程保持运行
try:
    while True:
        time.sleep(1)
except KeyboardInterrupt:
    client.loop_stop()
    client.disconnect()

其他Python程序(客户端)代码:

import socket

client_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
client_socket.connect(('localhost', 12345))

try:
    while True:
        data = client_socket.recv(1024)
        if not data:
            break
        print(f"从MQTT程序获取的消息: {data.decode('utf-8').strip()}")
except KeyboardInterrupt:
    client_socket.close()

方案2:共享内存(同机器进程间快速访问)

用multiprocessing模块的共享内存实现进程间数据共享:

import paho.mqtt.client as paho 
import time 
from multiprocessing import Array, Process

def onMessage(client, userdata, msg):
    message = str(msg.payload.decode())
    # 更新共享内存(限制最大100字符)
    userdata[:len(message)] = message.encode('utf-8')
    userdata[len(message)] = b'\0'  # 字符串结束符

def mqtt_client(shared_msg):
    client = paho.Client(userdata=shared_msg)
    client.on_message = onMessage       
    client.connect("broker.mqtt-dashboard.com", 1883)
    client.subscribe("AGV1/posisi") 
    client.loop_forever()

if __name__ == "__main__":
    # 创建共享内存,存储最多100个字符
    shared_msg = Array('c', b'\0' * 100)
    # 启动MQTT子进程
    mqtt_process = Process(target=mqtt_client, args=(shared_msg,))
    mqtt_process.start()

    # 主进程读取共享内存中的消息
    try:
        while True:
            msg = shared_msg.value.decode('utf-8').strip('\x00')
            if msg:
                print(f"主进程读取到消息: {msg}")
            time.sleep(1)
    except KeyboardInterrupt:
        mqtt_process.terminate()
        mqtt_process.join()

内容的提问来源于stack exchange,提问作者41_ Andi Yuda

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 07:20:35