首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >从 RJ45 到 MQTT Broker:以太网温湿度传感器数据透传至云端的协议转换网关

从 RJ45 到 MQTT Broker:以太网温湿度传感器数据透传至云端的协议转换网关

原创
作者头像
BJ盛世宏博小程
发布2026-09-21 17:13:56
发布2026-09-21 17:13:56
500
举报

从 RJ45 到 MQTT Broker:以太网温湿度传感器数据透传至云端的协议转换网关

前面我们把采集、布线、控制逻辑、TCP 保活、CKafka 异步解耦、InfluxDB 降采样都讲透了。

这一篇回到链路中间——那台不起眼、但决定整个系统能不能"活"的协议转换网关

先说一个残酷的事实:

你买的以太网温湿度传感器,99% 只支持 Modbus TCP 从站。 它不会主动连云、不会发 MQTT、不会注册到 IoT 平台。 它能做的只有一件事:等着被读。

所以"数据上云"这件事,不是传感器自己完成的,而是协议转换网关在中间"翻译"的结果。


一、为什么需要协议转换网关

1. 协议鸿沟

设备侧

网关

云端

Modbus TCP(从站,被动)

协议转换

MQTT(主动发布)

寄存器地址 40001–40014

映射

JSON 物模型属性

大端/小端、INT16/FLOAT

解析

类型化字段

无心跳、无重连

保活

Keepalive + 会话

无加密、无认证

安全

TLS + SASL/证书

传感器说的是"寄存器语言",云端说的是"消息语言"。网关就是翻译官

2. 网关的三种形态

形态

代表

优点

缺点

硬件网关(工业 ARM 盒子)

纵横、有人、摩莎

工业级、宽温、断电离线缓存

贵(¥800–3000/台)、算力有限

软件网关(x86/容器)

自建 Python/Go 服务

灵活、可扩展、成本低

需要自己保证高可用

一体机内置网关​

高端净化恒湿一体机

省设备、集成度高

绑定厂商、不可独立运维

中小型项目(20–100 个监测点)最常用的是:一台 x86 工控机/迷你 PC + Docker 容器跑软件网关


二、网关核心架构

代码语言:javascript
复制
┌─────────────────────────────────────────────────────────────┐
│                    协议转换网关                               │
│                                                             │
│  ┌──────────────┐    ┌──────────────┐    ┌──────────────┐ │
│  │  Modbus TCP  │    │  数据解析与   │    │   MQTT 客户端 │ │
│  │  采集引擎     │───▶│  点表映射     │───▶│  发布/订阅    │ │
│  │              │    │              │    │              │ │
│  │  · 轮询调度  │    │  · 寄存器解码 │    │  · TLS 加密   │ │
│  │  · 超时重试  │    │  · 字节序转换 │    │  · 自动重连   │ │
│  │  · 并发控制  │    │  · 质量位标记 │    │  · QoS 1/2   │ │
│  │  · 断线检测  │    │  · 变化过滤   │    │  · 遗嘱消息   │ │
│  └──────────────┘    └──────────────┘    └──────────────┘ │
│         │                    │                    │         │
│  ┌──────▼──────┐   ┌────────▼────────┐  ┌──────▼──────┐  │
│  │  本地缓存    │   │  连接管理器      │  │  配置热加载  │  │
│  │  · SQLite   │   │  · TCP Keepalive│  │  · 点表更新  │  │
│  │  · WAL      │   │  · MQTT 会话    │  │  · 参数调整  │  │
│  │  · 补传队列 │   │  · 指数退避     │  │  · 不重启    │  │
│  └─────────────┘   └─────────────────┘  └─────────────┘  │
└─────────────────────────────────────────────────────────────┘

三、Modbus TCP 采集引擎设计

1. 轮询调度策略

不能简单地"for 循环挨个读"——20 台设备串行轮询,每台读 10 个寄存器,RTU 超时设 3 秒,一轮下来就是 60 秒。数据新鲜度完全不够。

