
关键词:楼宇自控、IoT网关、POE温湿度变送器、UDP、SNMP、双链路冗余、链路切换、数据一致性、BAS集成、高可用 标签:#物联网 #Modbus #TCP/IP #UDP #POE供电 #Python #以太网温湿度传感器 #网口温湿度变送器 #机房监控 #楼宇自控 #SNMP
楼宇自控系统(BAS)对温湿度数据的连续性要求,介于机房动环和电力监控之间:机房可以容忍短暂中断(有备用采集),电力监控要求毫秒级联动,而楼控要求分钟级连续趋势数据,用于能耗分析、舒适度评估和调度优化。
POE温湿度变送器在楼控场景中,通常采用单一通信方式:要么UDP主动上报,要么SNMP轮询。两种方式各有致命弱点:
通信方式 | 优点 | 致命弱点 |
|---|---|---|
UDP主动上报 | 实时性好,设备端简单,网络负载低 | 丢包不可知,设备异常时数据静默丢失 |
SNMP轮询 | 可靠性高,支持Trap告警,可获取设备状态 | 轮询间隔受限,网络拥塞时超时,设备离线无法区分 |
双链路冗余的核心思路:同时启用UDP和SNMP两条独立通信路径,互为备份,任意一条链路故障不影响数据采集连续性。
┌─────────────────────────────────────────────────────────────┐
│ 楼宇自控IoT网关 │
├─────────────────────────────────────────────────────────────┤
│ │
│ ┌──────────────────┐ ┌──────────────────┐ │
│ │ UDP接收线程 │ │ SNMP轮询线程 │ │
│ │ (端口 9000) │ │ (每30秒) │ │
│ │ │ │ │ │
│ │ - 接收UDP数据报 │ │ - GET温度/湿度 │ │
│ │ - 解析二进制帧 │ │ - GET设备状态 │ │
│ │ - 时间戳标记 │ │ - 超时重试 │ │
│ └────────┬─────────┘ └────────┬──────────┘ │
│ │ │ │
│ ▼ ▼ │
│ ┌──────────────────────────────────────────────────┐ │
│ │ 数据融合与链路选择引擎 │ │
│ │ │ │
│ │ - 时间戳对齐(±2s窗口内两条链路数据匹配) │ │
│ │ - 质量标记(primary/secondary/stale/invalid) │ │
│ │ - 链路健康度评估(丢包率、延迟、错误率) │ │
│ │ - 自动切换策略(主链路故障→备用链路) │ │
│ │ - 数据一致性校验(两条链路数据偏差超过阈值告警) │ │
│ └───────────────────────┬──────────────────────────┘ │
│ │ │
│ ▼ │
│ ┌──────────────────────────────────────────────────┐ │
│ │ 北向接口层 │ │
│ │ │ │
│ │ - BACnet/IP 服务端(对象 present-value 更新) │ │
│ │ - MQTT 发布(可选,用于云平台) │ │
│ │ - REST API(本地调试) │ │
│ └──────────────────────────────────────────────────┘ │
│ │
└─────────────────────────────────────────────────────────────┘
│
┌─────────────┼─────────────┐
│ │
┌───────▼───────┐ ┌────────▼────────┐
│ POE交换机 │ │ POE交换机 │
│ (VLAN 200) │ │ (VLAN 200) │
│ UDP链路 │ │ SNMP链路 │
└───────┬───────┘ └────────┬────────┘
│ │
└──────────┬────────────────┘
│
┌──────────▼──────────┐
│ POE温湿度变送器 │
│ (双协议固件) │
│ IP: 192.168.200.x │
└─────────────────────┘双链路冗余的前提是两条链路在网络层面真正独立,否则交换机故障会导致两条链路同时中断:
隔离维度 | 方案 |
|---|---|
VLAN隔离 | UDP和SNMP走不同VLAN(如VLAN 200/201),通过不同逻辑接口 |
物理路径 | 两台独立POE交换机,分别上联到网关的不同网卡 |
协议端口 | UDP: 9000/UDP,SNMP: 161/UDP(不同目的端口,可走不同QoS队列) |
源IP | 网关配置两个IP地址,分别用于UDP接收和SNMP轮询 |
网关双网卡配置:
eth0: 192.168.200.10/24 (VLAN 200, UDP链路)
eth1: 192.168.200.11/24 (VLAN 201, SNMP链路)
变送器配置:
UDP目标IP: 192.168.200.10, 端口 9000
SNMP Trap目标: 192.168.200.11 (或SNMP GET源IP)"""
udp_receiver.py - UDP数据接收服务
"""
import asyncio
import struct
from dataclasses import dataclass
from datetime import datetime
@dataclass
class UDPReading:
dev_id: str
temperature: float
humidity: float
dewpoint: float
recv_ts: float
sequence: int
rssi: int = 0 # 部分设备支持信号强度
quality: str = "good"
class UDPReceiver:
def __init__(self, host="0.0.0.0", port=9000, device_registry=None):
self.host = host
self.port = port
self.device_registry = device_registry # IP→设备ID映射
self.cache = {} # dev_id → UDPReading
self.stats = {"packets_recv": 0, "packets_parse_err": 0, "bytes_recv": 0}
self.transport = None
def start(self, loop):
"""启动UDP接收"""
listen = loop.create_datagram_endpoint(
lambda: self,
local_addr=(self.host, self.port)
)
self.transport, _ = loop.run_until_complete(listen)
print(f"UDP Receiver listening on {self.host}:{self.port}")
def datagram_received(self, data, addr):
"""接收UDP数据报"""
self.stats["packets_recv"] += 1
self.stats["bytes_recv"] += len(data)
try:
reading = self._parse_packet(data, addr)
if reading:
self.cache[reading.dev_id] = reading
except Exception as e:
self.stats["packets_parse_err"] += 1
print(f"UDP parse error from {addr}: {e}")
def _parse_packet(self, data, addr):
"""解析设备UDP数据包(示例格式,需按厂家协议调整)"""
# 假设格式:包头(2B) + 设备ID(4B) + 温度(2B,×0.1) + 湿度(2B,×0.1) + 露点(2B,×0.1) + 序列号(2B) + CRC(2B)
if len(data) < 14:
raise ValueError("Packet too short")
header, dev_id_int = struct.unpack_from(">HH", data, 0)
if header != 0xAA55:
raise ValueError(f"Invalid header: {header:#x}")
temp_raw, hum_raw, dp_raw, seq = struct.unpack_from(">hhhH", data, 4)
# CRC校验
crc_received = struct.unpack_from(">H", data, len(data)-2)[0]
crc_calc = self._crc16(data[:-2])
if crc_received != crc_calc:
raise ValueError("CRC mismatch")
dev_id = self.device_registry.get(addr[0], f"dev_{dev_id_int}")
return UDPReading(
dev_id=dev_id,
temperature=temp_raw / 10.0,
humidity=hum_raw / 10.0,
dewpoint=dp_raw / 10.0,
recv_ts=asyncio.get_event_loop().time(),
sequence=seq
)
def _crc16(self, data):
"""CRC-16/MODBUS"""
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 get_latest(self, dev_id):
"""获取最新UDP数据"""
return self.cache.get(dev_id)
def connection_made(self, transport):
self.transport = transport
def error_received(self, exc):
print(f"UDP error: {exc}")"""
udp_health.py - UDP链路健康度监控
"""
import time
from collections import deque
class UDPHealthMonitor:
def __init__(self, window_size=100, expected_interval=5.0, jitter_threshold=2.0):
self.window_size = window_size
self.expected_interval = expected_interval
self.jitter_threshold = jitter_threshold
self.dev_stats = {}
def on_packet(self, dev_id, recv_ts, sequence):
"""记录收到的包,计算健康指标"""
stats = self.dev_stats.setdefault(dev_id, {
"recv_times": deque(maxlen=self.window_size),
"sequences": deque(maxlen=self.window_size),
"packet_loss": 0,
"last_seq": None,
})
stats["recv_times"].append(recv_ts)
stats["sequences"].append(sequence)
# 检测丢包(序列号不连续)
if stats["last_seq"] is not None:
expected = (stats["last_seq"] + 1) % 65536
if sequence != expected:
stats["packet_loss"] += 1
stats["last_seq"] = sequence
def get_health(self, dev_id):
"""返回链路健康度评分(0~100,100为最佳)"""
stats = self.dev_stats.get(dev_id)
if not stats or len(stats["recv_times"]) < 2:
return 0
recv_times = list(stats["recv_times"])
intervals = [recv_times[i] - recv_times[i-1] for i in range(1, len(recv_times))]
# 丢包率
loss_rate = stats["packet_loss"] / max(len(stats["sequences"]), 1)
# 抖动(间隔标准差)
import statistics
jitter = statistics.stdev(intervals) if len(intervals) > 1 else 0
# 延迟(平均间隔与期望的偏差)
avg_interval = sum(intervals) / len(intervals)
delay_deviation = abs(avg_interval - self.expected_interval) / self.expected_interval
# 综合评分
score = 100
score -= loss_rate * 50 # 丢包率权重 50
score -= min(jitter / self.jitter_threshold, 1) * 30 # 抖动权重 30
score -= min(delay_deviation, 1) * 20 # 延迟偏差权重 20
return max(0, min(100, score))"""
snmp_poller.py - SNMP轮询服务(与之前文章类似,此处聚焦冗余设计)
"""
from pysnmp.hlapi import *
from pysnmp.entity import config
import asyncio
class SNMPPoller:
def __init__(self, devices, poll_interval=30, oid_base=".1.3.6.1.4.1.12345.1"):
self.devices = devices # list of {ip, community, version, ...}
self.poll_interval = poll_interval
self.oid_base = oid_base
self.cache = {}
self.stats = {"polls_total": 0, "polls_failed": 0, "timeouts": 0}
async def poll_all(self):
"""轮询所有设备"""
while True:
for dev in self.devices:
asyncio.create_task(self._poll_single(dev))
await asyncio.sleep(self.poll_interval)
async def _poll_single(self, dev):
"""轮询单台设备"""
self.stats["polls_total"] += 1
try:
result = await self._snmp_get(dev)
if result:
self.cache[dev["ip"]] = {
"temperature": result.get("temperature"),
"humidity": result.get("humidity"),
"dewpoint": result.get("dewpoint"),
"alarm_status": result.get("alarm_status"),
"poll_ts": asyncio.get_event_loop().time(),
"quality": "good"
}
except asyncio.TimeoutError:
self.stats["timeouts"] += 1
self.stats["polls_failed"] += 1
if dev["ip"] in self.cache:
self.cache[dev["ip"]]["quality"] = "stale"
except Exception as e:
self.stats["polls_failed"] += 1
print(f"SNMP poll {dev['ip']} failed: {e}")
async def _snmp_get(self, dev):
"""执行SNMP GET"""
# 实现与之前文章类似,此处省略具体pysnmp调用
pass
def get_latest(self, dev_ip):
"""获取最新SNMP数据"""
return self.cache.get(dev_ip)"""
snmp_health.py - SNMP链路健康度监控
"""
import time
class SNMPHealthMonitor:
def __init__(self, timeout_threshold=3.0, consecutive_fail_threshold=3):
self.timeout_threshold = timeout_threshold
self.consecutive_fail_threshold = consecutive_fail_threshold
self.dev_stats = {}
def on_poll_result(self, dev_id, success, response_time_ms):
"""记录轮询结果"""
stats = self.dev_stats.setdefault(dev_id, {
"consecutive_fails": 0,
"total_polls": 0,
"failed_polls": 0,
"response_times": [],
"last_success_ts": None
})
stats["total_polls"] += 1
if success:
stats["consecutive_fails"] = 0
stats["last_success_ts"] = time.time()
stats["response_times"].append(response_time_ms)
if len(stats["response_times"]) > 100:
stats["response_times"].pop(0)
else:
stats["consecutive_fails"] += 1
stats["failed_polls"] += 1
def get_health(self, dev_id):
"""返回链路健康度评分(0~100)"""
stats = self.dev_stats.get(dev_id)
if not stats or stats["total_polls"] == 0:
return 0
# 连续失败次数
if stats["consecutive_fails"] >= self.consecutive_fail_threshold:
return 0
# 失败率
fail_rate = stats["failed_polls"] / stats["total_polls"]
# 响应时间(如果可用)
if stats["response_times"]:
avg_rtt = sum(stats["response_times"]) / len(stats["response_times"])
rtt_score = max(0, 1 - avg_rtt / (self.timeout_threshold * 1000))
else:
rtt_score = 0.5
score = 100
score -= fail_rate * 60
score -= (1 - rtt_score) * 40
return max(0, min(100, score))"""
fusion_engine.py - 双链路数据融合与切换
"""
import time
from enum import Enum
from dataclasses import dataclass
class LinkStatus(Enum):
PRIMARY = "primary" # 主链路正常
DEGRADED = "degraded" # 主链路降级,备用链路正常
FAILOVER = "failover" # 主链路故障,使用备用链路
BOTH_DOWN = "both_down" # 两条链路都故障
@dataclass
class FusedReading:
dev_id: str
temperature: float
humidity: float
dewpoint: float
timestamp: float
source: str # "udp" / "snmp" / "fused"
link_status: LinkStatus
quality: str
udp_health: float
snmp_health: float
divergence: float = 0.0 # 两条链路数据偏差
class FusionEngine:
def __init__(self, udp_receiver, snmp_poller,
udp_health_monitor, snmp_health_monitor,
divergence_threshold=0.5):
self.udp = udp_receiver
self.snmp = snmp_poller
self.udp_health = udp_health_monitor
self.snmp_health = snmp_health_monitor
self.divergence_threshold = divergence_threshold
self.fused_cache = {}
def get_reading(self, dev_id, dev_ip=None):
"""获取融合后的数据"""
udp_data = self.udp.get_latest(dev_id)
snmp_data = self.snmp.get_latest(dev_ip) if dev_ip else None
udp_health = self.udp_health.get_health(dev_id)
snmp_health = self.snmp_health.get_health(dev_ip) if dev_ip else 0
# 判断链路状态
link_status = self._determine_link_status(udp_data, snmp_data, udp_health, snmp_health)
# 数据融合
if link_status == LinkStatus.PRIMARY:
# 主链路正常,使用UDP数据
reading = self._build_reading(dev_id, udp_data, "udp", link_status, udp_health, snmp_health)
elif link_status == LinkStatus.FAILOVER:
# 主链路故障,使用SNMP数据
reading = self._build_reading(dev_id, snmp_data, "snmp", link_status, udp_health, snmp_health)
elif link_status == LinkStatus.DEGRADED:
# 两条链路都可用,检查一致性
divergence = self._calc_divergence(udp_data, snmp_data)
if divergence > self.divergence_threshold:
# 数据不一致,标记质量异常,优先使用UDP
reading = self._build_reading(dev_id, udp_data, "udp", link_status, udp_health, snmp_health)
reading.divergence = divergence
reading.quality = "divergence_alert"
else:
# 数据一致,使用UDP
reading = self._build_reading(dev_id, udp_data, "fused", link_status, udp_health, snmp_health)
reading.divergence = divergence
else:
# 两条链路都故障
reading = self._build_reading(dev_id, None, "none", link_status, udp_health, snmp_health)
reading.quality = "invalid"
self.fused_cache[dev_id] = reading
return reading
def _determine_link_status(self, udp_data, snmp_data, udp_health, snmp_health):
"""判断链路状态"""
udp_ok = udp_data is not None and udp_health > 50
snmp_ok = snmp_data is not None and snmp_health > 50
if udp_ok and snmp_ok:
return LinkStatus.DEGRADED if udp_health < 80 else LinkStatus.PRIMARY
elif udp_ok and not snmp_ok:
return LinkStatus.PRIMARY
elif not udp_ok and snmp_ok:
return LinkStatus.FAILOVER
else:
return LinkStatus.BOTH_DOWN
def _build_reading(self, dev_id, data, source, link_status, udp_health, snmp_health):
"""构建融合读数"""
if data is None:
return FusedReading(
dev_id=dev_id,
temperature=None,
humidity=None,
dewpoint=None,
timestamp=time.time(),
source=source,
link_status=link_status,
quality="invalid",
udp_health=udp_health,
snmp_health=snmp_health
)
return FusedReading(
dev_id=dev_id,
temperature=data.temperature if hasattr(data, 'temperature') else data.get('temperature'),
humidity=data.humidity if hasattr(data, 'humidity') else data.get('humidity'),
dewpoint=data.dewpoint if hasattr(data, 'dewpoint') else data.get('dewpoint'),
timestamp=data.recv_ts if hasattr(data, 'recv_ts') else data.get('poll_ts', time.time()),
source=source,
link_status=link_status,
quality="good",
udp_health=udp_health,
snmp_health=snmp_health
)
def _calc_divergence(self, udp_data, snmp_data):
"""计算两条链路数据的偏差"""
if not udp_data or not snmp_data:
return 0.0
udp_temp = udp_data.temperature if hasattr(udp_data, 'temperature') else udp_data.get('temperature', 0)
snmp_temp = snmp_data.get('temperature', 0) if isinstance(snmp_data, dict) else snmp_data.temperature
udp_hum = udp_data.humidity if hasattr(udp_data, 'humidity') else udp_data.get('humidity', 0)
snmp_hum = snmp_data.get('humidity', 0) if isinstance(snmp_data, dict) else snmp_data.humidity
temp_diff = abs(udp_temp - snmp_temp)
hum_diff = abs(udp_hum - snmp_hum)
return max(temp_diff, hum_diff)链路切换决策树:
1. UDP健康度 > 80 且 SNMP健康度 > 50
→ 主链路UDP,SNMP作为热备
→ 数据一致性检查,偏差<阈值时标记"fused"
2. UDP健康度降至 50~80
→ 仍然使用UDP,但标记"degraded"
→ 增加SNMP轮询频率(从30s降到10s)
3. UDP健康度 < 50 或 UDP数据超时(>2倍上报周期未收到)
→ 切换到SNMP链路
→ 标记"failover"
→ 触发告警:UDP链路异常,已切换至SNMP
4. SNMP也故障(连续3次超时)
→ 标记"both_down"
→ 触发紧急告警
→ 尝试重启UDP接收服务
→ 记录最后有效数据的时间戳
5. UDP恢复(连续收到5个有效包)
→ 切回UDP主链路
→ 标记"primary_restored"
→ 触发恢复通知
融合引擎输出的数据,需要更新到BACnet对象的present-value中:
"""
bacnet_updater.py - 更新BACnet对象
"""
class BACnetUpdater:
def __init__(self, fusion_engine, bacnet_app):
self.fusion = fusion_engine
self.app = bacnet_app
self.obj_map = {} # dev_id → {temp_obj, hum_obj, dp_obj, status_obj}
def update_all(self):
"""更新所有设备的BACnet对象"""
for dev_id, objs in self.obj_map.items():
reading = self.fusion.get_reading(dev_id)
if reading.quality == "good":
objs["temp"].presentValue = reading.temperature
objs["hum"].presentValue = reading.humidity
objs["dp"].presentValue = reading.dewpoint
objs["status"].presentValue = reading.link_status.value
elif reading.quality == "stale":
# SNMP数据过期,但UDP正常
objs["temp"].reliability = "UNRELIABLE_OTHER"
elif reading.quality == "invalid":
# 两条链路都故障
objs["temp"].reliability = "FAULT"
objs["status"].presentValue = "COMM_FAILURE"楼宇自控IoT网关接入POE温湿度变送器时,UDP和SNMP双链路冗余设计的核心价值在于:UDP提供实时性,SNMP提供可靠性,两者互为备份,确保数据连续性。设计要点:网络路径物理隔离、链路健康度量化评估、数据融合与自动切换、北向接口正确标记数据质量。融合引擎是核心,它不仅要判断链路状态,还要检测数据一致性,防止错误数据进入BAS系统。双链路冗余不是简单的"两条路都走",而是有主备、有切换、有校验的完整高可用方案。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。