

近期在调试一套工业环境监控系统时,遇到了关于Modbus连接保持的问题。现场部署了一批网口温湿度变送器,PoE取电、Modbus TCP上云,同时利旧接入存量RS485探头,并开启了SNMP和UDP Trap服务。档案馆项目里往往不止一个库房——特藏库、纸质档案库、音像档案库、过渡库,每个库房的温湿度要求不同,物理位置分散在不同楼层甚至不同建筑。多库房场景下,网络需要分区隔离、数据需要跨子网汇聚、告警需要按区域分级——这篇从子网划分、路由策略、汇聚架构到数据去重与断点续传,完整记录多库房分区监测的工程落地方案。
库房类型 | 温湿度要求 | 监测密度 | 告警阈值 | 合规标准 |
|---|---|---|---|---|
特藏库 | 温度16~20℃,RH 45%~55% | 每50㎡一个 | ±1℃ / ±3%RH | DA/T 42-2022 |
纸质档案库 | 温度14~24℃,RH 45%~60% | 每100㎡一个 | ±2℃ / ±5%RH | DA/T 42-2022 |
音像档案库 | 温度18~22℃,RH 40%~50% | 每80㎡一个 | ±1.5℃ / ±4%RH | DA/T 42-2022 |
过渡/接收库 | 温度≤30℃,RH ≤65% | 每150㎡一个 | ±3℃ / ±8%RH | 内部规范 |
档案馆建筑平面示意:
┌────────────────────────────────────────────────────────────┐
│ 3F: 特藏库(80㎡) │ 纸质库A(300㎡) │ 纸质库B(250㎡) │
│ ┌───┐ │ ┌───┐ ┌───┐ │ ┌───┐ ┌───┐ │
│ │S1 │ │ │S2│ │S3│ │ │ │S4│ │S5│ │
│ └───┘ │ └───┘ └───┘ │ └───┘ └───┘ │
├────────────────────────────────────────────────────────────┤
│ 2F: 音像库(120㎡)│ 数字化车间(200㎡)│ 办公区 │
│ ┌───┐ ┌───┐ │ ┌───┐ ┌───┐ │ │
│ │S6 │ │S7 │ │ │S8│ │S9│ │ │ │
│ └───┘ └───┘ │ └───┘ └───┘ │ │
├────────────────────────────────────────────────────────────┤
│ 1F: 过渡库(400㎡)│ 接收整理区(150㎡)│ 监控中心 │
│ ┌───┐┌───┐ │ ┌───┐ │ ┌───────────────┐ │
│ │S10││S11│ │ │S12│ │ │ SCADA/InfluxDB│ │
│ └───┘└───┘ │ └───┘ │ └───────────────┘ │
└────────────────────────────────────────────────────────────┘挑战 | 描述 | 影响 |
|---|---|---|
跨子网通信 | 各库房在不同VLAN/子网 | 监控中心需跨网段汇聚数据 |
广播域隔离 | 库房间不希望互相广播 | IGMP/Modbus广播需控制 |
带宽竞争 | 数字化车间大流量影响传感器 | 数据延迟、丢包 |
单点故障 | 汇聚交换机故障导致全馆断监 | 合规风险 |
时间同步 | 各库房NTP源不同步 | 数据时间戳不一致 |
安全隔离 | 办公网不应访问传感器 | 未授权访问风险 |
档案馆整体网络:10.10.0.0/16
子网划分:
VLAN 10 - 管理网段 10.10.10.0/24 交换机/网关管理
VLAN 20 - 特藏库 10.10.20.0/24 传感器×8
VLAN 30 - 纸质库A 10.10.30.0/24 传感器×16
VLAN 40 - 纸质库B 10.10.40.0/24 传感器×12
VLAN 50 - 音像库 10.10.50.0/24 传感器×10
VLAN 60 - 数字化车间 10.10.60.0/24 传感器×8 + 其他设备
VLAN 70 - 过渡库 10.10.70.0/24 传感器×14
VLAN 80 - 接收整理区 10.10.80.0/24 传感器×6
VLAN 90 - 监控中心 10.10.90.0/24 SCADA/InfluxDB
VLAN 99 - 设备发现 10.10.99.0/24 预留/DHCP
传感器IP分配规则:
10.10.{vlan_id}.{sensor_id}
例:特藏库第3个传感器 → 10.10.20.3┌─────────────────────────────────────────────────────────────────┐
│ 核心三层交换机 │
│ ┌───────────────────────────────────────────────────────────┐ │
│ │ VLAN间路由表 │ │
│ │ 10.10.20.0/24 → VLAN 20 (特藏库) │ │
│ │ 10.10.30.0/24 → VLAN 30 (纸质库A) │ │
│ │ ... │ │
│ │ 10.10.90.0/24 → VLAN 90 (监控中心) │ │
│ └───────────────────────────────────────────────────────────┘ │
│ │
│ 访问控制列表 (ACL): │
│ permit tcp 10.10.90.0/24 10.10.20.0/24 eq 502 # 监控中心→Modbus│
│ permit udp 10.10.20.0/24 10.10.90.0/24 eq 8089 # UDP上报 │
│ deny ip any 10.10.20.0/24 # 其他网段禁止│
│ permit ip 10.10.90.0/24 any # 监控中心全通│
└─────────────────────────────────────────────────────────────────┘方式 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|
DHCP | 部署快、IP集中管理 | 租约过期可能IP变化 | 临时测试 |
静态IP | 地址固定、不依赖DHCP服务 | 配置繁琐 | 生产环境(推荐) |
DHCP Reservation | 兼顾两者 | 需要交换机支持 | 有条件时可用 |
工程决策:传感器全部使用静态IP,通过配置文件固化。网关设备(汇聚服务器)使用DHCP Reservation。
┌─────────────────────────────────────────────────────────────────┐
│ 监控中心 │
│ ┌───────────────────────────────────────────────────────────┐ │
│ │ 数据汇聚网关 (Docking Gateway) │ │
│ │ ┌─────────────┐ ┌──────────────┐ ┌─────────────────┐ │ │
│ │ │ Modbus TCP │ │ UDP Collector│ │ SNMP Trap Receiver│ │ │
│ │ │ 连接池管理 │ │ 多播/单播 │ │ 事件告警 │ │ │
│ │ │ 16个子网 │ │ 239.255.x.x │ │ UDP/162 │ │ │
│ │ └──────┬──────┘ └──────┬───────┘ └────────┬────────┘ │ │
│ │ │ │ │ │ │
│ │ ┌──────▼─────────────────▼────────────────────▼────────┐ │ │
│ │ │ 统一数据总线 (ZeroMQ/Redis) │ │ │
│ │ │ Topic: archive/room/{room_id}/sensor/{sensor_id} │ │ │
│ │ └──────┬─────────────────┬────────────────────┬────────┘ │ │
│ │ │ │ │ │ │
│ │ ┌──────▼──────┐ ┌──────▼──────┐ ┌─────────▼────────┐ │ │
│ │ │ InfluxDB │ │ Prometheus │ │ 告警引擎 │ │ │
│ │ │ 时序存储 │ │ 指标监控 │ │ 阈值/趋势/突变 │ │ │
│ │ └─────────────┘ └─────────────┘ └──────────────────┘ │ │
│ └───────────────────────────────────────────────────────────┘ │
│ │
│ ┌─────────────┐ ┌──────────────┐ ┌───────────────────────┐ │
│ │ Grafana │ │ BA系统接口 │ │ 合规报表生成 │ │
│ │ 可视化大屏 │ │ BACnet/IP │ │ 日报/月报/年报 │ │
│ └─────────────┘ └──────────────┘ └───────────────────────┘ │
└─────────────────────────────────────────────────────────────────┘
▲ ▲ ▲
│ │ │
═════════╪══════════════╪═══════════════╪═════════════════════════
│ │ │
┌─────┴─────┐ ┌────┴────┐ ┌─────┴─────┐
│ 3F汇聚 │ │ 2F汇聚 │ │ 1F汇聚 │
│ 交换机 │ │ 交换机 │ │ 交换机 │
└─────┬─────┘ └────┬────┘ └─────┬─────┘
│ │ │
┌─────┴─────┐ ┌────┴────┐ ┌─────┴─────┐
│ VLAN 20 │ │ VLAN 50 │ │ VLAN 70 │
│ VLAN 30 │ │ VLAN 60 │ │ VLAN 80 │
│ VLAN 40 │ │ │ │ │
└───────────┘ └─────────┘ └───────────┘#!/usr/bin/env python3
"""
多库房数据汇聚网关
- 管理多个子网的Modbus TCP连接
- 统一UDP数据接收
- 数据归一化与路由
"""
import asyncio
import zmq
import json
from typing import Dict, List
from dataclasses import dataclass
from modbus_pool import ModbusConnectionPool
from udp_collector import UDPCollector
@dataclass
class RoomConfig:
room_id: str
vlan_id: int
subnet: str
sensors: List[dict]
poll_interval: int
alert_thresholds: dict
class AggregationGateway:
"""跨区域数据汇聚网关"""
def __init__(self, rooms: List[RoomConfig]):
self.rooms = {r.room_id: r for r in rooms}
self.modbus_pools: Dict[str, ModbusConnectionPool] = {}
self.udp_collector = UDPCollector()
self.zmq_ctx = zmq.Context()
self.zmq_pub = self.zmq_ctx.socket(zmq.PUB)
self.zmq_pub.bind("tcp://*:5555")
# 初始化各库房Modbus连接池
for room in rooms:
pool = ModbusConnectionPool(reconnect_interval=10)
for sensor in room.sensors:
dev = ModbusDevice(
dev_addr=sensor["dev_addr"],
ip=sensor["ip"],
port=502,
poll_interval=room.poll_interval,
bacnet_device_id=sensor["bacnet_device_id"],
objects=sensor["objects"]
)
pool.add_device(dev)
self.modbus_pools[room.room_id] = pool
async def start(self):
"""启动汇聚网关"""
# 连接所有子网的Modbus设备
for room_id, pool in self.modbus_pools.items():
asyncio.create_task(pool.connect_all())
# 启动Modbus轮询
for room_id, pool in self.modbus_pools.items():
asyncio.create_task(self._poll_room(room_id, pool))
# 启动UDP接收
asyncio.create_task(self.udp_collector.start())
# 启动数据归一化与发布
asyncio.create_task(self._normalize_and_publish())
logger.info("汇聚网关启动完成,管理%d个库房", len(self.rooms))
await asyncio.Event().wait()
async def _poll_room(self, room_id: str, pool: ModbusConnectionPool):
"""轮询单个库房的所有传感器"""
while True:
for dev in pool.devices.values():
if asyncio.get_event_loop().time() - dev.last_poll >= dev.poll_interval:
results = await pool.poll_device(dev)
if results:
# 发布到ZeroMQ
topic = f"archive/room/{room_id}/sensor/{dev.dev_addr}"
payload = {
"room_id": room_id,
"sensor_id": dev.dev_addr,
"timestamp": asyncio.get_event_loop().time(),
"data": results,
"source": "modbus_tcp"
}
self.zmq_pub.send_multipart([
topic.encode(),
json.dumps(payload).encode()
])
dev.last_poll = asyncio.get_event_loop().time()
await asyncio.sleep(1.0)
async def _normalize_and_publish(self):
"""数据归一化:统一时间戳、单位、格式"""
while True:
# 从UDP collector获取数据
udp_data = self.udp_collector.get_latest()
for pkt in udp_data:
room_id = self._lookup_room_by_ip(pkt["src_ip"])
if room_id:
topic = f"archive/room/{room_id}/sensor/{pkt['sensor_id']}"
payload = {
"room_id": room_id,
"sensor_id": pkt["sensor_id"],
"timestamp": pkt["collect_ts"],
"data": {
"temperature": pkt["temperature"],
"humidity": pkt["humidity"],
"voltage": pkt.get("voltage", 0)
},
"source": "udp"
}
self.zmq_pub.send_multipart([
topic.encode(),
json.dumps(payload).encode()
])
await asyncio.sleep(0.5)
def _lookup_room_by_ip(self, ip: str) -> str:
"""根据IP地址查找所属库房"""
for room in self.rooms.values():
if ip.startswith(room.subnet.rsplit('.', 1)[0]):
return room.room_id
return None
# 主程序
async def main():
# 加载配置
rooms = load_room_configs("rooms.yaml")
gateway = AggregationGateway(rooms)
await gateway.start()
if __name__ == "__main__":
asyncio.run(main()){
"room_id": "special_collection",
"sensor_id": 3,
"timestamp": 1704067200,
"data": {
"temperature": 18.5,
"humidity": 52.0,
"voltage": 48.2,
"status": 0
},
"source": "modbus_tcp",
"raw_register": {
"40001": 1850,
"40002": 520,
"40003": 4820
}
}# 核心交换机配置(华为S5735示例)
# 创建VLAN
vlan batch 10 20 30 40 50 60 70 80 90 99
# 配置VLAN接口IP
interface Vlanif20
ip address 10.10.20.1 255.255.255.0
description Special_Collection_Room
interface Vlanif30
ip address 10.10.30.1 255.255.255.0
description Paper_Archive_A
# ... 其他VLAN接口
# 配置ACL限制跨子网访问
acl number 3000
rule 10 permit tcp source 10.10.90.0 0.0.0.255 destination 10.10.20.0 0.0.0.255 destination-port eq 502
rule 20 permit udp source 10.10.20.0 0.0.0.255 destination 10.10.90.0 0.0.0.255 destination-port eq 8089
rule 30 deny ip source any destination 10.10.20.0 0.0.0.255
rule 40 permit ip source 10.10.90.0 0.0.0.255 destination any
# 应用ACL到VLAN接口
interface Vlanif20
traffic-filter inbound acl 3000场景 | 路由方式 | 理由 |
|---|---|---|
单栋建筑 | 静态路由 | 拓扑简单,无需动态协议 |
多栋建筑 | OSPF | 链路冗余,自动收敛 |
跨园区 | OSPF + 静态 | 核心用OSPF,出口用静态 |
工程选择:单栋档案馆建筑,使用静态路由。汇聚交换机配置默认路由指向监控中心。
┌─────────────────────────────────────────────────────────────────┐
│ Stratum 1: GPS/NTP服务器(监控中心) │
│ IP: 10.10.90.100 │
│ 同步源: GPS/北斗 + 互联网NTP │
└───────────────────────────┬─────────────────────────────────────┘
│ NTP广播/单播
┌───────────────────┼───────────────────┐
│ │ │
┌───────▼───────┐ ┌───────▼───────┐ ┌───────▼───────┐
│ 3F汇聚交换机 │ │ 2F汇聚交换机 │ │ 1F汇聚交换机 │
│ Stratum 2 │ │ Stratum 2 │ │ Stratum 2 │
│ 10.10.20.1 │ │ 10.10.50.1 │ │ 10.10.70.1 │
└───────┬───────┘ └───────┬───────┘ └───────┬───────┘
│ │ │
┌───────▼───────┐ ┌───────▼───────┐ ┌───────▼───────┐
│ 传感器群 │ │ 传感器群 │ │ 传感器群 │
│ Stratum 3 │ │ Stratum 3 │ │ Stratum 3 │
│ 10.10.20.x │ │ 10.10.50.x │ │ 10.10.70.x │
└───────────────┘ └───────────────┘ └───────────────┘// sensor_ntp.c
#include "lwip/opt.h"
#include "lwip/udp.h"
#include "lwip/dns.h"
#define NTP_SERVER_IP "10.10.90.100"
#define NTP_PORT 123
#define NTP_SYNC_INTERVAL 3600 // 1小时同步一次
#define NTP_RETRY_INTERVAL 10 // 失败重试10秒
static uint32_t g_ntp_last_sync = 0;
static int64_t g_ntp_offset = 0; // 与NTP服务器的时间偏差(秒)
void ntp_sync_task(void *pvParameters) {
struct netconn *conn;
struct netbuf *buf;
ip_addr_t server_ip;
uint32_t tx_ts, rx_ts;
while (1) {
vTaskDelay(pdMS_TO_TICKS(NTP_SYNC_INTERVAL * 1000));
conn = netconn_new(NETCONN_UDP);
if (!conn) continue;
ipaddr_aton(NTP_SERVER_IP, &server_ip);
netconn_connect(conn, &server_ip, NTP_PORT);
// 构造NTP请求包
uint8_t ntp_req[48] = {0};
ntp_req[0] = 0x1B; // LI=0, VN=3, Mode=3 (Client)
tx_ts = sys_now();
netconn_write(conn, ntp_req, 48, NETCONN_COPY);
// 等待响应(超时3秒)
err_t err = netconn_recv(conn, &buf);
if (err == ERR_OK) {
rx_ts = sys_now();
uint8_t *data = (uint8_t *)buf->p->payload;
// 解析NTP响应(Transmit Timestamp, bytes 40-47)
uint32_t tx_time = (data[40] << 24) | (data[41] << 16) | (data[42] << 8) | data[43];
// NTP时间戳从1900年开始,Unix时间戳从1970年开始,相差2208988800秒
uint32_t unix_time = tx_time - 2208988800UL;
// 计算偏差
uint32_t local_time = get_rtc_time();
g_ntp_offset = (int64_t)unix_time - (int64_t)local_time;
g_ntp_last_sync = local_time;
// 校正RTC
rtc_set_time(unix_time);
netbuf_delete(buf);
}
netconn_close(conn);
netconn_delete(conn);
}
}
uint32_t get_corrected_timestamp(void) {
return get_rtc_time() + g_ntp_offset;
}汇聚网关接收来自多个子网的UDP数据,可能存在重复包(网络重传、多播复制等)。使用滑动窗口去重:
# dedup.py
from collections import deque
import time
class DedupWindow:
"""序列号去重窗口"""
def __init__(self, window_size: int = 1000, ttl: int = 300):
self.window_size = window_size
self.ttl = ttl # 序列号有效期(秒)
self.seen = {} # {seq_num: timestamp}
self.max_seq = 0
def is_duplicate(self, seq_num: int, source_ip: str) -> bool:
"""检查是否为重复包"""
key = f"{source_ip}:{seq_num}"
# 清理过期条目
now = time.time()
expired = [k for k, v in self.seen.items() if now - v > self.ttl]
for k in expired:
del self.seen[k]
# 检查重复
if key in self.seen:
return True
# 记录新序列号
self.seen[key] = now
self.max_seq = max(self.max_seq, seq_num)
# 限制窗口大小
if len(self.seen) > self.window_size:
oldest_key = min(self.seen.keys(), key=lambda k: self.seen[k])
del self.seen[oldest_key]
return False当某个库房网络中断恢复后,传感器需要补传离线期间的数据:
补传协议流程:
传感器(客户端) 汇聚网关(服务端)
│ │
│── REPLAY_REQ ────────────────►│ 请求补传
│ {start_ts, end_ts} │ 指定时间范围
│ │
│◄── REPLAY_GRANT ─────────────│ 允许补传
│ {batch_size, interval} │ 批次大小、间隔
│ │
│── REPLAY_DATA ───────────────►│ 数据批次
│ {records: [...]} │ 最多32条
│ │
│◄── REPLAY_ACK ───────────────│ 确认
│ {last_ts} │ 最后一条时间戳
│ │
│ (循环直到补传完成) │# replay_server.py
class ReplayServer:
"""断点续传服务端"""
def __init__(self, storage):
self.storage = storage # InfluxDB或本地存储
self.pending_replays = {}
async def handle_replay_req(self, sensor_id: str, start_ts: int, end_ts: int):
"""处理补传请求"""
# 查询数据库中该传感器在指定时间范围内的数据
data = await self.storage.query(
sensor_id=sensor_id,
start=start_ts,
end=end_ts
)
if not data:
return {"status": "no_data"}
# 分批发送
batch_size = 32
batches = [data[i:i+batch_size] for i in range(0, len(data), batch_size)]
return {
"status": "granted",
"total_batches": len(batches),
"batch_size": batch_size,
"total_records": len(data)
}
async def send_batch(self, sensor_id: str, batch: list):
"""发送数据批次"""
payload = {
"sensor_id": sensor_id,
"batch_id": batch[0]["collect_ts"],
"records": batch
}
# 通过Modbus TCP或UDP发送
await self._send_replay_data(payload)级别 | 条件 | 响应时间 | 通知方式 |
|---|---|---|---|
紧急 | 温度>28℃或湿度>70%RH(纸质库) | 即时 | 短信+电话+平台弹窗 |
重要 | 温度>25℃或湿度>65%RH | 5分钟内 | 平台弹窗+邮件 |
一般 | 温度偏离设定值±2℃ | 30分钟内 | 平台记录 |
提示 | 传感器离线>5分钟 | 10分钟内 | 平台记录 |
告警引擎 → 根据room_id路由:
特藏库告警 → 馆长 + 技术部主任 + 值班人员
纸质库告警 → 技术部主任 + 值班人员
音像库告警 → 技术部主任 + 值班人员
过渡库告警 → 值班人员
通知渠道:
紧急 → 短信网关 + 电话语音 + 微信企业号
重要 → 微信企业号 + 邮件
一般 → 平台内消息
提示 → 日志记录项目 | 标准 | 方法 |
|---|---|---|
子网间连通性 | 监控中心可访问所有库房传感器 | ping测试 |
ACL有效性 | 非授权网段无法访问传感器 | nmap扫描 |
路由收敛 | 链路切换<3秒 | 断开测试 |
NTP同步 | 所有传感器时间偏差<1秒 | 时间戳比对 |
项目 | 标准 | 方法 |
|---|---|---|
数据完整性 | 汇聚网关入库率≥99.9% | 对比传感器本地存储 |
延迟 | 从采集到汇聚<5秒 | 时间戳差值统计 |
去重有效性 | 重复包率0% | 注入重复包测试 |
断点续传 | 恢复后数据完整补齐 | 模拟网络中断30分钟 |
多库房分区监测、子网划分、跨区域数据汇聚、VLAN、三层路由、Modbus TCP、UDP采集、ZeroMQ、数据归一化、NTP时间同步、序列号去重、断点续传、区域告警分级、档案库房、环境监控
多库房分区监测的核心在于"分区隔离、统一汇聚"——通过VLAN子网划分实现各库房网络隔离与安全控制,三层交换机路由打通监控中心与传感器子网,汇聚网关统一管理多区域Modbus TCP连接与UDP数据流,ZeroMQ消息总线实现数据归一化分发,NTP层级同步确保全馆时间一致,序列号去重窗口与断点续传协议保障数据完整性,最终形成一套可扩展、高可靠的多库房环境感知网络。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。