一、系统背景

排水泵站远程监控是工业物联网在市政领域的典型应用。本文介绍一套基于Modbus TCP和MQTT协议的泵站远程监控系统,实现多泵站数据采集、远程控制、报警推送和历史存储。系统已在多个县域泵站项目中实际部署。

二、系统架构

现场设备层          边缘控制层           云平台层

┌─────────┐    Modbus TCP    ┌──────────┐   MQTT/TLS  ┌──────────┐

│ PLC/变频器│◄──────────────►│ 边缘控制器 │◄──────────►│ EMQX集群  │

│ 仪表/传感器│   RS485/4-20mA │ (Edge)   │             │ 业务服务  │

└─────────┘                 └──────────┘             └──────────┘

边缘控制器作为协议网关和控制核心,南向通过Modbus TCP/RTU采集PLC和仪表数据,北向通过MQTT over TLS与云端通信。

三、边缘端设计

3.1 数据采集

采用轮询方式采集Modbus设备数据,关键代码逻辑:

class ModbusPoller:

    def __init__(self, host, port=502, interval=1000):

        self.client = ModbusTcpClient(host, port)

        self.interval = interval  # ms

        self.registers = []       # 寄存器配置表

        

    def add_register(self, name, addr, count, dtype='float'):

        self.registers.append({

            'name': name, 'addr': addr,

            'count': count, 'dtype': dtype

        })

    

    def poll(self):

        while True:

            for reg in self.registers:

                result = self.client.read_holding_registers(

                    reg['addr'], reg['count']

                )

                if not result.isError():

                    value = self.decode(result.registers, reg['dtype'])

                    self.on_data(reg['name'], value)

            time.sleep(self.interval / 1000)

采集周期根据数据重要性分级:液位、电流等关键参数1秒采集一次,电能累计量10秒,环境温度60秒。

3.2 数据过滤与压缩

原始数据不全部上传,采用死区压缩策略:

class DeadbandFilter:

    def __init__(self, deadband=0.5):

        self.last_value = None

        self.deadband = deadband

    

    def should_send(self, value):

        if self.last_value is None:

            self.last_value = value

            return True

        if abs(value - self.last_value) >= self.deadband:

            self.last_value = value

            return True

        return False

仅当数据变化超过死区阈值时才上报,正常运行时数据上传量降低约60%。

3.3 MQTT通信设计

Topic设计:

数据上报:/pump/{station_id}/data

状态上报:/pump/{station_id}/status

指令下发:/pump/{station_id}/cmd

指令响应:/pump/{station_id}/cmd_ack

数据上报Payload(JSON):

{

  "ts": 1726032000000,

  "seq": 10234,

  "data": {

    "level": 3.25,

    "pump1_current": 45.2,

    "pump1_status": 1,

    "pump2_status": 0,

    "flow_rate": 120.5

  }

}

关键设计:

QoS 1用于数据上报,QoS 2用于控制指令

每条消息携带设备时间戳和递增序列号,云端去重

遗嘱消息(LWT)通知异常离线

TLS 1.2加密,设备证书认证

3.4 断网缓存与续传

使用SQLite本地缓存断网期间数据:

class OfflineCache:

    def __init__(self, db_path='cache.db'):

        self.conn = sqlite3.connect(db_path)

        self.conn.execute('''CREATE TABLE IF NOT EXISTS data

            (id INTEGER PRIMARY KEY AUTOINCREMENT,

             topic TEXT, payload TEXT, ts INTEGER)''')

    

    def save(self, topic, payload):

        self.conn.execute(

            'INSERT INTO data (topic, payload, ts) VALUES (?,?,?)',

            (topic, json.dumps(payload), payload['ts'])

        )

        self.conn.commit()

    

    def flush(self, mqtt_client):

        rows = self.conn.execute(

            'SELECT id, topic, payload FROM data ORDER BY id LIMIT 500'

        ).fetchall()

        for row_id, topic, payload in rows:

            info = mqtt_client.publish(topic, payload, qos=1)

            if info.rc == 0:

                self.conn.execute('DELETE FROM data WHERE id=?', (row_id,))

        self.conn.commit()

网络恢复后按时间顺序批量续传,确保数据零丢失。

四、云端设计

云端采用EMQX作为MQTT Broker,后端服务订阅所有数据Topic,写入TDengine时序数据库。控制指令通过REST API下发,经权限校验后发布到指令Topic。

报警规则引擎实时分析数据流,液位超限、设备故障等事件触发报警,通过APP推送、短信和电话分级通知。

五、部署效果

系统在某县15座泵站部署后:

数据采集延迟<2秒,控制指令响应<1秒

断网72小时本地控制正常,数据完整续传

云端单集群支持5000+泵站并发接入

运维人员从45人减至8人

六、总结

本文介绍的泵站远程监控系统采用"边缘采集+MQTT通信+云端处理"的经典IIoT架构,核心设计要点在于:边缘端本地控制保证可靠性、死区压缩降低带宽、断网缓存保证数据完整、MQTT QoS机制保证消息可达。该架构同样适用于其他分布式工业监控场景。

更多推荐