
关键词:存量档案馆、智能化升级、温湿度传感、恒温恒湿设备、对接改造、Modbus TCP、SNMP、POE供电、增量部署 标签:#物联网 #Modbus #TCP/IP #UDP #POE供电 #腾讯云 #Wireshark #Python #InfluxDB #以太网温湿度传感器 #网口温湿度变送器 #机房监控
新建档案馆可以从零规划网络、供电、点位,但存量档案馆的智能化升级面临完全不同的约束条件:
维度 | 新建档案馆 | 存量档案馆改造 |
|---|---|---|
网络基础设施 | 全新部署,按需设计 | 已有网络,可能无冗余,交换机端口紧张 |
供电条件 | 提前规划POE交换机 | 多数点位附近无网线,更无POE |
施工窗口 | 装修阶段同步施工 | 档案不能搬离,只能在非工作时间施工 |
设备利旧 | 全部新购 | 已有恒温恒湿设备需评估是否可接入 |
布线方式 | 桥架+明管 | 尽量利旧,新增走线需隐蔽,不能破坏装修 |
标准符合性 | 一步到位 | 分阶段实施,每阶段需独立可用 |
核心命题:在不影响档案安全、不中断现有业务的前提下,用最小侵入的方式完成温湿度感知层的增量部署和恒温恒湿设备的协议对接。
改造的第一步不是买设备,而是摸清楚现有网络"家底":
摸底清单:
□ 交换机品牌/型号/固件版本
□ 可用端口数量(含预留)
□ 是否支持POE(802.3af/at)
□ VLAN划分情况
□ 现有IP地址规划(网段、网关、DHCP范围)
□ 网络拓扑(核心-汇聚-接入层级)
□ 链路带宽利用率(峰值/均值)
□ 是否有独立的设备网(与办公网隔离)
□ 机柜空间和PDU剩余容量
□ 现有UPS续航能力关键决策点:传感器走独立网络还是复用现有办公/业务网?
方案 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|
独立VLAN | 成本低,利用现有交换机 | 依赖现有网络稳定性 | 现有网络质量好、有网管能力 |
独立物理网络 | 完全隔离,不依赖现有网络 | 成本高,需重新布线 | 现有网络老旧、不可控 |
混合方案 | 核心独立,接入复用 | 平衡成本与可靠性 | 大多数存量改造 |
推荐方案:在现有网络上为传感设备新建一个VLAN(如VLAN 100),配置独立的IP网段(如192.168.100.0/24),通过ACL限制该VLAN只能与采集服务器通信,不能访问其他网段。这样既有隔离性,又不需要重新布线。
存量档案馆的接入交换机端口通常已经用满。扩展方案:
方案A:级联POE交换机(推荐)
在现有接入交换机下挂一台POE交换机
优点:只需1个上联端口,新增8~24个POE端口
缺点:增加了单点故障
实施要点:
- 选择工业级POE交换机(档案馆环境可能有灰尘、温度波动)
- 上联端口配置为Trunk,允许传感VLAN通过
- POE预算计算:每台传感器约3~5W,8口交换机总POE功率≥60W
方案B:更换现有交换机为POE机型
优点:不增加设备层级
缺点:需要停机更换,影响现有业务
适用:有停机窗口(如节假日)
方案C:无线桥接(仅限无网线可达区域)
优点:无需布线
缺点:可靠性不如有线,不适合关键监测点
适用:临时或辅助点位存量档案馆传感器IP规划示例:
网段:192.168.100.0/24
网关:192.168.100.1
DHCP范围:192.168.100.100 ~ 192.168.100.200(临时调试用)
静态分配范围:192.168.100.10 ~ 192.168.100.99
分配规则:
192.168.100.1x → 库房A(10~19)
192.168.100.2x → 库房B(20~29)
192.168.100.3x → 库房C(30~39)
...
子网掩码:255.255.255.0
DNS:不需要(传感器不访问外网)
VLAN ID:100
交换机端口配置:
access端口 → 直接连接传感器
trunk端口 → 上联核心交换机存量档案馆不能像新建那样预埋线管,安装方式需要因地制宜:
安装场景 | 推荐方式 | 说明 |
|---|---|---|
密集架通道顶部 | 磁吸+POE | 利用密集架顶部金属面,磁吸固定,沿架体走线至最近网络端口 |
墙面(已有插座附近) | 壁挂+POE | 86盒面板式或直接壁挂,网线沿踢脚线走 |
库房中央(无墙面) | 吊装+POE | 从天花板桥架引下,美观但施工量大 |
密集架内部 | 超薄型+短网线 | 选择厚度<2cm的传感器,夹在架体隔板间 |
通道入口 | 立柱安装 | 利用通道立柱,扎带或胶粘固定 |
存量档案馆走线原则:
1. 优先利旧:查看是否有废弃的电话线、同轴电缆管道可复用
2. 沿边走线:沿踢脚线、墙角、桥架边缘,避免明线横跨通道
3. 隐蔽处理:使用与墙面同色的线槽,或用装饰条遮盖
4. 桥架借用:如果上方有桥架,优先从桥架走,下方用软管引下
5. 长度控制:单根网线≤80m(POE供电衰减),超长需加POE中继器
走线材料清单:
□ 超五类/六类非屏蔽网线(室外用需选室外线)
□ PVC线槽(宽度20~40mm,颜色匹配墙面)
□ 磁吸底座(传感器用)
□ 扎带、线卡、膨胀螺丝
□ 水晶头、接线子(备用)
□ 防水盒(如果有潮湿区域)存量档案馆施工约束:
□ 档案不能搬离 → 施工人员不能进入密集架内部操作(除非架体打开且有人监护)
□ 工作时间不能断网 → 网络配置在周末或夜间进行
□ 不能产生大量灰尘 → 打孔需使用吸尘式电锤,走线尽量免打孔
□ 噪音限制 → 避免在档案利用高峰期施工
推荐施工排期(以100个传感器为例):
第1周:网络评估、方案设计、设备采购
第2周:交换机配置、VLAN划分(周末夜间)
第3-4周:传感器安装(每天安装10~15个,非工作时间)
第5周:设备对接调试(与现有恒温恒湿设备联调)
第6周:系统联调、验收测试存量档案馆通常已有恒温恒湿设备,但品牌型号各异,通信能力参差不齐:
设备摸底清单:
□ 品牌/型号/出厂年份
□ 是否有通信接口(RS485/RS232/以太网)
□ 通信协议(Modbus RTU/TCP、BACnet、 proprietary)
□ 是否有控制接口(能远程启停/调速)还是仅监测
□ 设备当前运行状态(正常/故障/老化)
□ 是否有原厂提供的通信协议文档
□ 是否还在保修期内(改造是否影响保修)设备分类与对接策略:
设备类型 | 通信能力 | 对接方案 |
|---|---|---|
有以太网口+标准协议 | 直接支持Modbus TCP | 直接接入,配置IP和寄存器映射 |
有RS485+Modbus RTU | 串口通信 | 加装串口服务器(RS485转TCP) |
仅有干接点/模拟量 | 开关量/4-20mA | 加装I/O模块(如ADAM-6000系列) |
无通信接口 | 仅面板操作 | 加装外挂控制器(红外遥控/继电器模拟按键) |
协议不开放 | 私有协议 | 联系原厂获取SDK/协议文档,或加装第三方控制器 |
"""
legacy_device_adapter.py - 存量设备协议适配层
"""
import socket
import struct
import time
from typing import Dict, Any, Optional
from enum import Enum
# ============================================================
# Modbus TCP 客户端(用于直接支持Modbus的设备)
# ============================================================
class ModbusTCPClient:
"""
Modbus TCP 客户端,用于对接支持Modbus TCP的恒温恒湿设备
"""
def __init__(self, ip: str, port: int = 502, slave_id: int = 1,
timeout: float = 3.0):
self.ip = ip
self.port = port
self.slave_id = slave_id
self.timeout = timeout
self.sock: Optional[socket.socket] = None
self.transaction_id = 0
def connect(self) -> bool:
try:
self.sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
self.sock.settimeout(self.timeout)
self.sock.connect((self.ip, self.port))
return True
except Exception as e:
print(f"连接失败 {self.ip}:{self.port} - {e}")
return False
def disconnect(self):
if self.sock:
self.sock.close()
self.sock = None
def _next_transaction_id(self) -> int:
self.transaction_id = (self.transaction_id + 1) % 65536
return self.transaction_id
def read_holding_registers(self, address: int, count: int) -> Optional[list]:
"""
读取保持寄存器 (Function Code 0x03)
"""
if not self.sock:
if not self.connect():
return None
tid = self._next_transaction_id()
# Modbus TCP MBAP Header + PDU
# Transaction ID (2) + Protocol ID (2) + Length (2) + Unit ID (1) + FC (1) + Addr (2) + Count (2)
length = 6 # 1(Unit) + 1(FC) + 2(Addr) + 2(Count)
packet = struct.pack('>HHHBH', tid, 0, length, self.slave_id, 0x03, address, count)
try:
self.sock.send(packet)
response = self.sock.recv(1024)
if len(response) < 9:
return None
# 解析响应
header = response[:8]
resp_tid, _, resp_len, resp_unit = struct.unpack('>HHHB', header)
if resp_tid != tid:
return None
byte_count = response[8]
data = response[9:]
registers = []
for i in range(byte_count // 2):
val = struct.unpack('>H', data[i*2:i*2+2])[0]
registers.append(val)
return registers
except socket.timeout:
print(f"读取超时: {self.ip}")
self.disconnect()
return None
except Exception as e:
print(f"读取错误: {e}")
self.disconnect()
return None
def write_single_register(self, address: int, value: int) -> bool:
"""
写入单个寄存器 (Function Code 0x06)
"""
if not self.sock:
if not self.connect():
return False
tid = self._next_transaction_id()
length = 6
packet = struct.pack('>HHHHBHH', tid, 0, length, self.slave_id, 0x06, address, value)
try:
self.sock.send(packet)
response = self.sock.recv(1024)
# 写寄存器响应与请求相同
return len(response) >= 12
except Exception as e:
print(f"写入错误: {e}")
self.disconnect()
return False
# ============================================================
# 串口服务器桥接(用于RS485 Modbus RTU设备)
# ============================================================
class SerialServerBridge:
"""
通过串口服务器(如MOXA NPort)将RS485 Modbus RTU转为TCP
串口服务器配置为TCP Server模式,本类作为TCP Client连接
发送Modbus RTU帧,串口服务器透明传输
"""
def __init__(self, ip: str, port: int = 4001, slave_id: int = 1):
self.ip = ip
self.port = port
self.slave_id = slave_id
self.sock: Optional[socket.socket] = None
def connect(self) -> bool:
try:
self.sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
self.sock.settimeout(3.0)
self.sock.connect((self.ip, self.port))
return True
except Exception as e:
print(f"串口服务器连接失败: {e}")
return False
def disconnect(self):
if self.sock:
self.sock.close()
self.sock = None
def _modbus_rtu_crc(self, data: bytes) -> bytes:
"""计算Modbus RTU CRC16"""
crc = 0xFFFF
for byte in data:
crc ^= byte
for _ in range(8):
if crc & 0x0001:
crc = (crc >> 1) ^ 0xA001
else:
crc >>= 1
return struct.pack('<H', crc)
def read_holding_registers_rtu(self, address: int, count: int) -> Optional[list]:
"""
通过串口服务器发送Modbus RTU读寄存器命令
"""
if not self.sock:
if not self.connect():
return None
# Modbus RTU帧: SlaveID + FC + Addr + Count + CRC
frame = struct.pack('>BBHH', self.slave_id, 0x03, address, count)
frame += self._modbus_rtu_crc(frame)
try:
self.sock.send(frame)
# RTU响应需要等待完整帧
time.sleep(0.1)
response = self.sock.recv(1024)
if len(response) < 5:
return None
# 验证CRC
data = response[:-2]
crc_received = response[-2:]
crc_calculated = self._modbus_rtu_crc(data)
if crc_received != crc_calculated:
print("CRC校验失败")
return None
# 解析数据
byte_count = response[2]
registers = []
for i in range(byte_count // 2):
offset = 3 + i * 2
val = struct.unpack('>H', response[offset:offset+2])[0]
registers.append(val)
return registers
except Exception as e:
print(f"RTU读取错误: {e}")
return None
# ============================================================
# 存量设备抽象适配器
# ============================================================
class LegacyDeviceAdapter:
"""
存量设备统一适配器
将不同通信方式的设备统一为相同的接口
"""
def __init__(self, device_id: str, name: str, adapter_type: str, **kwargs):
self.device_id = device_id
self.name = name
self.adapter_type = adapter_type
self.client = None
self.reg_map = kwargs.get('reg_map', {})
self._init_client(**kwargs)
def _init_client(self, **kwargs):
if self.adapter_type == 'modbus_tcp':
self.client = ModbusTCPClient(
ip=kwargs.get('ip'),
port=kwargs.get('port', 502),
slave_id=kwargs.get('slave_id', 1)
)
elif self.adapter_type == 'serial_server':
self.client = SerialServerBridge(
ip=kwargs.get('ip'),
port=kwargs.get('port', 4001),
slave_id=kwargs.get('slave_id', 1)
)
def read_status(self) -> Dict[str, Any]:
"""读取设备状态"""
if not self.client:
return {"error": "client not initialized"}
result = {
"device_id": self.device_id,
"name": self.name,
"timestamp": time.time()
}
# 根据寄存器映射读取各项参数
for key, reg_info in self.reg_map.items():
addr = reg_info.get('address')
count = reg_info.get('count', 1)
scale = reg_info.get('scale', 1.0)
unit = reg_info.get('unit', '')
if isinstance(self.client, ModbusTCPClient):
raw = self.client.read_holding_registers(addr, count)
else:
raw = self.client.read_holding_registers_rtu(addr, count)
if raw:
value = raw[0] * scale
result[key] = f"{value}{unit}"
else:
result[key] = "N/A"
return result
def control(self, command: str, **params) -> bool:
"""发送控制命令"""
if not self.client:
return False
cmd_map = self.reg_map.get('_commands', {})
if command not in cmd_map:
print(f"未知命令: {command}")
return False
reg_addr = cmd_map[command].get('address')
value = cmd_map[command].get('value')
if isinstance(self.client, ModbusTCPClient):
return self.client.write_single_register(reg_addr, value)
else:
# 串口服务器模式也用相同接口
# 实际实现中需要发送RTU写命令
print("串口写命令需实现RTU帧")
return False
# ============================================================
# 使用示例
# ============================================================
if __name__ == "__main__":
# 场景1:直接Modbus TCP设备(较新的精密空调)
device1 = LegacyDeviceAdapter(
device_id="AC-OLD-01",
name="精密空调-01(存量)",
adapter_type="modbus_tcp",
ip="192.168.100.50",
port=502,
slave_id=1,
reg_map={
"run_status": {"address": 0x1000, "scale": 1, "unit": ""},
"return_temp": {"address": 0x1001, "scale": 0.1, "unit": "℃"},
"return_hum": {"address": 0x1002, "scale": 0.1, "unit": "%RH"},
"set_temp": {"address": 0x2000, "scale": 0.1, "unit": "℃"},
"set_hum": {"address": 0x2001, "scale": 0.1, "unit": "%RH"},
"_commands": {
"start": {"address": 0x3000, "value": 1},
"stop": {"address": 0x3000, "value": 0},
"set_temp_23": {"address": 0x2000, "value": 230},
}
}
)
print("读取设备状态:")
status = device1.read_status()
for k, v in status.items():
print(f" {k}: {v}")
# 场景2:RS485 Modbus RTU设备(通过串口服务器)
device2 = LegacyDeviceAdapter(
device_id="DEH-OLD-01",
name="除湿机-01(存量RS485)",
adapter_type="serial_server",
ip="192.168.100.51",
port=4001,
slave_id=2,
reg_map={
"run_status": {"address": 0x0001, "scale": 1, "unit": ""},
"current_hum": {"address": 0x0002, "scale": 0.1, "unit": "%RH"},
"setpoint_hum": {"address": 0x1001, "scale": 0.1, "unit": "%RH"},
}
)
print("\n读取RS485设备状态:")
status2 = device2.read_status()
for k, v in status2.items():
print(f" {k}: {v}")对于完全没有通信接口的老旧设备(仅面板操作),可通过以下方式实现远程控制:
方案A:红外遥控模拟
适用:设备有红外遥控器的
组件:红外发射模块(如Broadlink RM4 Pro)+ 学习遥控码
优点:不改设备硬件
缺点:红外需对准,可靠性一般
方案B:继电器模拟按键
适用:设备面板有物理按键
组件:继电器模块 + 机械臂/电磁铁
优点:物理模拟,可靠
缺点:需打开面板,可能失去保修
方案C:干接点并联
适用:设备有外部干接点控制端子(如消防联动端子)
组件:I/O模块输出干接点
优点:利用设备原有控制接口,安全
缺点:不是所有设备都有
方案D:加装第三方控制器
适用:设备有压缩机/风机控制电路
组件:加装温控器/湿度控制器,替换原控制板
优点:完全可控
缺点:改动大,需专业人员,失去原厂保修推荐优先级:干接点并联 > 红外遥控模拟 > 继电器模拟按键 > 加装第三方控制器
┌─────────────────────────────────────────────────────────────────────┐
│ 存量档案馆数据采集架构 │
├─────────────────────────────────────────────────────────────────────┤
│ │
│ ┌────────────────────────────────────────────────────────────────┐│
│ │ 感知层(新增传感器) ││
│ │ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────────────┐││
│ │ │ 传感器01 │ │ 传感器02 │ │ 传感器03 │ │ ... 传感器N │││
│ │ │ UDP+SNMP │ │ UDP+SNMP │ │ UDP+SNMP │ │ UDP+SNMP │││
│ │ └────┬─────┘ └────┬─────┘ └────┬─────┘ └────────┬─────────┘││
│ └───────┼─────────────┼─────────────┼───────────────┼──────────┘│
│ │ │ │ │ │
│ ┌───────▼─────────────▼─────────────▼───────────────▼──────────┐│
│ │ 采集服务器(边缘网关) ││
│ │ ┌───────────────────────────────────────────────────────────┐││
│ │ │ 采集服务(Python) │││
│ │ │ · UDP监听(传感器主动上报) │││
│ │ │ · SNMP轮询(定时读取传感器) │││
│ │ │ · Modbus TCP客户端(读取恒温恒湿设备) │││
│ │ │ · 串口服务器桥接(RS485设备) │││
│ │ └───────────────────────────┬───────────────────────────────┘││
│ │ │ ││
│ │ ┌───────────────────────────▼───────────────────────────────┐││
│ │ │ 数据处理与转发 │││
│ │ │ · 数据校验与异常剔除 │││
│ │ │ · 滑动窗口滤波 │││
│ │ │ · 多传感器融合判定 │││
│ │ │ · 转发至InfluxDB(时序存储) │││
│ │ │ · 转发至MQTT Broker(实时订阅) │││
│ │ │ · 转发至腾讯云IoT Hub(云端备份) │││
│ │ └───────────────────────────────────────────────────────────┘││
│ └──────────────────────────────────────────────────────────────────┘│
│ │ │
│ ┌───────────────────────────▼──────────────────────────────────┐ │
│ │ 应用层 │ │
│ │ ┌─────────────┐ ┌─────────────┐ ┌─────────────────────────┐│ │
│ │ │ SCADA组态 │ │ 联动引擎 │ │ 八防环境健康度平台 ││ │
│ │ │ (可视化) │ │ (自动控制) │ │ (综合评分+报表) ││ │
│ │ └─────────────┘ └─────────────┘ └─────────────────────────┘│ │
│ └──────────────────────────────────────────────────────────────────┘│
└─────────────────────────────────────────────────────────────────────┘"""
collector_service.py - 存量档案馆数据采集服务
"""
import socket
import struct
import threading
import time
from datetime import datetime
from typing import Dict, List, Callable
import queue
# 尝试导入可选依赖
try:
from pysnmp.hlapi import (
getCmd, SnmpEngine, CommunityData, UdpTransportTarget,
ContextData, ObjectType, ObjectIdentity
)
HAS_SNMP = True
except ImportError:
HAS_SNMP = False
print("pysnmp 未安装,SNMP功能不可用")
try:
from influxdb_client import InfluxDBClient, Point
from influxdb_client.client.write_api import SYNCHRONOUS
HAS_INFLUXDB = True
except ImportError:
HAS_INFLUXDB = False
print("influxdb_client 未安装,InfluxDB功能不可用")
try:
import paho.mqtt.client as mqtt
HAS_MQTT = True
except ImportError:
HAS_MQTT = False
print("paho-mqtt 未安装,MQTT功能不可用")
class SensorData:
"""传感器数据结构"""
def __init__(self, sensor_id: str, temperature: float, humidity: float,
timestamp: float = None):
self.sensor_id = sensor_id
self.temperature = temperature
self.humidity = humidity
self.timestamp = timestamp or time.time()
def to_dict(self) -> Dict:
return {
"sensor_id": self.sensor_id,
"temperature": self.temperature,
"humidity": self.humidity,
"timestamp": self.timestamp
}
class UDPServer:
"""
UDP服务器,接收传感器主动上报的数据
"""
def __init__(self, host: str = "0.0.0.0", port: int = 9000,
data_callback: Callable[[SensorData], None] = None):
self.host = host
self.port = port
self.data_callback = data_callback
self.sock = None
self.running = False
self.thread = None
def start(self):
self.sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
self.sock.bind((self.host, self.port))
self.running = True
self.thread = threading.Thread(target=self._listen_loop, daemon=True)
self.thread.start()
print(f"UDP服务器启动在 {self.host}:{self.port}")
def stop(self):
self.running = False
if self.sock:
self.sock.close()
def _listen_loop(self):
while self.running:
try:
data, addr = self.sock.recvfrom(1024)
sensor_data = self._parse_udp_data(data, addr)
if sensor_data and self.data_callback:
self.data_callback(sensor_data)
except Exception as e:
print(f"UDP接收错误: {e}")
def _parse_udp_data(self, data: bytes, addr) -> SensorData:
"""
解析传感器UDP上报数据
假设数据格式: sensor_id(4字节) + temp(4字节float) + hum(4字节float)
"""
try:
if len(data) < 12:
return None
sensor_id = data[:4].decode('ascii', errors='ignore').strip('\x00')
temp = struct.unpack('>f', data[4:8])[0]
hum = struct.unpack('>f', data[8:12])[0]
return SensorData(sensor_id=sensor_id, temperature=temp, humidity=hum)
except Exception:
return None
class SNMPScanner:
"""
SNMP轮询器,定时读取传感器数据
"""
def __init__(self, targets: List[Dict], interval: int = 30,
data_callback: Callable[[SensorData], None] = None):
"""
targets: [{"ip": "192.168.100.10", "community": "public",
"temp_oid": "...", "hum_oid": "..."}]
"""
self.targets = targets
self.interval = interval
self.data_callback = data_callback
self.running = False
self.thread = None
def start(self):
if not HAS_SNMP:
print("SNMP不可用,跳过SNMP轮询")
return
self.running = True
self.thread = threading.Thread(target=self._scan_loop, daemon=True)
self.thread.start()
print(f"SNMP轮询器启动,间隔 {self.interval}s")
def stop(self):
self.running = False
def _scan_loop(self):
while self.running:
for target in self.targets:
try:
data = self._query_snmp(target)
if data and self.data_callback:
self.data_callback(data)
except Exception as e:
print(f"SNMP查询失败 {target['ip']}: {e}")
time.sleep(self.interval)
def _query_snmp(self, target: Dict) -> SensorData:
# SNMP查询实现(需要pysnmp)
# 这里简化示意
pass
class InfluxDBWriter:
"""
InfluxDB写入器
"""
def __init__(self, url: str, token: str, org: str, bucket: str):
self.url = url
self.token = token
self.org = org
self.bucket = bucket
self.client = None
self.write_api = None
if HAS_INFLUXDB:
self.client = InfluxDBClient(url=url, token=token, org=org)
self.write_api = self.client.write_api(write_options=SYNCHRONOUS)
def write(self, data: SensorData):
if not HAS_INFLUXDB or not self.write_api:
return
point = Point("environment") \
.tag("sensor_id", data.sensor_id) \
.field("temperature", data.temperature) \
.field("humidity", data.humidity) \
.time(datetime.fromtimestamp(data.timestamp))
self.write_api.write(bucket=self.bucket, org=self.org, record=point)
def close(self):
if self.client:
self.client.close()
class MQTTPublisher:
"""
MQTT发布器
"""
def __init__(self, broker: str, port: int = 1883, topic_prefix: str = "archive/env"):
self.broker = broker
self.port = port
self.topic_prefix = topic_prefix
self.client = None
if HAS_MQTT:
self.client = mqtt.Client()
try:
self.client.connect(broker, port, 60)
self.client.loop_start()
except Exception as e:
print(f"MQTT连接失败: {e}")
def publish(self, data: SensorData):
if not HAS_MQTT or not self.client:
return
import json
topic = f"{self.topic_prefix}/{data.sensor_id}"
payload = json.dumps(data.to_dict())
self.client.publish(topic, payload, qos=1)
def close(self):
if self.client:
self.client.loop_stop()
self.client.disconnect()
class CollectorService:
"""
数据采集服务主类
"""
def __init__(self):
self.udp_server = None
self.snmp_scanner = None
self.influx_writer = None
self.mqtt_publisher = None
self.data_queue = queue.Queue(maxsize=1000)
self.running = False
self.processor_thread = None
def setup_udp(self, host: str = "0.0.0.0", port: int = 9000):
self.udp_server = UDPServer(
host=host, port=port,
data_callback=self._on_sensor_data
)
def setup_snmp(self, targets: List[Dict], interval: int = 30):
self.snmp_scanner = SNMPScanner(
targets=targets, interval=interval,
data_callback=self._on_sensor_data
)
def setup_influxdb(self, url: str, token: str, org: str, bucket: str):
self.influx_writer = InfluxDBWriter(url, token, org, bucket)
def setup_mqtt(self, broker: str, port: int = 1883):
self.mqtt_publisher = MQTTPublisher(broker, port)
def _on_sensor_data(self, data: SensorData):
"""数据回调"""
try:
self.data_queue.put_nowait(data)
except queue.Full:
print("数据队列满,丢弃数据")
def _process_data(self):
"""数据处理线程"""
while self.running:
try:
data = self.data_queue.get(timeout=1.0)
# 写入InfluxDB
if self.influx_writer:
self.influx_writer.write(data)
# 发布到MQTT
if self.mqtt_publisher:
self.mqtt_publisher.publish(data)
# 打印日志
print(f"[{datetime.fromtimestamp(data.timestamp).strftime('%H:%M:%S')}] "
f"{data.sensor_id}: {data.temperature:.1f}℃, {data.humidity:.1f}%RH")
except queue.Empty:
continue
except Exception as e:
print(f"数据处理错误: {e}")
def start(self):
"""启动所有组件"""
self.running = True
if self.udp_server:
self.udp_server.start()
if self.snmp_scanner:
self.snmp_scanner.start()
self.processor_thread = threading.Thread(
target=self._process_data, daemon=True
)
self.processor_thread.start()
print("数据采集服务已启动")
def stop(self):
"""停止所有组件"""
self.running = False
if self.udp_server:
self.udp_server.stop()
if self.snmp_scanner:
self.snmp_scanner.stop()
if self.influx_writer:
self.influx_writer.close()
if self.mqtt_publisher:
self.mqtt_publisher.close()
print("数据采集服务已停止")
# ============================================================
# 使用示例
# ============================================================
if __name__ == "__main__":
service = CollectorService()
# 配置UDP接收
service.setup_udp(host="0.0.0.0", port=9000)
# 配置SNMP轮询
snmp_targets = [
{
"ip": "192.168.100.10",
"community": "public",
"temp_oid": "1.3.6.1.4.1.12345.1.1.0",
"hum_oid": "1.3.6.1.4.1.12345.1.2.0"
},
{
"ip": "192.168.100.11",
"community": "public",
"temp_oid": "1.3.6.1.4.1.12345.1.1.0",
"hum_oid": "1.3.6.1.4.1.12345.1.2.0"
}
]
service.setup_snmp(targets=snmp_targets, interval=30)
# 配置InfluxDB(可选)
# service.setup_influxdb(
# url="http://localhost:8086",
# token="your-token",
# org="archive",
# bucket="environment"
# )
# 配置MQTT(可选)
# service.setup_mqtt(broker="localhost")
# 启动服务
service.start()
# 主循环
try:
while True:
time.sleep(1)
except KeyboardInterrupt:
print("\n正在停止服务...")
service.stop()阶段一:感知层覆盖(第1~2个月)
目标:完成所有库房温湿度传感器的增量部署
交付物:
□ 传感器安装到位,全部在线
□ 数据采集服务运行,数据写入InfluxDB
□ 基础可视化(Grafana仪表盘)
□ 告警功能(短信/邮件通知)
验收标准:所有传感器数据连续7天无中断
阶段二:设备对接(第3~4个月)
目标:完成存量恒温恒湿设备的协议对接
交付物:
□ 设备通信测试报告
□ 寄存器映射表
□ 设备状态监控画面
□ 手动远程控制功能
验收标准:能远程读取所有设备状态,能远程启停
阶段三:联动自动化(第5~6个月)
目标:实现基于八防规则的自动联动
交付物:
□ 联动规则引擎部署
□ 八防规则配置完成
□ 自动启停功能上线
□ 防抖/死区参数调优
验收标准:模拟异常条件下,系统能正确触发联动
阶段四:优化与完善(第7~12个月)
目标:系统稳定运行,持续优化
交付物:
□ 运行数据分析报告
□ 阈值和死区优化调整
□ 运维管理功能完善
□ 年度运行报告
验收标准:系统全年可用率>99.5%,八防评分>80阶段一回退:传感器数据采集失败
→ 回退到人工巡检+手持式温湿度计
→ 不影响现有设备和业务
阶段二回退:设备对接失败
→ 保持原有手动控制方式
→ 仅做状态监测,不做自动控制
阶段三回退:联动异常
→ 一键切换到手动模式
→ 所有设备恢复本地控制
阶段四回退:优化参数不合适
→ 恢复到上一版本参数
→ 参数调整需经过审批流程存量档案馆的智能化升级,核心挑战不是技术本身,而是在受限条件下找到最优的增量部署方案。网络评估先行,传感器增量部署以最小侵入为原则,设备对接以"能接尽接、不能接则外挂"为策略。分阶段实施确保每步都可验证、可回退。系统上线后,持续运行观察和数据积累比一次性完美部署更重要——先用起来,再持续优化。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。