正确做法:并发 + 批量

代码语言:javascript
复制
import asyncio
import struct
from pymodbus.client import AsyncModbusTcpClient

class ModbusPoller:
    def __init__(self, devices, poll_interval=5):
        self.devices = devices  # [{ip, port, slave_id, registers: [{addr, count, type}]}]
        self.poll_interval = poll_interval
        self.clients = {}

    async def get_client(self, ip, port):
        key = f"{ip}:{port}"
        if key not in self.clients or not self.clients[key].connected:
            client = AsyncModbusTcpClient(ip, port=port)
            await client.connect()
            self.clients[key] = client
        return self.clients[key]

    async def poll_device(self, device):
        try:
            client = await self.get_client(device['ip'], device['port'])
            results = {}

            for reg_block in device['registers']:
                # 批量读取连续寄存器(一次请求读多个)
                resp = await client.read_holding_registers(
                    address=reg_block['addr'],
                    count=reg_block['count'],
                    slave=device['slave_id']
                )
                if resp.isError():
                    results[reg_block['name']] = {'value': None, 'quality': 2}  # bad
                else:
                    # 字节序转换
                    values = self.decode_registers(resp.registers, reg_block['type'])
                    results[reg_block['name']] = {'value': values, 'quality': 0}  # good

            return {
                'device_id': device['id'],
                'timestamp': asyncio.get_event_loop().time(),
                'data': results
            }
        except Exception as e:
            return {
                'device_id': device['id'],
                'timestamp': asyncio.get_event_loop().time(),
                'error': str(e),
                'quality': 2
            }

    async def poll_all(self):
        """并发轮询所有设备"""
        tasks = [self.poll_device(d) for d in self.devices]
        return await asyncio.gather(*tasks, return_exceptions=True)

    def decode_registers(self, registers, data_type):
        """根据数据类型解码寄存器"""
        if data_type == 'INT16':
            return registers[0]
        elif data_type == 'FLOAT32_BE':
            # 大端浮点:两个 16 位寄存器拼成 32 位浮点
            raw = (registers[0] << 16) | registers[1]
            return struct.unpack('>f', struct.pack('>I', raw))[0]
        elif data_type == 'FLOAT32_LE':
            raw = (registers[1] << 16) | registers[0]
            return struct.unpack('<f', struct.pack('<I', raw))[0]

2. 关键参数

参数

推荐值

说明

轮询间隔

5s

温湿度变化慢,5s 足够

并发数

10–20

同时轮询的设备数

单次读寄存器数

≤ 125

Modbus TCP 协议限制(单帧 256 字节)

超时

2s

超过 2s 视为设备不可达

重试次数

2

超时后重试 2 次

重试间隔

500ms

指数退避:500ms → 1s

3. 变化过滤(Deadband)

不是每次轮询结果都要上报:

代码语言:javascript
复制
class ChangeFilter:
    def __init__(self, deadband_temp=0.2, deadband_humid=0.5):
        self.last_values = {}
        self.deadband_temp = deadband_temp
        self.deadband_humid = deadband_humid

    def should_report(self, device_id, temp, humid):
        key = device_id
        if key not in self.last_values:
            self.last_values[key] = (temp, humid)
            return True

        last_temp, last_humid = self.last_values[key]
        temp_changed = abs(temp - last_temp) >= self.deadband_temp
        humid_changed = abs(humid - last_humid) >= self.deadband_humid

        if temp_changed or humid_changed:
            self.last_values[key] = (temp, humid)
            return True
        return False

这样:

  • 稳态时(温度波动 < 0.2℃ 且湿度波动 < 0.5%RH),不上报
  • 异常时(温湿度快速变化),立即上报
  • 网络流量和云端写入量降低 60–80%

四、点表映射:从寄存器到物模型

1. 设备点表(传感器侧)

