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

Python多线程实现:一个线程定时执行,另一个持续运行

Python多线程整合MQTT定时发送与Web服务的实现方案

嘿,这个需求其实挺常见的,用Python自带的threading模块就能完美解决——把两个独立的任务拆分成不同线程,让它们各自运行互不干扰。我给你写个完整的可落地方案,直接替换你的业务逻辑就能用:

核心思路

  1. 将原MQTT定时发送逻辑封装成一个线程函数,内部做无限循环,每2分钟执行一次发送操作
  2. 将Web请求服务(比如用Flask/FastAPI)封装成另一个线程函数,持续监听请求
  3. 主线程负责启动这两个子线程,并等待它们持续运行

完整代码示例

import threading
import time
import paho.mqtt.client as mqtt
from flask import Flask, request

# ------------------- 配置参数(根据你的实际情况修改) -------------------
MQTT_BROKER = "your_mqtt_broker_ip"
MQTT_PORT = 1883
MQTT_TOPIC = "sensor/data"
WEB_SERVICE_PORT = 5000
SEND_INTERVAL = 120  # 2分钟,单位秒
# ---------------------------------------------------------------------

def mqtt_periodic_sender():
    """MQTT定时发送线程函数"""
    # 初始化MQTT客户端
    client = mqtt.Client()
    # 如果Broker需要认证,取消下面两行注释并填写账号密码
    # client.username_pw_set("mqtt_username", "mqtt_password")
    
    try:
        # 连接MQTT Broker并启动后台消息循环
        client.connect(MQTT_BROKER, MQTT_PORT, 60)
        client.loop_start()
        
        while True:
            # ------------------- 替换成你的数据生成逻辑 -------------------
            # 示例:生成模拟数据
            sensor_data = {"temperature": 25.6, "humidity": 62}
            # -------------------------------------------------------------
            
            # 发送数据到MQTT主题
            print(f"[{time.strftime('%Y-%m-%d %H:%M:%S')}] 发送MQTT数据: {sensor_data}")
            client.publish(MQTT_TOPIC, str(sensor_data))
            
            # 等待2分钟
            time.sleep(SEND_INTERVAL)
            
    except Exception as e:
        print(f"MQTT线程异常: {str(e)}")
        # 可选:添加重连逻辑,比如捕获连接异常后重试
    finally:
        # 清理资源
        client.loop_stop()
        client.disconnect()

def web_request_handler():
    """Web服务线程函数(以Flask为例)"""
    app = Flask(__name__)
    
    # ------------------- 替换成你的Web路由逻辑 -------------------
    @app.route("/api/receive", methods=["POST"])
    def receive_data():
        if request.is_json:
            received_data = request.get_json()
            print(f"[{time.strftime('%Y-%m-%d %H:%M:%S')}] 收到Web请求数据: {received_data}")
            return {"status": "success", "message": "数据已接收"}
        return {"status": "error", "message": "仅接受JSON格式数据"}, 400
    
    @app.route("/health", methods=["GET"])
    def health_check():
        return {"status": "running", "service": "web_handler"}
    # -------------------------------------------------------------
    
    # 启动Web服务,必须设置use_reloader=False避免重复创建线程
    app.run(host="0.0.0.0", port=WEB_SERVICE_PORT, use_reloader=False)

if __name__ == "__main__":
    # 创建两个线程
    mqtt_thread = threading.Thread(target=mqtt_periodic_sender)
    web_thread = threading.Thread(target=web_request_handler)
    
    # 可选:设置守护线程(主线程退出时子线程自动终止)
    # mqtt_thread.daemon = True
    # web_thread.daemon = True
    
    # 启动线程
    print("启动MQTT定时发送线程...")
    mqtt_thread.start()
    print("启动Web服务线程...")
    web_thread.start()
    
    # 等待线程结束(因为两个线程都是无限循环,主线程会一直阻塞在这里)
    mqtt_thread.join()
    web_thread.join()

关键细节说明

  • MQTT客户端的loop_start():启动后台线程处理MQTT的网络交互,避免阻塞定时发送的主逻辑
  • Flask的use_reloader=False:默认自动重载会创建新进程,导致线程重复启动,必须关闭
  • 异常处理:在MQTT线程中添加了异常捕获,你可以根据需要扩展重连逻辑,避免网络波动导致线程挂掉
  • 线程安全:如果两个线程需要共享数据(比如Web接收的数据要通过MQTT发送),记得用threading.Lock加锁,避免竞争条件
  • 守护线程设置:如果希望主线程退出时子线程也跟着终止,就打开daemon=True的注释;否则主线程会一直等待子线程运行

适配其他Web框架

如果你用的是FastAPI,只需要修改web_request_handler函数:

from fastapi import FastAPI
import uvicorn

def web_request_handler():
    app = FastAPI()
    
    @app.post("/api/receive")
    def receive_data(data: dict):
        print(f"收到Web请求数据: {data}")
        return {"status": "success"}
    
    uvicorn.run(app, host="0.0.0.0", port=WEB_SERVICE_PORT)

内容的提问来源于stack exchange,提问作者Alexandru Tarna

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 09:29:51