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

Python连接MQTT网关时,如何每分钟执行一次指定函数?

问题分析与解决方案

你的问题出在两个关键地方:

  1. 每次收到MQTT消息就重复添加定时任务,导致任务堆积,且定时任务的触发逻辑没有被正确执行。
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 17:17:46