寄存器

名称

类型

缩放

单位

40001

温度

INT16

×0.1

40002

湿度

INT16

×0.1

%RH

40003

PM2.5

UINT16

×1

μg/m³

40004–40005

累计运行时间

UINT32

×1

小时

40006

设备状态字

UINT16

位域

40007–40014

保留

2. 映射配置(JSON)

代码语言:javascript
复制
{
  "device_type": "ethernet_th_sensor_v2",
  "poll_interval": 5,
  "registers": [
    {
      "name": "temperature",
      "address": 40001,
      "count": 1,
      "type": "INT16",
      "scale": 0.1,
      "unit": "°C",
      "target_field": "temp"
    },
    {
      "name": "humidity",
      "address": 40002,
      "count": 1,
      "type": "INT16",
      "scale": 0.1,
      "unit": "%RH",
      "target_field": "rh"
    },
    {
      "name": "pm25",
      "address": 40003,
      "count": 1,
      "type": "UINT16",
      "scale": 1,
      "unit": "μg/m³",
      "target_field": "pm25"
    },
    {
      "name": "uptime_hours",
      "address": 40004,
      "count": 2,
      "type": "UINT32_BE",
      "scale": 1,
      "unit": "h",
      "target_field": "uptime"
    },
    {
      "name": "status_word",
      "address": 40006,
      "count": 1,
      "type": "UINT16",
      "bitmap": {
        "0": "sensor_error",
        "1": "eeprom_error",
        "2": "communication_error",
        "3": "calibration_needed"
      },
      "target_field": "fault"
    }
  ]
}

3. 转换后的 MQTT 消息

代码语言:javascript
复制
{
  "message_id": "msg-a3-01-20260921-143052-001",
  "device_id": "TH-A3-01",
  "site_id": "site-tianjin-001",
  "room_id": "archive-3f-a3",
  "sample_time": "2026-09-21T14:30:52+08:00",
  "report_time": "2026-09-21T14:30:52.156+08:00",
  "schema_version": "1.0",
  "metrics": {
    "temp": 22.4,
    "rh": 55.8,
    "pm25": 12,
    "uptime": 8760
  },
  "fault": {
    "sensor_error": false,
    "eeprom_error": false,
    "communication_error": false,
    "calibration_needed": true
  },
  "quality": "good"
}

五、MQTT 客户端设计

1. 连接参数

代码语言:javascript
复制
import paho.mqtt.client as mqtt
import ssl

def create_mqtt_client(config):
    client = mqtt.Client(
        client_id=config['client_id'],
        clean_session=False,       # 持久会话,断线重连后不丢消息
        protocol=mqtt.MQTTv311
    )

    # 认证
    if config.get('username'):
        client.username_pw_set(config['username'], config['password'])
    elif config.get('certfile'):
        client.tls_set(
            ca_certs=config['ca_certs'],
            certfile=config['certfile'],
            keyfile=config['keyfile'],
            cert_reqs=ssl.CERT_REQUIRED,
            tls_version=ssl.PROTOCOL_TLSv1_2
        )

    # 遗嘱消息(LWT)—— 网关异常断开时,云端知道
    client.will_set(
        topic=f"gateway/{config['gateway_id']}/status",
        payload='{"status":"offline","reason":"unexpected_disconnect"}',
        qos=1,
        retain=True
    )

    # 连接参数
    client.connect_async(
        host=config['broker_host'],
        port=config.get('port', 8883),
        keepalive=60
    )

    return client

2. 发布 Topic 设计

Topic

方向

用途

env/telemetry/{site_id}/{room_id}/{device_id}

上行

遥测数据

env/event/{site_id}/{device_id}

上行

设备事件(离线、故障)

env/control/{site_id}/{device_id}

下行

控制指令(如果有写寄存器能力)

gateway/{gateway_id}/status

上行

网关在线状态

gateway/{gateway_id}/config

