Python多线程实现:一个线程定时执行,另一个持续运行
Python多线程整合MQTT定时发送与Web服务的实现方案
嘿,这个需求其实挺常见的,用Python自带的threading模块就能完美解决——把两个独立的任务拆分成不同线程,让它们各自运行互不干扰。我给你写个完整的可落地方案,直接替换你的业务逻辑就能用:
核心思路
- 将原MQTT定时发送逻辑封装成一个线程函数,内部做无限循环,每2分钟执行一次发送操作
- 将Web请求服务(比如用Flask/FastAPI)封装成另一个线程函数,持续监听请求
- 主线程负责启动这两个子线程,并等待它们持续运行
完整代码示例
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
相关产品推荐
相关产品推荐

