
关键词:电厂变电站、机房温湿度采集、UDP上报、SNMP轮询、双协议采集、系统架构、网关设计、数据融合、高可用 标签:#物联网 #TCP/IP #UDP #SNMP #POE供电 #以太网温湿度传感器 #网口温湿度变送器 #机房监控
电厂变电站机房(升压站继电器室、保护小室、通信机房、直流系统室)与传统数据中心和工业厂房相比,有其独特的环境与运行约束:
维度 | 特征 | 对采集系统的影响 |
|---|---|---|
电磁环境 | 隔离开关操作产生强瞬态电磁脉冲(可达数十kV/m) | 传感器需金属屏蔽外壳,网线需屏蔽 |
分布形态 | 多个小室分散布置,间距 50~300m | 网络需分层汇聚,不能简单扁平化 |
供电 | 操作直流系统(DC 110V/220V)+ UPS | POE交换机需双电源输入 |
运维 | 无人值守,定期巡检 | 数据必须远程可访问,告警实时推送 |
通信网络 | 综合数据网(生产控制大区/管理信息大区) | 采集系统需适配电力安防分区 |
环境 | 无精密空调,仅有工业空调或通风 | 温湿度波动大,夏季可能超温 |
设备寿命 | 要求 10~15 年不更换 | 硬件选型需工业级,固件需长期可维护 |
核心需求:连续、可靠、标准化的温湿度数据采集,既能满足实时监控(告警联动),又能对接电力综合监控系统(SCADA/动环平台)。
┌─────────────────────────────────────────────────────────────────────┐
│ 管理信息大区(或独立动环网) │
├─────────────────────────────────────────────────────────────────────┤
│ │
│ ┌──────────────────────────────────────────────────────────────┐ │
│ │ 采集网关服务器(双机热备) │ │
│ │ ┌─────────────────┐ ┌─────────────────┐ ┌──────────────┐ │ │
│ │ │ UDP接收服务 │ │ SNMP轮询引擎 │ │ 数据融合引擎 │ │ │
│ │ │ (端口9000) │ │ (每30s) │ │ (链路选择) │ │ │
│ │ └────────┬────────┘ └────────┬────────┘ └──────┬───────┘ │ │
│ │ │ │ │ │ │
│ │ ┌────────▼─────────────────────▼───────────────────▼───────┐│ │
│ │ │ 统一数据缓存 (Redis/共享内存) ││ │
│ │ └───────────────────────────────┬───────────────────────────┘│ │
│ │ │ │ │
│ │ ┌───────────────────────────────▼───────────────────────────┐│ │
│ │ │ 北向接口层 ││ │
│ │ │ - Modbus TCP (对接SCADA) ││ │
│ │ │ - IEC 60870-5-104 (电力规约,可选) ││ │
│ │ │ - MQTT (对接云平台/动环平台) ││ │
│ │ │ - REST API (本地调试/Web展示) ││ │
│ │ └───────────────────────────────────────────────────────────┘│ │
│ └────────────────────────────────────────────────────────────────┘ │
│ │ │
│ ┌───────────────┼───────────────┐ │
│ │ │ │ │
│ ┌────────▼───────┐ ┌─────▼───────┐ ┌────▼──────────┐ │
│ │ 继电器室A │ │ 继电器室B │ │ 通信机房 │ │
│ │ 接入交换机 │ │ 接入交换机 │ │ 接入交换机 │ │
│ │ (工业级,POE) │ │ (工业级,POE) │ │ (工业级,POE) │ │
│ └───────┬───────┘ └──────┬──────┘ └──────┬────────┘ │
│ │ │ │ │
└─────────────┼─────────────────┼───────────────┼────────────────────┘
│ │ │
┌──────────▼───┐ ┌────────▼───┐ ┌──────▼────────┐
│ 传感器1~8 │ │ 传感器9~16 │ │ 传感器17~24 │
│ (UDP+SNMP) │ │ (UDP+SNMP) │ │ (UDP+SNMP) │
│ 机柜内/墙挂 │ │ 机柜内/墙挂 │ │ 机柜内/墙挂 │
└──────────────┘ └──────────────┘ └───────────────┘核心层:采集网关(双机热备,连接核心交换机)
│
汇聚层:工业环网交换机(可选,用于多个小室互联)
│
接入层:POE接入交换机(每小室1~2台,24口/16口)
│
终端层:以太网温湿度变送器(每机柜1~2个,每小室8~16个)网络隔离要求(电力二次系统安全防护规定):
"""
udp_receiver.py - 网关侧UDP接收服务
"""
import asyncio
import struct
from dataclasses import dataclass
from typing import Dict, Optional
@dataclass
class SensorReading:
dev_id: str
temperature: float
humidity: float
dewpoint: float
sequence: int
status: int
recv_ts: float
source_ip: str
class UDPReceiver:
def __init__(self, host="0.0.0.0", port=9000, cache_ttl=60):
self.host = host
self.port = port
self.cache_ttl = cache_ttl
self.cache: Dict[str, SensorReading] = {}
self.stats = {"packets_recv": 0, "parse_err": 0, "last_reset": asyncio.get_event_loop().time()}
self.transport = None
def start(self, loop):
listen = loop.create_datagram_endpoint(
lambda: self,
local_addr=(self.host, self.port)
)
self.transport, _ = loop.run_until_complete(listen)
print(f"[UDP] Listening on {self.host}:{self.port}")
def datagram_received(self, data, addr):
self.stats["packets_recv"] += 1
try:
reading = self._parse(data, addr)
if reading:
self.cache[reading.dev_id] = reading
# 清理过期数据
self._cleanup()
except Exception as e:
self.stats["parse_err"] += 1
def _parse(self, data, addr) -> Optional[SensorReading]:
if len(data) < 15:
return None
header, dev_id_int = struct.unpack_from(">HH", data, 0)
if header != 0xAA55:
return None
temp, hum, dp, seq = struct.unpack_from(">hhhH", data, 4)
status = struct.unpack_from(">B", data, 12)[0]
crc_received = struct.unpack_from(">H", data, 13)[0]
crc_calc = self._crc16(data[:13])
if crc_received != crc_calc:
return None
dev_id = f"sub-{dev_id_int:04d}"
return SensorReading(
dev_id=dev_id,
temperature=temp / 10.0,
humidity=hum / 10.0,
dewpoint=dp / 10.0,
sequence=seq,
status=status,
recv_ts=asyncio.get_event_loop().time(),
source_ip=addr[0]
)
def _crc16(self, data: bytes) -> int:
crc = 0xFFFF
for b in data:
crc ^= b
for _ in range(8):
if crc & 1:
crc = (crc >> 1) ^ 0xA001
else:
crc >>= 1
return crc
def _cleanup(self):
now = asyncio.get_event_loop().time()
expired = [k for k, v in self.cache.items() if now - v.recv_ts > self.cache_ttl]
for k in expired:
del self.cache[k]
def get_latest(self, dev_id: str) -> Optional[SensorReading]:
return self.cache.get(dev_id)
def connection_made(self, transport):
self.transport = transport
def error_received(self, exc):
print(f"[UDP] Error: {exc}")"""
snmp_poller.py - SNMP轮询引擎
"""
import asyncio
from pysnmp.hlapi import *
from typing import Dict, Optional
import time
class SNMPEngine:
def __init__(self, devices: list, poll_interval=30, timeout=3, retries=2):
self.devices = devices
self.poll_interval = poll_interval
self.timeout = timeout
self.retries = retries
self.cache: Dict[str, dict] = {}
self.stats = {"polls": 0, "failures": 0, "timeouts": 0}
self._running = False
async def start(self):
self._running = True
while self._running:
tasks = [self._poll(dev) for dev in self.devices]
await asyncio.gather(*tasks, return_exceptions=True)
await asyncio.sleep(self.poll_interval)
async def _poll(self, dev):
self.stats["polls"] += 1
ip = dev["ip"]
community = dev.get("community", "public")
oids = dev.get("oids", {
"temperature": "1.3.6.1.4.1.12345.2.1.0",
"humidity": "1.3.6.1.4.1.12345.2.2.0",
"dewpoint": "1.3.6.1.4.1.12345.2.3.0",
"alarm": "1.3.6.1.4.1.12345.2.6.0",
})
try:
result = await self._snmp_get_bulk(ip, community, oids)
if result:
result["poll_ts"] = time.time()
self.cache[ip] = result
except asyncio.TimeoutError:
self.stats["timeouts"] += 1
self.stats["failures"] += 1
except Exception:
self.stats["failures"] += 1
async def _snmp_get_bulk(self, ip, community, oids):
loop = asyncio.get_event_loop()
# 使用线程池执行同步pysnmp调用
def sync_get():
from pysnmp.hlapi import getCmd, CommunityData, UdpTransportTarget, ContextData, ObjectType, ObjectIdentity
results = {}
for name, oid in oids.items():
iterator = getCmd(
CommunityData(community),
UdpTransportTarget((ip, 161), timeout=self.timeout, retries=self.retries),
ContextData(),
ObjectType(ObjectIdentity(oid))
)
errorIndication, errorStatus, errorIndex, varBinds = next(iterator)
if not errorIndication and not errorStatus:
for varBind in varBinds:
results[name] = int(varBind[1])
return results
return await loop.run_in_executor(None, sync_get)
def get_latest(self, ip: str) -> Optional[dict]:
return self.cache.get(ip)
def stop(self):
self._running = False"""
fusion.py - 数据融合与链路选择
"""
import time
from enum import Enum
from dataclasses import dataclass, field
class LinkState(Enum):
UDP_ONLY = "udp_only"
SNMP_ONLY = "snmp_only"
BOTH_OK = "both_ok"
BOTH_DOWN = "both_down"
@dataclass
class FusedData:
dev_id: str
temperature: float
humidity: float
dewpoint: float
timestamp: float
source: str
link_state: LinkState
quality: str
udp_age: float = 0.0
snmp_age: float = 0.0
divergence: float = 0.0
class FusionEngine:
def __init__(self, udp_receiver, snmp_engine, max_age=10.0, divergence_threshold=0.5):
self.udp = udp_receiver
self.snmp = snmp_engine
self.max_age = max_age
self.divergence_threshold = divergence_threshold
self.fused_cache: dict = {}
def fuse(self, dev_id: str, dev_ip: str = None) -> FusedData:
now = time.time()
udp_data = self.udp.get_latest(dev_id)
snmp_data = self.snmp.get_latest(dev_ip) if dev_ip else None
udp_age = now - udp_data.recv_ts if udp_data else float('inf')
snmp_age = now - snmp_data["poll_ts"] if snmp_data else float('inf')
udp_valid = udp_data and udp_age <= self.max_age
snmp_valid = snmp_data and snmp_age <= self.max_age
if udp_valid and snmp_valid:
div = max(
abs(udp_data.temperature - snmp_data.get("temperature", 0) / 10.0),
abs(udp_data.humidity - snmp_data.get("humidity", 0) / 10.0)
)
if div > self.divergence_threshold:
# 数据不一致,优先UDP,标记告警
return self._build(udp_data, "udp", LinkState.BOTH_OK, "divergence", udp_age, snmp_age, div)
return self._build(udp_data, "fused", LinkState.BOTH_OK, "good", udp_age, snmp_age, div)
elif udp_valid:
return self._build(udp_data, "udp", LinkState.UDP_ONLY, "good", udp_age, snmp_age)
elif snmp_valid:
return self._build_snmp(dev_ip, snmp_data, snmp_age)
else:
return FusedData(
dev_id=dev_id, temperature=None, humidity=None, dewpoint=None,
timestamp=now, source="none", link_state=LinkState.BOTH_DOWN,
quality="invalid"
)
def _build(self, udp_data, source, link_state, quality, udp_age, snmp_age, div=0.0):
return FusedData(
dev_id=udp_data.dev_id,
temperature=udp_data.temperature,
humidity=udp_data.humidity,
dewpoint=udp_data.dewpoint,
timestamp=udp_data.recv_ts,
source=source,
link_state=link_state,
quality=quality,
udp_age=udp_age,
snmp_age=snmp_age,
divergence=div
)
def _build_snmp(self, dev_ip, snmp_data, snmp_age):
return FusedData(
dev_id=dev_ip,
temperature=snmp_data.get("temperature", 0) / 10.0,
humidity=snmp_data.get("humidity", 0) / 10.0,
dewpoint=snmp_data.get("dewpoint", 0) / 10.0,
timestamp=snmp_data.get("poll_ts", time.time()),
source="snmp",
link_state=LinkState.SNMP_ONLY,
quality="good",
snmp_age=snmp_age
)"""
modbus_server.py - Modbus TCP 服务
"""
from pymodbus.server import StartTcpServer
from pymodbus.datastore import ModbusSlaveContext, ModbusServerContext
from pymodbus.datastore import ModbusSequentialDataBlock
import threading
class ModbusService:
def __init__(self, fusion_engine, port=5020):
self.fusion = fusion_engine
self.port = port
self.store = ModbusSlaveContext(
di=ModbusSequentialDataBlock(0, [0]*1000),
co=ModbusSequentialDataBlock(0, [0]*1000),
hr=ModbusSequentialDataBlock(0, [0]*1000),
ir=ModbusSequentialDataBlock(0, [0]*1000)
)
self.context = ModbusServerContext(slaves=self.store, single=True)
def update_register(self, dev_id, fused_data):
"""将融合数据写入Modbus寄存器"""
base_addr = int(dev_id.split('-')[1]) * 10 if '-' in dev_id else 0
self.store.setValues(3, base_addr, [int(fused_data.temperature * 10)])
self.store.setValues(3, base_addr + 1, [int(fused_data.humidity * 10)])
self.store.setValues(3, base_addr + 2, [int(fused_data.dewpoint * 10)])
self.store.setValues(3, base_addr + 3, [1 if fused_data.quality == "good" else 0])
def start(self):
thread = threading.Thread(target=StartTcpServer, args=(self.context,), kwargs={"address":("0.0.0.0", self.port)})
thread.daemon = True
thread.start()
print(f"[Modbus] Server started on port {self.port}")"""
mqtt_publisher.py - MQTT发布
"""
import paho.mqtt.client as mqtt
import json
import time
class MQTTPublisher:
def __init__(self, broker="mqtt.local", port=1883, topic_prefix="substation/env"):
self.broker = broker
self.port = port
self.topic_prefix = topic_prefix
self.client = mqtt.Client()
self.client.on_connect = self._on_connect
self.connected = False
def _on_connect(self, client, userdata, flags, rc):
if rc == 0:
self.connected = True
print("[MQTT] Connected")
else:
print(f"[MQTT] Connect failed: {rc}")
def connect(self):
self.client.connect(self.broker, self.port, 60)
self.client.loop_start()
def publish(self, dev_id, fused_data):
if not self.connected:
return
payload = {
"dev_id": dev_id,
"temperature": fused_data.temperature,
"humidity": fused_data.humidity,
"dewpoint": fused_data.dewpoint,
"source": fused_data.source,
"link_state": fused_data.link_state.value,
"quality": fused_data.quality,
"ts": int(fused_data.timestamp * 1000)
}
topic = f"{self.topic_prefix}/{dev_id}"
self.client.publish(topic, json.dumps(payload), qos=1)电厂变电站机房的以太网温湿度采集系统,采用UDP上报+SNMP轮询双协议架构,核心设计要点:网络分层隔离(符合电力安防)、网关双机热备、数据融合引擎自动选源、北向多协议适配(Modbus TCP/MQTT/IEC104)。UDP提供实时性,SNMP提供可靠性,两者互为补充。部署时需特别注意电力二次系统安全防护规定,确保采集网络与SCADA网络的安全隔离。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。