下行

配置热更新

3. QoS 选择

消息类型

QoS

说明

遥测数据

1(至少一次)

允许重复,不允许丢失

设备事件

1

同上

控制指令

2(恰好一次)

不允许重复执行

配置更新

1

允许重试


六、本地缓存与断网补传

1. SQLite WAL 缓存

代码语言:javascript
复制
import sqlite3
import json

class LocalCache:
    def __init__(self, db_path):
        self.conn = sqlite3.connect(db_path, check_same_thread=False)
        self.conn.execute('''
            CREATE TABLE IF NOT EXISTS message_queue (
                id INTEGER PRIMARY KEY AUTOINCREMENT,
                device_id TEXT,
                topic TEXT,
                payload TEXT,
                sample_time TEXT,
                created_at REAL,
                sent INTEGER DEFAULT 0,
                retry_count INTEGER DEFAULT 0
            )
        ''')
        self.conn.commit()

    def store(self, device_id, topic, payload, sample_time):
        self.conn.execute(
            'INSERT INTO message_queue (device_id, topic, payload, sample_time, created_at) VALUES (?, ?, ?, ?, ?)',
            (device_id, topic, json.dumps(payload), sample_time, time.time())
        )
        self.conn.commit()

    def get_unsent(self, limit=100):
        cursor = self.conn.execute(
            'SELECT id, topic, payload FROM message_queue WHERE sent = 0 ORDER BY sample_time ASC LIMIT ?',
            (limit,)
        )
        return [(row[0], row[1], json.loads(row[2])) for row in cursor.fetchall()]

    def mark_sent(self, msg_id):
        self.conn.execute('UPDATE message_queue SET sent = 1 WHERE id = ?', (msg_id,))
        self.conn.commit()

    def cleanup_sent(self, older_than_days=7):
        """清理 7 天前已发送的消息"""
        cutoff = time.time() - (older_than_days * 86400)
        self.conn.execute('DELETE FROM message_queue WHERE sent = 1 AND created_at < ?', (cutoff,))
        self.conn.commit()

2. 补传逻辑

代码语言:javascript
复制
class BackfillManager:
    def __init__(self, cache, mqtt_client):
        self.cache = cache
        self.mqtt_client = mqtt_client
        self.backfill_batch_size = 50

    def try_backfill(self):
        """尝试补传缓存的消息"""
        if not self.mqtt_client.is_connected():
            return  # 还没连上,不补传

        unsent = self.cache.get_unsent(self.backfill_batch_size)
        if not unsent:
            return

        for msg_id, topic, payload in unsent:
            try:
                result = self.mqtt_client.publish(
                    topic,
                    json.dumps(payload),
                    qos=1,
                    retain=False
                )
                if result.rc == mqtt.MQTT_ERR_SUCCESS:
                    self.cache.mark_sent(msg_id)
            except Exception:
                break  # 发送失败,停止补传,等下次重试

3. 补传时的"时间语义"

关键:补传消息的时间戳必须是 sample_time(传感器采样时间),不是 report_time(网关发送时间)。

代码语言:javascript
复制
{
  "message_id": "msg-a3-01-20260921-143052-001",
  "sample_time": "2026-09-21T14:30:52+08:00",
  "report_time": "2026-09-21T18:45:12+08:00",
  "is_backfill": true,
  "metrics": { "temp": 22.4, "rh": 55.8 }
}

云端消费者看到 is_backfill: true,就知道这是补传数据,不会触发实时告警(但可以触发"离线期间异常"分析)。


七、高可用设计

1. 网关自身保活

代码语言:javascript
复制
# systemd 服务,自动重启
[Unit]
Description=Env Gateway Protocol Converter
After=network.target

[Service]
Type=simple
User=envgw
WorkingDirectory=/opt/env-gateway
ExecStart=/opt/env-gateway/venv/bin/python main.py --config /etc/env-gateway/config.yaml
Restart=always
RestartSec=10
StandardOutput=journal
StandardError=journal

