
前面我们把采集、布线、控制逻辑、TCP 保活、CKafka 异步解耦、InfluxDB 降采样都讲透了。
这一篇回到链路中间——那台不起眼、但决定整个系统能不能"活"的协议转换网关。
先说一个残酷的事实:
你买的以太网温湿度传感器,99% 只支持 Modbus TCP 从站。 它不会主动连云、不会发 MQTT、不会注册到 IoT 平台。 它能做的只有一件事:等着被读。
所以"数据上云"这件事,不是传感器自己完成的,而是协议转换网关在中间"翻译"的结果。
设备侧 | 网关 | 云端 |
|---|---|---|
Modbus TCP(从站,被动) | 协议转换 | MQTT(主动发布) |
寄存器地址 40001–40014 | 映射 | JSON 物模型属性 |
大端/小端、INT16/FLOAT | 解析 | 类型化字段 |
无心跳、无重连 | 保活 | Keepalive + 会话 |
无加密、无认证 | 安全 | TLS + SASL/证书 |
传感器说的是"寄存器语言",云端说的是"消息语言"。网关就是翻译官。
形态 | 代表 | 优点 | 缺点 |
|---|---|---|---|
硬件网关(工业 ARM 盒子) | 纵横、有人、摩莎 | 工业级、宽温、断电离线缓存 | 贵(¥800–3000/台)、算力有限 |
软件网关(x86/容器) | 自建 Python/Go 服务 | 灵活、可扩展、成本低 | 需要自己保证高可用 |
一体机内置网关 | 高端净化恒湿一体机 | 省设备、集成度高 | 绑定厂商、不可独立运维 |
中小型项目(20–100 个监测点)最常用的是:一台 x86 工控机/迷你 PC + Docker 容器跑软件网关。
┌─────────────────────────────────────────────────────────────┐
│ 协议转换网关 │
│ │
│ ┌──────────────┐ ┌──────────────┐ ┌──────────────┐ │
│ │ Modbus TCP │ │ 数据解析与 │ │ MQTT 客户端 │ │
│ │ 采集引擎 │───▶│ 点表映射 │───▶│ 发布/订阅 │ │
│ │ │ │ │ │ │ │
│ │ · 轮询调度 │ │ · 寄存器解码 │ │ · TLS 加密 │ │
│ │ · 超时重试 │ │ · 字节序转换 │ │ · 自动重连 │ │
│ │ · 并发控制 │ │ · 质量位标记 │ │ · QoS 1/2 │ │
│ │ · 断线检测 │ │ · 变化过滤 │ │ · 遗嘱消息 │ │
│ └──────────────┘ └──────────────┘ └──────────────┘ │
│ │ │ │ │
│ ┌──────▼──────┐ ┌────────▼────────┐ ┌──────▼──────┐ │
│ │ 本地缓存 │ │ 连接管理器 │ │ 配置热加载 │ │
│ │ · SQLite │ │ · TCP Keepalive│ │ · 点表更新 │ │
│ │ · WAL │ │ · MQTT 会话 │ │ · 参数调整 │ │
│ │ · 补传队列 │ │ · 指数退避 │ │ · 不重启 │ │
│ └─────────────┘ └─────────────────┘ └─────────────┘ │
└─────────────────────────────────────────────────────────────┘不能简单地"for 循环挨个读"——20 台设备串行轮询,每台读 10 个寄存器,RTU 超时设 3 秒,一轮下来就是 60 秒。数据新鲜度完全不够。
正确做法:并发 + 批量
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]参数 | 推荐值 | 说明 |
|---|---|---|
轮询间隔 | 5s | 温湿度变化慢,5s 足够 |
并发数 | 10–20 | 同时轮询的设备数 |
单次读寄存器数 | ≤ 125 | Modbus TCP 协议限制(单帧 256 字节) |
超时 | 2s | 超过 2s 视为设备不可达 |
重试次数 | 2 | 超时后重试 2 次 |
重试间隔 | 500ms | 指数退避:500ms → 1s |
不是每次轮询结果都要上报:
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这样:
寄存器 | 名称 | 类型 | 缩放 | 单位 |
|---|---|---|---|---|
40001 | 温度 | INT16 | ×0.1 | ℃ |
40002 | 湿度 | INT16 | ×0.1 | %RH |
40003 | PM2.5 | UINT16 | ×1 | μg/m³ |
40004–40005 | 累计运行时间 | UINT32 | ×1 | 小时 |
40006 | 设备状态字 | UINT16 | 位域 | — |
40007–40014 | 保留 | — | — | — |
{
"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"
}
]
}{
"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"
}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 clientTopic | 方向 | 用途 |
|---|---|---|
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 | 下行 | 配置热更新 |
消息类型 | QoS | 说明 |
|---|---|---|
遥测数据 | 1(至少一次) | 允许重复,不允许丢失 |
设备事件 | 1 | 同上 |
控制指令 | 2(恰好一次) | 不允许重复执行 |
配置更新 | 1 | 允许重试 |
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()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 # 发送失败,停止补传,等下次重试关键:补传消息的时间戳必须是 sample_time(传感器采样时间),不是 report_time(网关发送时间)。
{
"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,就知道这是补传数据,不会触发实时告警(但可以触发"离线期间异常"分析)。
# 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对于关键机房,可以部署两台网关做主备:
主网关(192.168.10.100)──┐
├── 虚拟 IP(192.168.10.10,VRRP/Keepalived)
备网关(192.168.10.101)──┘不需要重启网关就能更新配置:
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 删除。