
一、为什么选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 12.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部署bashdocker 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:18083admin/public3.2 认证配置启用内置数据库认证创建用户- 设备用户名设备ID密码一机一密- 应用端用户名app密码强密码3.3 规则引擎配置数据直接写入InfluxDBsqlSELECTpayload.DO as DO,payload.ORP as ORP,payload.MLSS as MLSS,clientid as stationFROM wtps//telemetry动作写入InfluxDB的wtps数据库measurementtelemetry报警转发sqlSELECT * FROM wtps//event WHERE payload.level critical动作触发Webhook通知运维人员3.4 保留消息和会话- 设备状态Topic设为保留消息新订阅者立即获取最新状态- Clean Sessionfalse设备断线期间消息Broker保存上线后投递四、IntBoxIO端MQTT客户端实现Pythonpythonimport paho.mqtt.client as mqttimport json, time, sslclass WtpsMQTTClient:def __init__(self, station_id, broker, port1883):self.station_id station_idself.client mqtt.Client(client_idstation_id,clean_sessionFalse,protocolmqtt.MQTTv311)self.client.username_pw_set(station_id, device_secret_key)self.client.will_set(fwtps/{station_id}/status,json.dumps({status:offline,ts:time.time()}),qos1, retainTrue)self.client.on_connect self.on_connectself.client.on_message self.on_cmdself.client.connect_async(broker, port, keepalive30)def on_connect(self, client, userdata, flags, rc):client.publish(fwtps/{self.station_id}/status,json.dumps({status:online,ts:time.time()}),qos1, retainTrue)client.subscribe(fwtps/{self.station_id}/cmd/, qos1)def on_cmd(self, client, userdata, msg):# 处理下行指令cmd json.loads(msg.payload)# 执行指令...client.publish(fwtps/{self.station_id}/cmd/ack,json.dumps({cmd_id:cmd[id],result:ok}), qos1)def publish_telemetry(self, data):payload json.dumps({ts:time.time(), **data})self.client.publish(fwtps/{self.station_id}/telemetry,payload, qos0)def publish_alarm(self, alarm):self.client.publish(fwtps/{self.station_id}/event,json.dumps(alarm), qos1)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库几行代码就能接入。