Python连接MQTT网关时,如何每分钟执行一次指定函数?
问题分析与解决方案
你的问题出在两个关键地方:
- 每次收到MQTT消息就重复添加定时任务,导致任务堆积,且定时任务的触发逻辑没有被正确执行。
schedule库需要持续调用schedule.run_pending()来检查并执行到期任务,你当前的代码没有这个触发逻辑。
修正步骤与代码示例
1. 修复定时任务初始化逻辑
定时任务只需要初始化一次,把 schedule.every(1).minutes.do(func) 从 on_message 函数中移到主程序里,避免重复创建任务。
2. 确保定时任务被触发
由于MQTT客户端的消息循环通常是阻塞或后台运行的,需要单独维护一个循环来调用 schedule.run_pending(),或者用线程分离定时任务逻辑。
方案一:使用MQTT后台消息循环 + 主循环处理定时任务
import schedule import time import paho.mqtt.client as mqtt # 全局缓冲区 buffer1 = [] buffer2 = [] def func(): # 这里写你的定时任务逻辑,比如处理缓冲区数据 print(f"定时任务执行:buffer1长度={len(buffer1)}, buffer2长度={len(buffer2)}") def on_message(client, userdata, msg): global buffer1 global buffer2 if msg.topic == "test/topic": buffer1.append(msg.payload) print("buffer1", buffer1) elif msg.topic == "test/topic2": buffer2.append(msg.payload) print("buffer2", buffer2) # MQTT客户端配置 client = mqtt.Client() client.on_message = on_message client.connect("你的MQTT broker地址", 1883, 60) client.subscribe([("test/topic", 0), ("test/topic2", 0)]) # 初始化定时任务(仅执行一次) schedule.every(1).minutes.do(func) # 启动MQTT后台消息处理线程 client.loop_start() # 主循环:持续检查定时任务 try: while True: schedule.run_pending() time.sleep(1) # 降低CPU占用 except KeyboardInterrupt: client.loop_stop() print("程序已退出")
方案二:使用线程分离定时任务(适配阻塞式MQTT循环)
如果你的代码使用 client.loop_forever() 这种阻塞式循环,可以用单独线程来运行定时任务逻辑:
import schedule import time import threading import paho.mqtt.client as mqtt buffer1 = [] buffer2 = [] def func(): print(f"定时任务执行:buffer1长度={len(buffer1)}, buffer2长度={len(buffer2)}") def on_message(client, userdata, msg): global buffer1 global buffer2 if msg.topic == "test/topic": buffer1.append(msg.payload) print("buffer1", buffer1) elif msg.topic == "test/topic2": buffer2.append(msg.payload) print("buffer2", buffer2) # 定时任务循环线程 def schedule_runner(): while True: schedule.run_pending() time.sleep(1) # MQTT客户端配置 client = mqtt.Client() client.on_message = on_message client.connect("你的MQTT broker地址", 1883, 60) client.subscribe([("test/topic", 0), ("test/topic2", 0)]) # 初始化定时任务并启动线程 schedule.every(1).minutes.do(func) threading.Thread(target=schedule_runner, daemon=True).start() # 启动阻塞式MQTT消息循环 client.loop_forever()
关键说明
schedule.run_pending()是触发定时任务的核心,必须定期调用,否则任务永远不会执行。- 不要在
on_message内重复创建定时任务,否则会导致同一任务被添加成百上千次,执行逻辑混乱。
内容的提问来源于stack exchange,提问作者user18430327
相关产品推荐
相关产品推荐