[Install]
WantedBy=multi-user.target

2. 双机热备(可选)

对于关键机房,可以部署两台网关做主备:

代码语言:javascript
复制
主网关(192.168.10.100)──┐
                           ├── 虚拟 IP(192.168.10.10,VRRP/Keepalived)
备网关(192.168.10.101)──┘
  • 主网关正常时,虚拟 IP 绑定在主网关
  • 主网关故障,Keepalived 将虚拟 IP 漂移到备网关
  • 备网关平时也轮询设备(只读不写),主备切换时无缝接管

八、配置热加载

不需要重启网关就能更新配置:

代码语言:javascript
复制
import signal

class ConfigManager:
    def __init__(self, config_path):
        self.config_path = config_path
        self.config = self.load()
        signal.signal(signal.SIGHUP, self.reload)  # kill -HUP pid

    def reload(self, signum, frame):
        logger.info("Reloading configuration...")
        new_config = self.load()
        if self.validate(new_config):
            self.config = new_config
            self.apply_changes()
            logger.info("Configuration reloaded successfully")
        else:
            logger.error("Invalid configuration, keeping old one")

    def load(self):
        with open(self.config_path) as f:
            return yaml.safe_load(f)

九、常见返工点

问题

后果

正确做法

轮询串行执行

20 台设备一轮 60 秒

并发 + 批量读寄存器

字节序搞反

温度显示 1717℃(0x06B5 当小端)

确认设备手册的字节序

不标记质量位

传感器断线,平台显示上一次的值

超时/错误时 quality=bad

断网不缓存

网络中断 4 小时,数据全丢

SQLite WAL 本地缓存

补传用当前时间

时间轴错乱,降采样结果错误

必须用 sample_time

MQTT clean_session=True

断线重连后,离线期间的指令全丢

持久会话

不设置遗嘱消息

网关崩溃,云端不知道

LWT 通知离线

网关单点

网关故障,所有设备失联

systemd 自启 + 可选双机热备

点表硬编码

换一批传感器就要改代码

JSON 配置文件,热加载

QoS 全用 0

网络抖动时消息静默丢失

遥测 QoS 1,控制 QoS 2


十、一句话总结

协议转换网关不是"中间人",而是整个动环监控系统的翻译官 + 守门人 + 保险箱——它把 Modbus TCP 的寄存器语言翻译成 MQTT 的消息语言,把不可靠的以太网变成可靠的消息通道,把断网期间的每一帧数据都锁进本地缓存等待补传。

当验收那天,你拔掉网关的网线 4 小时,插回去后云端自动补上所有缺失的数据点——曲线连续不断、审计链完整无缺——这才是网关的价值,而不是"能 ping 通"的价值

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

目录
  • 从 RJ45 到 MQTT Broker:以太网温湿度传感器数据透传至云端的协议转换网关
    • 一、为什么需要协议转换网关
      • 1. 协议鸿沟
      • 2. 网关的三种形态
    • 二、网关核心架构
    • 三、Modbus TCP 采集引擎设计
      • 1. 轮询调度策略
      • 2. 关键参数
      • 3. 变化过滤(Deadband)
    • 四、点表映射:从寄存器到物模型
      • 1. 设备点表(传感器侧)
      • 2. 映射配置(JSON)
      • 3. 转换后的 MQTT 消息
    • 五、MQTT 客户端设计
      • 1. 连接参数
      • 2. 发布 Topic 设计
      • 3. QoS 选择
    • 六、本地缓存与断网补传
      • 1. SQLite WAL 缓存
      • 2. 补传逻辑
      • 3. 补传时的"时间语义"
    • 七、高可用设计
      • 1. 网关自身保活
      • 2. 双机热备(可选)
    • 八、配置热加载
    • 九、常见返工点
    • 十、一句话总结
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档