MQTT在污水处理物联网中的应用:从协议分析到EMQX部署
一、为什么选MQTT
污水处理物联网场景特点:
- 设备数量多(一个区域几十上百个站)
- 网络环境复杂(4G为主,可能不稳定)
- 数据量不大但要求实时(秒级)
- 需要主动推送(报警)
- 设备计算资源有限(边缘控制器)
MQTT vs HTTP:
| 特性 | MQTT | HTTP |
|------|------|------|
| 模式 | 发布/订阅 | 请求/响应 |
| 头部开销 | 最小2字节 | 数百字节 |
| 实时推送 | 原生支持 | 需要轮询 |
| 连接保持 | 长连接+心跳 | 短连接 |
| QoS | 0/1/2三级 | 无 |
| 断线重连 | 原生支持 | 需自己实现 |
| 适用 | 物联网 | Web API |
二、MQTT协议核心概念
2.1 报文结构
```
固定头(1字节+) + 可变头 + 有效载荷
```
固定头第一个字节:报文类型(4bit)+标志(4bit)
- CONNECT(1), CONNACK(2), PUBLISH(3), PUBACK(4), SUBSCRIBE(8), SUBACK(9), PINGREQ(12), PINGRESP(13), DISCONNECT(14)
2.2 QoS级别
- QoS 0:最多一次,发完不管
- QoS 1:至少一次,PUBACK确认,可能重复
- QoS 2:恰好一次,四次握手,开销大
污水处理场景推荐:
- 实时数据:QoS 0(丢一两个点无所谓)
- 报警和控制指令:QoS 1(必须收到)
- 历史数据补传:QoS 1
2.3 Topic设计
```
wtps/{station_id}/telemetry # 遥测数据
wtps/{station_id}/event # 事件/报警
wtps/{station_id}/status # 设备状态(online/offline via LWT)
wtps/{station_id}/cmd/+ # 下行指令
wtps/{station_id}/cmd/ack # 指令确认
```
通配符:+单层,#多层
2.4 遗嘱消息(LWT)
设备CONNECT时指定遗嘱Topic和内容,异常断线时Broker自动发布,平台可实时感知设备离线。
三、EMQX部署
3.1 Docker部署
```bash
docker run -d --name emqx -p 1883:1883 -p 8083:8083 -p 8883:8883 -p 18083:18083 -v /opt/emqx/data:/opt/emqx/data emqx/emqx:5.0
```
默认控制台http://ip:18083,admin/public
3.2 认证配置
启用内置数据库认证,创建用户:
- 设备用户名:设备ID,密码:一机一密
- 应用端用户名:app,密码:强密码
3.3 规则引擎配置
数据直接写入InfluxDB:
```sql
SELECT
payload.DO as DO,
payload.ORP as ORP,
payload.MLSS as MLSS,
clientid as station
FROM "wtps/+/telemetry"
```
动作:写入InfluxDB的wtps数据库,measurement=telemetry
报警转发:
```sql
SELECT * FROM "wtps/+/event" WHERE payload.level = 'critical'
```
动作:触发Webhook通知运维人员
3.4 保留消息和会话
- 设备状态Topic设为保留消息,新订阅者立即获取最新状态
- Clean Session=false,设备断线期间消息Broker保存,上线后投递
四、IntBoxIO端MQTT客户端实现(Python)
```python
import paho.mqtt.client as mqtt
import json, time, ssl
class WtpsMQTTClient:
def __init__(self, station_id, broker, port=1883):
self.station_id = station_id
self.client = mqtt.Client(
client_id=station_id,
clean_session=False,
protocol=mqtt.MQTTv311
)
self.client.username_pw_set(station_id, "device_secret_key")
self.client.will_set(
f"wtps/{station_id}/status",
json.dumps({"status":"offline","ts":time.time()}),
qos=1, retain=True
)
self.client.on_connect = self.on_connect
self.client.on_message = self.on_cmd
self.client.connect_async(broker, port, keepalive=30)
def on_connect(self, client, userdata, flags, rc):
client.publish(f"wtps/{self.station_id}/status",
json.dumps({"status":"online","ts":time.time()}),
qos=1, retain=True)
client.subscribe(f"wtps/{self.station_id}/cmd/+", qos=1)
def on_cmd(self, client, userdata, msg):
# 处理下行指令
cmd = json.loads(msg.payload)
# 执行指令...
client.publish(f"wtps/{self.station_id}/cmd/ack",
json.dumps({"cmd_id":cmd["id"],"result":"ok"}), qos=1)
def publish_telemetry(self, data):
payload = json.dumps({"ts":time.time(), **data})
self.client.publish(f"wtps/{self.station_id}/telemetry",
payload, qos=0)
def publish_alarm(self, alarm):
self.client.publish(f"wtps/{self.station_id}/event",
json.dumps(alarm), qos=1)
def start(self):
self.client.loop_start()
```
五、性能优化
- 单EMQX节点支持5万+连接,集群支持百万级
- 数据打包:5秒上报一次,每次打包多个数据点,减少报文数
- Payload用CBOR代替JSON可减少40%流量(但调试不便)
- TLS开销大,4G场景可用PSK或VPN代替全量TLS
六、安全
- 一机一密认证
- Topic级ACL:设备只能发布自己的Topic,不能订阅其他设备
- 外部访问用WSS或8883 TLS
- 定期轮换密码
七、总结
MQTT是污水处理物联网的最佳通讯协议。EMQX开源版即可满足中小规模部署,规则引擎可以直接把数据写入时序数据库,省去自己写订阅服务。IntBoxIO通过paho-mqtt库几行代码就能接入。
