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/temperaturesensor/2/temperature,但不匹配sensor/1/floor/temperature

  • 多级通配符#:匹配任意层级,必须放在最后。例如sensor/#可匹配sensor/1/temperaturesensor/1sensor/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模块或直接使用socketthreading

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 启动顺序

  1. 启动MQTT broker(Mosquitto或自实现)。

  2. 启动Flask应用:python app.py

  3. 打开浏览器访问 http://localhost:5000

  4. 运行模拟设备脚本(发布传感器数据)。

  5. 观察仪表盘实时更新,尝试发布消息。

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等其他物联网协议。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