用Python打造MQTT服务器与可视化管理平台
1. 引言
1.1 什么是MQTT?
MQTT(Message Queuing Telemetry Transport,消息队列遥测传输)是一种基于发布/订阅(Publish/Subscribe)模式的轻量级物联网通信协议。它运行在TCP/IP协议栈之上,由IBM在1999年发布,目前已成为ISO标准(ISO/IEC 20922)。MQTT的设计目标是实现简单、开放、轻量,非常适合在计算能力有限、网络带宽低、延迟高的环境下工作,因此被广泛应用于物联网、移动互联网、车联网等场景。
MQTT的核心组件包括:
-
客户端(Client):任何连接到MQTT服务器的设备或应用程序,可以是发布者(Publisher)或订阅者(Subscriber)。
-
服务器(Broker):负责接收所有客户端的消息,并根据订阅关系将消息路由给对应的订阅者。
MQTT的主要特性:
-
发布/订阅模式:解耦了消息的发送者和接收者,客户端无需知道对方的存在。
-
三种服务质量(QoS):最多一次(0)、至少一次(1)、 exactly一次(2),满足不同可靠性需求。
-
主题(Topic):消息的分类标识,支持层次结构(如
home/livingroom/temperature)和通配符(+单级,#多级)。 -
保留消息(Retained Message):broker可以为每个主题保留最后一条消息,新订阅的客户端能立即收到。
-
遗嘱消息(Will Message):客户端异常断开时,broker代为发布预设的消息,通知其他客户端。
1.2 为什么用Python打造MQTT服务器与管理平台?
Python以其简洁的语法、丰富的第三方库和快速的开发能力,成为物联网原型开发的首选语言之一。使用Python我们可以:
-
深入理解MQTT协议:通过自己编写简单的broker,可以更清晰地理解协议细节。
-
快速构建管理平台:利用Flask/Django等Web框架,配合前端技术,快速搭建可视化管理界面。
-
灵活扩展:Python生态中有大量与MQTT交互的库(如paho-mqtt),可以方便地集成到各种应用中。
本教程将带领读者从零开始,先了解MQTT协议细节,然后搭建一个基于Python的MQTT broker(可选使用开源Mosquitto或自实现),接着开发一个功能完善的可视化管理平台,实现对连接设备、消息流、系统状态的实时监控和管理。整个项目包含详细的代码示例、架构设计和部署指导,最终完成一个可运行的物联网基础平台。
2. MQTT协议详解
在动手实现之前,深入理解MQTT协议的工作原理至关重要。本章将详细解析MQTT协议的报文结构、连接流程、主题匹配规则以及QoS机制。
2.1 MQTT协议报文结构
MQTT协议基于TCP,数据包由三部分组成:
-
固定头(Fixed Header):所有报文都包含,长度为2~5字节。
-
可变头(Variable Header):部分报文包含,长度可变。
-
有效载荷(Payload):部分报文包含,如PUBLISH消息内容。
固定头格式
| Bit | 7-4 | 3-0 |
|---|---|---|
| Byte1 | 报文类型 | 标志位(DUP, QoS, RETAIN) |
| Byte2 | 剩余长度(Remaining Length) |
-
报文类型(4位):共14种类型,如CONNECT(1)、CONNACK(2)、PUBLISH(3)、SUBSCRIBE(8)等。
-
标志位(4位):不同报文的含义不同。例如PUBLISH报文使用DUP、QoS、RETAIN标志。
-
剩余长度:表示当前报文剩余部分的字节数(可变头+载荷)。编码使用可变长度机制:每个字节的低7位表示数值,最高位为1表示后续还有字节。范围0~268,435,455。
可变头与载荷
可变头的内容因报文类型而异。以PUBLISH报文为例,可变头包含:
-
主题名(Topic Name):UTF-8编码字符串,长度由两个字节指定。
-
报文标识符(Packet ID):仅当QoS>0时存在,用于确认消息。
载荷即消息内容(二进制数据)。
2.2 连接与认证
客户端与broker建立TCP连接后,第一个报文必须是CONNECT报文。CONNECT报文包含协议名(MQTT)、协议级别(v3.1.1为4)、连接标志(用户名/密码、遗嘱、清除会话等)、保活时间等。broker收到后响应CONNACK报文,包含连接返回码(0表示成功)。
安全考虑:MQTT支持用户名/密码认证,也可通过TLS加密传输。
2.3 主题与通配符
主题是UTF-8字符串,可以使用斜杠/分隔层级。例如:sensor/1/temperature。
订阅主题时可以使用通配符:
-
单级通配符
+:匹配一个层级。例如sensor/+/temperature可匹配sensor/1/temperature、sensor/2/temperature,但不匹配sensor/1/floor/temperature。 -
多级通配符
#:匹配任意层级,必须放在最后。例如sensor/#可匹配sensor/1/temperature、sensor/1、sensor/1/floor/humidity等。
2.4 服务质量(QoS)
MQTT定义了三种消息服务质量等级,影响消息传递的可靠性:
-
QoS 0(最多一次):消息只发送一次,不确认,可能丢失。
-
QoS 1(至少一次):消息保证到达,但可能重复。通过PUBACK确认。
-
QoS 2( exactly一次):消息保证到达且不重复。通过两轮确认(PUBREC, PUBREL, PUBCOMP)实现。
broker和客户端都需要根据QoS等级处理消息确认和重传。
2.5 保留消息与遗嘱消息
-
保留消息:发布者可以将消息标记为“保留”。broker会为每个主题保存最新的一条保留消息,当新订阅者订阅该主题时,立即收到该保留消息。
-
遗嘱消息:客户端在连接时指定遗嘱主题和消息。当客户端异常断开(如网络中断、心跳超时)时,broker会代为发布遗嘱消息,通知其他客户端。
3. 搭建Python MQTT Broker
实现一个功能完善的MQTT broker非常复杂,涉及并发、状态管理、消息路由、QoS处理等。在实际项目中,通常直接使用成熟的开源broker,如Mosquitto、EMQX、VerneMQ等。但为了深入理解协议,我们也可以自己实现一个简单的broker,处理基本的CONNECT、PUBLISH、SUBSCRIBE、PINGREQ等报文。
本章先介绍如何使用Mosquitto快速搭建生产可用的broker,然后给出一个自实现简易broker的Python代码示例,帮助理解内部原理。
3.1 使用Mosquitto搭建MQTT Broker
Mosquitto是Eclipse基金会下的开源MQTT broker,支持MQTT v3.1/v3.1.1/v5.0,轻量且易于安装。以下是在Ubuntu上的安装配置步骤。
3.1.1 安装Mosquitto
bash
sudo apt update sudo apt install -y mosquitto mosquitto-clients
安装完成后,Mosquitto会自动以系统服务启动,默认监听1883端口(非加密)和8883端口(TLS)。可以通过以下命令查看服务状态:
bash
sudo systemctl status mosquitto
3.1.2 配置Mosquitto
默认配置文件位于/etc/mosquitto/mosquitto.conf。我们可以修改配置以支持用户认证、WebSocket、日志等。
常用配置项示例:
conf
# 监听1883端口 listener 1883 0.0.0.0 protocol mqtt # 监听WebSocket端口(用于浏览器连接) listener 9001 0.0.0.0 protocol websockets # 允许匿名访问(生产环境应关闭) allow_anonymous true # 密码文件(需用mosquitto_passwd创建) # password_file /etc/mosquitto/passwd # 日志 log_dest file /var/log/mosquitto/mosquitto.log log_type all
如果需要用户认证,先创建密码文件:
bash
sudo mosquitto_passwd -c /etc/mosquitto/passwd user1 # 按提示输入密码
然后取消配置文件中password_file行的注释,并设置allow_anonymous false。重启Mosquitto生效:
bash
sudo systemctl restart mosquitto
3.1.3 测试连接
使用mosquitto客户端工具测试:
bash
# 订阅主题(在一个终端) mosquitto_sub -h localhost -t "test/topic" # 发布消息(在另一个终端) mosquitto_pub -h localhost -t "test/topic" -m "Hello MQTT"
应能看到订阅终端打印出消息。
3.2 自实现一个简易Python MQTT Broker
为了学习,我们可以用Python的socket编程实现一个最基本的broker,处理部分报文类型。这里不追求完整实现所有特性,而是展示核心逻辑。
3.2.1 设计思路
-
使用
socket创建TCP服务器,监听1883端口。 -
为每个客户端连接创建一个线程(或使用asyncio)处理。
-
解析客户端发送的报文,根据报文类型执行相应操作。
-
维护一个订阅关系字典:
topic -> set(client_sockets)。 -
维护一个保留消息字典:
topic -> (payload, qos)。 -
处理CONNECT、PUBLISH、SUBSCRIBE、PINGREQ、DISCONNECT等报文。
3.2.2 代码实现
下面给出一个简化版的broker代码,仅支持QoS 0,忽略用户名/密码认证,无遗嘱消息。为了清晰,我们使用Python的socketserver模块或直接使用socket和threading。
python
import socket
import threading
import struct
# MQTT报文类型常量
CONNECT = 1
CONNACK = 2
PUBLISH = 3
PUBACK = 4
SUBSCRIBE = 8
SUBACK = 9
PINGREQ = 12
PINGRESP = 13
DISCONNECT = 14
class MQTTServer:
def __init__(self, host='0.0.0.0', port=1883):
self.host = host
self.port = port
self.sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
self.sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
self.sock.bind((self.host, self.port))
self.sock.listen(5)
self.clients = [] # 保存所有客户端socket,用于广播
self.subscriptions = {} # 主题 -> set(客户端socket)
self.retained_messages = {} # 主题 -> (payload, qos)
self.lock = threading.Lock()
print(f"MQTT Broker listening on {self.host}:{self.port}")
def start(self):
while True:
client_sock, addr = self.sock.accept()
print(f"New connection from {addr}")
with self.lock:
self.clients.append(client_sock)
client_thread = threading.Thread(target=self.handle_client, args=(client_sock, addr))
client_thread.daemon = True
client_thread.start()
def handle_client(self, client_sock, addr):
# 客户端处理循环
try:
while True:
# 读取固定头第一字节
header = client_sock.recv(1)
if not header:
break
# 解析固定头
byte1 = header[0]
msg_type = (byte1 & 0xF0) >> 4
flags = byte1 & 0x0F
# 读取剩余长度
remaining_length = self.decode_remaining_length(client_sock)
# 根据消息类型处理
if msg_type == CONNECT:
self.handle_connect(client_sock, remaining_length)
elif msg_type == PUBLISH:
self.handle_publish(client_sock, flags, remaining_length)
elif msg_type == SUBSCRIBE:
self.handle_subscribe(client_sock, remaining_length)
elif msg_type == PINGREQ:
self.handle_ping(client_sock)
elif msg_type == DISCONNECT:
break
else:
# 忽略未知类型
client_sock.recv(remaining_length)
except Exception as e:
print(f"Error handling client {addr}: {e}")
finally:
self.remove_client(client_sock, addr)
def decode_remaining_length(self, sock):
multiplier = 1
value = 0
while True:
byte = sock.recv(1)
if not byte:
raise Exception("Connection closed")
b = byte[0]
value += (b & 0x7F) * multiplier
multiplier *= 128
if multiplier > 128*128*128:
raise Exception("Malformed remaining length")
if (b & 0x80) == 0:
break
return value
def encode_remaining_length(self, length):
enc = []
while True:
digit = length % 128
length //= 128
if length > 0:
digit |= 0x80
enc.append(digit)
if length == 0:
break
return bytes(enc)
def handle_connect(self, sock, remaining_length):
# 读取可变头+载荷,但这里简化,只检查协议名
data = sock.recv(remaining_length)
# 解析协议名(前两个字节是长度)
proto_len = struct.unpack("!H", data[:2])[0]
proto_name = data[2:2+proto_len].decode()
if proto_name != "MQTT":
# 协议名错误,关闭连接
sock.close()
return
# 协议级别
proto_level = data[2+proto_len]
# 连接标志
connect_flags = data[2+proto_len+1]
# 保活时间
keep_alive = struct.unpack("!H", data[2+proto_len+2:2+proto_len+4])[0]
# 发送CONNACK
connack = bytes([CONNACK << 4, 2, 0, 0]) # 固定头+可变头(连接确认标志0,返回码0)
sock.send(connack)
print("Sent CONNACK")
def handle_publish(self, sock, flags, remaining_length):
# flags: DUP, QoS, RETAIN
data = sock.recv(remaining_length)
# 解析主题
topic_len = struct.unpack("!H", data[:2])[0]
topic = data[2:2+topic_len].decode()
payload_start = 2 + topic_len
# 如果QoS>0,还有报文标识符,这里假设QoS=0
payload = data[payload_start:].decode()
print(f"Received publish on {topic}: {payload}")
# 检查保留标志
retain = (flags & 0x01)
if retain:
with self.lock:
self.retained_messages[topic] = (payload, 0) # 假设qos=0
# 转发给订阅该主题的客户端
self.deliver_message(topic, payload)
# 如果QoS=1,需要回复PUBACK,这里省略
def deliver_message(self, topic, payload):
with self.lock:
# 获取订阅该主题的所有客户端(这里简化,不支持通配符)
if topic in self.subscriptions:
for client in self.subscriptions[topic].copy():
try:
# 构造PUBLISH报文
# 固定头:报文类型PUBLISH,QoS=0,保留标志=0
fixed_header = bytes([PUBLISH << 4])
# 主题
topic_encoded = struct.pack("!H", len(topic)) + topic.encode()
# 载荷
payload_encoded = payload.encode()
# 剩余长度 = 主题长度+载荷长度
remaining = len(topic_encoded) + len(payload_encoded)
remaining_enc = self.encode_remaining_length(remaining)
packet = fixed_header + remaining_enc + topic_encoded + payload_encoded
client.send(packet)
except:
# 发送失败,可能客户端断开,后续会移除
pass
def handle_subscribe(self, sock, remaining_length):
data = sock.recv(remaining_length)
# 报文标识符
packet_id = struct.unpack("!H", data[:2])[0]
# 解析主题过滤器列表
pos = 2
topics = []
while pos < len(data):
topic_len = struct.unpack("!H", data[pos:pos+2])[0]
pos += 2
topic = data[pos:pos+topic_len].decode()
pos += topic_len
qos = data[pos]
pos += 1
topics.append((topic, qos))
# 更新订阅关系
with self.lock:
for topic, qos in topics:
if topic not in self.subscriptions:
self.subscriptions[topic] = set()
self.subscriptions[topic].add(sock)
print(f"Client subscribed to {topic}")
# 如果有保留消息,立即发送
if topic in self.retained_messages:
payload, _ = self.retained_messages[topic]
# 构造保留消息的PUBLISH
fixed_header = bytes([(PUBLISH << 4) | 0x01]) # 保留标志置1
topic_encoded = struct.pack("!H", len(topic)) + topic.encode()
payload_encoded = payload.encode()
remaining = len(topic_encoded) + len(payload_encoded)
remaining_enc = self.encode_remaining_length(remaining)
packet = fixed_header + remaining_enc + topic_encoded + payload_encoded
sock.send(packet)
# 发送SUBACK
# 返回码根据订阅的每个主题的QoS(这里简化为0)
suback_payload = bytes([packet_id >> 8, packet_id & 0xFF] + [0] * len(topics))
fixed = bytes([(SUBACK << 4), len(suback_payload)])
sock.send(fixed + suback_payload)
def handle_ping(self, sock):
# 响应PINGRESP
sock.send(bytes([PINGRESP << 4, 0]))
def remove_client(self, sock, addr):
with self.lock:
if sock in self.clients:
self.clients.remove(sock)
# 从所有订阅中移除
for topic, clients in self.subscriptions.items():
clients.discard(sock)
try:
sock.close()
except:
pass
print(f"Client {addr} disconnected")
if __name__ == "__main__":
server = MQTTServer()
server.start()
3.2.3 运行与测试
保存上述代码为simple_mqtt_broker.py,运行:
bash
python3 simple_mqtt_broker.py
使用mosquitto客户端或paho-mqtt客户端测试连接、发布和订阅。
注意:此实现非常简陋,不支持QoS>0、通配符、用户名密码、心跳超时等,仅用于教学理解。生产环境请使用成熟broker。
4. MQTT客户端开发(Python)
无论使用哪种broker,我们需要编写客户端来模拟设备或应用程序与broker交互。Python中最常用的MQTT客户端库是paho-mqtt。
4.1 安装paho-mqtt
bash
pip install paho-mqtt
4.2 基础用法
连接与回调
python
import paho.mqtt.client as mqtt
# 连接回调
def on_connect(client, userdata, flags, rc):
print("Connected with result code " + str(rc))
# 订阅主题
client.subscribe("test/topic")
# 消息回调
def on_message(client, userdata, msg):
print(f"Received message on {msg.topic}: {msg.payload.decode()}")
client = mqtt.Client()
client.on_connect = on_connect
client.on_message = on_message
client.connect("localhost", 1883, 60)
# 阻塞循环(自动重连)
client.loop_forever()
发布消息
python
import paho.mqtt.client as mqtt
client = mqtt.Client()
client.connect("localhost", 1883, 60)
client.loop_start() # 启动网络循环
# 发布QoS=0消息
client.publish("test/topic", "Hello from Python")
# 发布保留消息
client.publish("test/topic", "Retained message", retain=True)
# 保持运行一段时间
import time
time.sleep(2)
client.loop_stop()
4.3 高级特性
-
遗嘱消息:在连接时设置
python
client.will_set("will/topic", "Client disconnected unexpectedly", qos=1, retain=False)
-
TLS加密:使用
client.tls_set()配置证书。 -
用户名密码:
client.username_pw_set("user", "password")。
4.4 模拟设备数据发送
我们可以编写一个简单的脚本,模拟温度传感器每隔5秒发布随机温度:
python
import paho.mqtt.client as mqtt
import random
import time
import json
client = mqtt.Client()
client.connect("localhost", 1883)
while True:
temperature = random.uniform(20.0, 30.0)
humidity = random.uniform(40.0, 60.0)
payload = json.dumps({"temp": round(temperature,2), "hum": round(humidity,2)})
client.publish("sensor/temperature", payload)
print(f"Published: {payload}")
time.sleep(5)
5. 可视化管理平台设计
现在进入核心部分:构建一个Web可视化管理平台,用于监控MQTT broker状态、管理设备、查看实时数据、发布消息等。我们将采用Flask作为后端框架,结合Bootstrap前端库和ECharts图表库,并通过WebSocket实现实时数据推送。
5.1 技术选型
-
后端:Flask(轻量、易扩展),Flask-SocketIO(实现WebSocket),Flask-SQLAlchemy(数据库ORM,可选),Celery(后台任务,可选)。
-
前端:Bootstrap 5(响应式UI),jQuery(简化DOM操作),ECharts(图表),Socket.IO客户端(实时通信)。
-
数据库:根据需求选择。本教程使用SQLite存储设备信息、历史消息等,方便演示。
-
MQTT客户端:paho-mqtt,用于与管理平台后端集成,订阅特定主题接收消息。
5.2 系统架构
管理平台作为一个MQTT客户端连接到broker,订阅所有需要监控的主题(如#或特定前缀),接收到的消息经过处理后存入数据库,同时通过WebSocket推送到前端页面。前端提供仪表盘、设备列表、消息历史、发布表单等界面。
https://%E5%BE%85%E8%A1%A5%E5%85%85
5.3 后端开发(Flask)
5.3.1 初始化项目
创建项目目录结构:
text
mqtt_manager/ ├── app.py # Flask主应用 ├── mqtt_client.py # MQTT客户端封装 ├── models.py # 数据库模型 ├── static/ # 静态文件(CSS, JS) ├── templates/ # HTML模板 ├── config.py # 配置文件 └── requirements.txt # 依赖
5.3.2 安装依赖
txt
Flask Flask-SocketIO Flask-SQLAlchemy paho-mqtt eventlet # 用于WebSocket生产环境
安装:pip install -r requirements.txt
5.3.3 配置文件 config.py
python
import os
class Config:
SECRET_KEY = os.environ.get('SECRET_KEY') or 'hard-to-guess-string'
SQLALCHEMY_DATABASE_URI = 'sqlite:///mqtt_manager.db'
SQLALCHEMY_TRACK_MODIFICATIONS = False
MQTT_BROKER_HOST = 'localhost'
MQTT_BROKER_PORT = 1883
MQTT_TOPIC_SUBSCRIBE = '#' # 订阅所有主题
5.3.4 数据库模型 models.py
使用SQLAlchemy定义两个简单模型:设备(Device)和消息(Message)。设备可以通过首次发布消息自动注册。
python
from flask_sqlalchemy import SQLAlchemy
from datetime import datetime
db = SQLAlchemy()
class Device(db.Model):
id = db.Column(db.Integer, primary_key=True)
client_id = db.Column(db.String(128), unique=True, nullable=False)
first_seen = db.Column(db.DateTime, default=datetime.utcnow)
last_seen = db.Column(db.DateTime, default=datetime.utcnow, onupdate=datetime.utcnow)
ip_address = db.Column(db.String(64))
# 可以扩展更多字段
class Message(db.Model):
id = db.Column(db.Integer, primary_key=True)
topic = db.Column(db.String(256), nullable=False)
payload = db.Column(db.Text)
qos = db.Column(db.Integer)
retain = db.Column(db.Boolean)
timestamp = db.Column(db.DateTime, default=datetime.utcnow)
# 可选关联设备
device_id = db.Column(db.Integer, db.ForeignKey('device.id'))
5.3.5 MQTT客户端封装 mqtt_client.py
该类负责连接broker,处理消息,并触发SocketIO事件。
python
import paho.mqtt.client as mqtt
from flask_socketio import SocketIO
from models import db, Device, Message
from datetime import datetime
import json
class MQTTClient:
def __init__(self, app, socketio):
self.app = app
self.socketio = socketio
self.client = mqtt.Client()
self.client.on_connect = self.on_connect
self.client.on_message = self.on_message
self.broker_host = app.config['MQTT_BROKER_HOST']
self.broker_port = app.config['MQTT_BROKER_PORT']
self.topic_sub = app.config['MQTT_TOPIC_SUBSCRIBE']
def on_connect(self, client, userdata, flags, rc):
print("MQTT connected with result code " + str(rc))
client.subscribe(self.topic_sub)
def on_message(self, client, userdata, msg):
# 在应用上下文中处理数据库操作
with self.app.app_context():
try:
payload_str = msg.payload.decode('utf-8', errors='ignore')
except:
payload_str = str(msg.payload)
# 保存消息到数据库
message = Message(topic=msg.topic, payload=payload_str,
qos=msg.qos, retain=msg.retain)
db.session.add(message)
# 尝试提取client_id(假设设备在连接时指定了client_id,但我们没有连接信息,这里简单从主题推断)
# 这里简化:从主题第一部分作为设备标识
parts = msg.topic.split('/')
if len(parts) > 0:
client_id = parts[0] # 假设主题格式 device_id/...
device = Device.query.filter_by(client_id=client_id).first()
if not device:
device = Device(client_id=client_id)
db.session.add(device)
device.last_seen = datetime.utcnow()
# 更新消息与设备关联
message.device_id = device.id
db.session.commit()
# 通过SocketIO推送到前端
self.socketio.emit('new_message', {
'topic': msg.topic,
'payload': payload_str,
'timestamp': message.timestamp.isoformat()
})
# 如果是特定主题,也可以解析JSON后推送更详细的数据
if msg.topic.endswith('/temperature'):
try:
data = json.loads(payload_str)
self.socketio.emit('sensor_data', data)
except:
pass
def start(self):
self.client.connect(self.broker_host, self.broker_port, 60)
# 使用loop_start启动后台线程
self.client.loop_start()
def stop(self):
self.client.loop_stop()
self.client.disconnect()
def publish(self, topic, payload, qos=0, retain=False):
self.client.publish(topic, payload, qos, retain)
5.3.6 Flask主应用 app.py
python
from flask import Flask, render_template, request, jsonify
from flask_socketio import SocketIO, emit
from models import db, Device, Message
from mqtt_client import MQTTClient
from config import Config
import json
app = Flask(__name__)
app.config.from_object(Config)
db.init_app(app)
socketio = SocketIO(app, cors_allowed_origins="*")
# 初始化MQTT客户端
mqtt_client = MQTTClient(app, socketio)
@app.before_first_request
def create_tables():
db.create_all()
@app.route('/')
def index():
return render_template('index.html')
@app.route('/api/devices')
def get_devices():
devices = Device.query.order_by(Device.last_seen.desc()).all()
return jsonify([{
'id': d.id,
'client_id': d.client_id,
'first_seen': d.first_seen.isoformat(),
'last_seen': d.last_seen.isoformat()
} for d in devices])
@app.route('/api/messages')
def get_messages():
limit = request.args.get('limit', 50, type=int)
messages = Message.query.order_by(Message.timestamp.desc()).limit(limit).all()
return jsonify([{
'id': m.id,
'topic': m.topic,
'payload': m.payload,
'timestamp': m.timestamp.isoformat()
} for m in messages])
@app.route('/api/publish', methods=['POST'])
def publish():
data = request.json
topic = data.get('topic')
payload = data.get('payload')
qos = data.get('qos', 0)
retain = data.get('retain', False)
if not topic or payload is None:
return jsonify({'error': 'Missing topic or payload'}), 400
mqtt_client.publish(topic, payload, qos, retain)
return jsonify({'status': 'ok'})
@socketio.on('connect')
def handle_connect():
print('Client connected to WebSocket')
@socketio.on('disconnect')
def handle_disconnect():
print('Client disconnected')
if __name__ == '__main__':
mqtt_client.start()
try:
socketio.run(app, debug=True, host='0.0.0.0', port=5000)
finally:
mqtt_client.stop()
5.4 前端开发
使用Bootstrap 5构建响应式界面,包含导航栏、侧边菜单、仪表盘、设备列表、消息历史、发布表单等页面。我们使用单页应用风格,通过Ajax加载数据,并用Socket.IO实时更新。
5.4.1 基础模板 templates/base.html
html
<!DOCTYPE html>
<html lang="en">
<head>
<meta charset="UTF-8">
<meta name="viewport" content="width=device-width, initial-scale=1.0">
<title>MQTT Manager</title>
<link href="https://cdn.jsdelivr.net/npm/bootstrap@5.1.3/dist/css/bootstrap.min.css" rel="stylesheet">
<script src="https://code.jquery.com/jquery-3.6.0.min.js"></script>
<script src="https://cdn.jsdelivr.net/npm/bootstrap@5.1.3/dist/js/bootstrap.bundle.min.js"></script>
<script src="https://cdn.socket.io/4.5.0/socket.io.min.js"></script>
<script src="https://cdn.jsdelivr.net/npm/echarts@5.4.3/dist/echarts.min.js"></script>
{% block head %}{% endblock %}
</head>
<body>
<nav class="navbar navbar-dark bg-dark">
<div class="container-fluid">
<span class="navbar-brand">MQTT 可视化管理平台</span>
</div>
</nav>
<div class="container-fluid">
<div class="row">
<nav class="col-md-2 d-md-block bg-light sidebar" style="min-height: calc(100vh - 56px);">
<div class="position-sticky pt-3">
<ul class="nav flex-column">
<li class="nav-item"><a class="nav-link" href="/">仪表盘</a></li>
<li class="nav-item"><a class="nav-link" href="/devices">设备列表</a></li>
<li class="nav-item"><a class="nav-link" href="/messages">消息历史</a></li>
<li class="nav-item"><a class="nav-link" href="/publish">发布消息</a></li>
</ul>
</div>
</nav>
<main class="col-md-10 ms-sm-auto px-md-4">
{% block content %}{% endblock %}
</main>
</div>
</div>
<script>
// 全局Socket.IO连接
var socket = io();
socket.on('connect', function() {
console.log('WebSocket connected');
});
</script>
{% block scripts %}{% endblock %}
</body>
</html>
5.4.2 仪表盘页面 templates/index.html
html
{% extends "base.html" %}
{% block content %}
<h1 class="h2">仪表盘</h1>
<div class="row">
<div class="col-md-4">
<div class="card text-white bg-primary mb-3">
<div class="card-header">在线设备</div>
<div class="card-body">
<h5 class="card-title" id="online-count">0</h5>
</div>
</div>
</div>
<div class="col-md-4">
<div class="card text-white bg-success mb-3">
<div class="card-header">今日消息</div>
<div class="card-body">
<h5 class="card-title" id="today-msg">0</h5>
</div>
</div>
</div>
<div class="col-md-4">
<div class="card text-white bg-info mb-3">
<div class="card-header">主题数量</div>
<div class="card-body">
<h5 class="card-title" id="topic-count">0</h5>
</div>
</div>
</div>
</div>
<div class="row">
<div class="col-md-6">
<div class="card">
<div class="card-header">最近消息</div>
<div class="card-body" style="max-height: 300px; overflow-y: auto;">
<ul class="list-group" id="recent-messages">
</ul>
</div>
</div>
</div>
<div class="col-md-6">
<div class="card">
<div class="card-header">温度实时曲线</div>
<div class="card-body">
<div id="tempChart" style="width:100%; height:300px;"></div>
</div>
</div>
</div>
</div>
{% endblock %}
{% block scripts %}
<script>
$(document).ready(function() {
// 初始化echarts
var chartDom = document.getElementById('tempChart');
var myChart = echarts.init(chartDom);
var option = {
title: { text: '温度实时数据' },
xAxis: { type: 'category', data: [] },
yAxis: { type: 'value' },
series: [{ data: [], type: 'line' }]
};
myChart.setOption(option);
// 加载最近消息
function loadRecent() {
$.get('/api/messages?limit=10', function(data) {
var list = $('#recent-messages');
list.empty();
data.forEach(function(msg) {
list.append('<li class="list-group-item"><strong>' + msg.topic + '</strong>: ' + msg.payload + ' <span class="text-muted float-end">' + new Date(msg.timestamp).toLocaleString() + '</span></li>');
});
});
}
loadRecent();
// Socket.IO监听新消息
socket.on('new_message', function(msg) {
// 更新最近消息列表
$('#recent-messages').prepend('<li class="list-group-item"><strong>' + msg.topic + '</strong>: ' + msg.payload + ' <span class="text-muted float-end">' + new Date(msg.timestamp).toLocaleString() + '</span></li>');
if ($('#recent-messages li').length > 10) {
$('#recent-messages li:last').remove();
}
});
// 监听传感器数据
socket.on('sensor_data', function(data) {
// 更新图表,此处简化,仅添加一个点
var xData = myChart.getOption().xAxis[0].data;
var yData = myChart.getOption().series[0].data;
var now = new Date().toLocaleTimeString();
xData.push(now);
yData.push(data.temp || 0);
if (xData.length > 20) {
xData.shift();
yData.shift();
}
myChart.setOption({
xAxis: { data: xData },
series: [{ data: yData }]
});
});
});
</script>
{% endblock %}
5.4.3 设备列表页面 templates/devices.html
html
{% extends "base.html" %}
{% block content %}
<h1>设备列表</h1>
<table class="table table-striped">
<thead>
<tr>
<th>ID</th>
<th>Client ID</th>
<th>首次上线</th>
<th>最后在线</th>
</tr>
</thead>
<tbody id="device-table">
</tbody>
</table>
{% endblock %}
{% block scripts %}
<script>
$(document).ready(function() {
function loadDevices() {
$.get('/api/devices', function(data) {
var tbody = $('#device-table');
tbody.empty();
data.forEach(function(dev) {
tbody.append('<tr><td>' + dev.id + '</td><td>' + dev.client_id + '</td><td>' + dev.first_seen + '</td><td>' + dev.last_seen + '</td></tr>');
});
});
}
loadDevices();
// 每10秒刷新设备列表
setInterval(loadDevices, 10000);
});
</script>
{% endblock %}
5.4.4 发布消息页面 templates/publish.html
html
{% extends "base.html" %}
{% block content %}
<h1>发布MQTT消息</h1>
<form id="publishForm">
<div class="mb-3">
<label for="topic" class="form-label">主题</label>
<input type="text" class="form-control" id="topic" required>
</div>
<div class="mb-3">
<label for="payload" class="form-label">消息内容</label>
<textarea class="form-control" id="payload" rows="3" required></textarea>
</div>
<div class="mb-3">
<label for="qos" class="form-label">QoS</label>
<select class="form-select" id="qos">
<option value="0">0</option>
<option value="1">1</option>
<option value="2">2</option>
</select>
</div>
<div class="mb-3 form-check">
<input type="checkbox" class="form-check-input" id="retain">
<label class="form-check-label" for="retain">保留消息</label>
</div>
<button type="submit" class="btn btn-primary">发布</button>
</form>
<div id="result" class="mt-3"></div>
{% endblock %}
{% block scripts %}
<script>
$('#publishForm').submit(function(e) {
e.preventDefault();
var data = {
topic: $('#topic').val(),
payload: $('#payload').val(),
qos: parseInt($('#qos').val()),
retain: $('#retain').is(':checked')
};
$.ajax({
url: '/api/publish',
method: 'POST',
contentType: 'application/json',
data: JSON.stringify(data),
success: function(res) {
$('#result').html('<div class="alert alert-success">消息已发布</div>');
},
error: function() {
$('#result').html('<div class="alert alert-danger">发布失败</div>');
}
});
});
</script>
{% endblock %}
5.5 WebSocket实时通信
我们已经在app.py中集成了Flask-SocketIO,并在前端连接了socket.io。当MQTT客户端收到消息时,通过socketio.emit推送给所有连接的浏览器。前端监听相应事件更新UI。
这种架构确保了低延迟的数据推送,用户体验良好。
6. 系统集成与测试
将上述所有组件整合起来,测试整个系统是否正常工作。
6.1 启动顺序
-
启动MQTT broker(Mosquitto或自实现)。
-
启动Flask应用:
python app.py。 -
打开浏览器访问
http://localhost:5000。 -
运行模拟设备脚本(发布传感器数据)。
-
观察仪表盘实时更新,尝试发布消息。
6.2 测试用例
-
连接测试:浏览器打开仪表盘,WebSocket应自动连接。
-
设备自动注册:模拟设备发布消息到主题
device1/temperature,设备列表应出现device1。 -
实时消息:在消息历史或最近消息列表看到新消息。
-
发布消息:通过发布表单向某主题发布消息,订阅了该主题的客户端(如mosquitto_sub)应收到。
6.3 常见问题与调试
-
WebSocket连接失败:检查Flask-SocketIO是否正常运行,浏览器控制台是否有跨域错误。
-
MQTT连接失败:确认broker地址和端口正确,防火墙允许。
-
数据库更新未提交:检查是否在应用上下文中操作数据库。
7. 部署与优化
7.1 使用Docker容器化部署
为简化部署,我们可以将应用打包成Docker镜像,并使用docker-compose同时启动broker和应用。
7.1.1 Dockerfile
dockerfile
FROM python:3.9-slim WORKDIR /app COPY requirements.txt . RUN pip install --no-cache-dir -r requirements.txt COPY . . EXPOSE 5000 CMD ["python", "app.py"]
7.1.2 docker-compose.yml
yaml
version: '3'
services:
mosquitto:
image: eclipse-mosquitto:2
ports:
- "1883:1883"
- "9001:9001"
volumes:
- ./mosquitto/config:/mosquitto/config
web:
build: .
ports:
- "5000:5000"
environment:
- MQTT_BROKER_HOST=mosquitto
depends_on:
- mosquitto
7.2 性能优化
-
数据库优化:使用索引、定期清理旧数据。
-
WebSocket并发:使用eventlet或gevent作为WSGI服务器,提高并发连接数。
-
MQTT客户端线程:确保paho的loop_start在后台运行,避免阻塞。
-
前端优化:使用虚拟滚动处理大量消息列表,限制图表数据点数量。
7.3 安全加固
-
启用broker认证:配置Mosquitto的密码文件,强制用户名密码。
-
TLS加密:为broker配置SSL证书,客户端使用TLS连接。
-
管理平台登录认证:添加Flask-Login,保护管理界面。
-
CSRF保护:启用Flask-WTF。
8. 总结与扩展
本教程详细介绍了如何使用Python打造一个MQTT服务器(包括使用现成broker和自实现简易版)以及一个功能完整的可视化管理平台。我们从MQTT协议基础入手,逐步构建了后端服务、前端界面,并通过WebSocket实现了实时数据推送。最终提供了一个可运行的物联网监控管理原型。
8.1 项目回顾
-
MQTT协议:理解了报文结构、QoS、主题通配符等核心概念。
-
Broker搭建:掌握了Mosquitto的安装配置,也通过自实现代码加深了协议理解。
-
管理平台:使用Flask + SQLite + Socket.IO实现了设备管理、消息展示、实时图表等功能。
-
集成部署:通过Docker实现了快速部署。
8.2 扩展方向
-
支持更多MQTT特性:如QoS 1/2、遗嘱消息、保留消息的完整处理。
-
增强设备管理:提供设备认证、分组、远程控制(如发送配置命令)。
-
数据分析:基于历史消息做趋势分析、异常检测。
-
告警系统:根据阈值触发告警(邮件、短信)。
-
多协议网关:集成CoAP、HTTP等其他物联网协议。
更多推荐



所有评论(0)