首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >楼宇自控IoT网关:POE温湿度变送器接入,UDP和SNMP双链路冗余设计

楼宇自控IoT网关:POE温湿度变送器接入,UDP和SNMP双链路冗余设计

原创
作者头像
HONSOR盛世宏博
发布于 2026-09-24 17:11:14
发布于 2026-09-24 17:11:14
870
举报

楼宇自控IoT网关:POE温湿度变送器接入,UDP和SNMP双链路冗余设计

关键词:楼宇自控、IoT网关、POE温湿度变送器、UDP、SNMP、双链路冗余、链路切换、数据一致性、BAS集成、高可用 标签:#物联网 #Modbus #TCP/IP #UDP #POE供电 #Python #以太网温湿度传感器 #网口温湿度变送器 #机房监控 #楼宇自控 #SNMP

一、为什么需要双链路冗余

楼宇自控系统(BAS)对温湿度数据的连续性要求,介于机房动环和电力监控之间:机房可以容忍短暂中断(有备用采集),电力监控要求毫秒级联动,而楼控要求分钟级连续趋势数据,用于能耗分析、舒适度评估和调度优化。

POE温湿度变送器在楼控场景中,通常采用单一通信方式:要么UDP主动上报,要么SNMP轮询。两种方式各有致命弱点:

通信方式

优点

致命弱点

UDP主动上报​

实时性好,设备端简单,网络负载低

丢包不可知,设备异常时数据静默丢失

SNMP轮询​

可靠性高,支持Trap告警,可获取设备状态

轮询间隔受限,网络拥塞时超时,设备离线无法区分

双链路冗余的核心思路:同时启用UDP和SNMP两条独立通信路径,互为备份,任意一条链路故障不影响数据采集连续性。


二、架构设计

代码语言:javascript
复制
┌─────────────────────────────────────────────────────────────┐
│                   楼宇自控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  │
            └─────────────────────┘

2.1 网络路径隔离

双链路冗余的前提是两条链路在网络层面真正独立,否则交换机故障会导致两条链路同时中断:

隔离维度

方案

VLAN隔离​

UDP和SNMP走不同VLAN(如VLAN 200/201),通过不同逻辑接口

物理路径​

两台独立POE交换机,分别上联到网关的不同网卡

协议端口​

UDP: 9000/UDP,SNMP: 161/UDP(不同目的端口,可走不同QoS队列)

源IP​

网关配置两个IP地址,分别用于UDP接收和SNMP轮询

代码语言:javascript
复制
网关双网卡配置:
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链路实现

3.1 UDP接收服务

代码语言:javascript
复制
"""
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}")

3.2 UDP链路健康度评估

代码语言:javascript
复制
"""
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链路实现

4.1 SNMP轮询服务

代码语言:javascript
复制
"""
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)

4.2 SNMP链路健康度评估

代码语言:javascript
复制
"""
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))

五、数据融合与链路切换引擎

5.1 融合策略

代码语言:javascript
复制
"""
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)

5.2 切换逻辑

代码语言:javascript
复制
链路切换决策树:

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"
   → 触发恢复通知

六、北向接口

6.1 BACnet/IP对象更新

融合引擎输出的数据,需要更新到BACnet对象的present-value中:

代码语言:javascript
复制
"""
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"

七、典型坑

  1. 两条链路走同一交换机:交换机故障导致双链路同时中断,冗余形同虚设。解决:物理路径分离,两台独立POE交换机。
  2. UDP和SNMP数据时间戳不对齐:UDP是设备主动上报,SNMP是网关轮询,时间差可能达数秒。解决:融合引擎使用时间窗口匹配(±2s)。
  3. SNMP轮询频率过高:大量设备频繁SNMP GET,网络拥塞,反而导致UDP丢包。解决:SNMP轮询间隔合理设置(30~60s),UDP为主。
  4. 链路切换震荡:UDP短暂抖动导致频繁切换。解决:切换需满足持续条件(连续5个包正常才切回),加防抖时间。
  5. 数据一致性误判:两条链路数据偏差超过阈值,但实际是传感器本身漂移。解决:区分链路故障和数据漂移,漂移由校准流程处理。
  6. SNMP超时设置过短:网络拥塞时大量SNMP超时,误判链路故障。解决:超时设置3~5s,重试2~3次。
  7. 融合引擎成为单点故障:网关进程崩溃,双链路数据无法输出。解决:网关进程systemd守护,自动重启;或部署双网关热备。
  8. BACnet对象reliability未正确设置:链路故障时present-value仍为旧值,BAS系统误以为数据正常。解决:链路故障时设置reliability=FAULT。
  9. UDP端口冲突:多个网关实例监听同一UDP端口,导致数据重复或丢失。解决:每个网关实例使用不同端口,或通过负载均衡分发。
  10. 忘记配置SNMP Trap:SNMP链路只做轮询,未配置Trap,设备异常时无法快速感知。解决:同时启用SNMP Trap接收,作为轮询的补充。

八、小结

楼宇自控IoT网关接入POE温湿度变送器时,UDP和SNMP双链路冗余设计的核心价值在于:UDP提供实时性,SNMP提供可靠性,两者互为备份,确保数据连续性。设计要点:网络路径物理隔离、链路健康度量化评估、数据融合与自动切换、北向接口正确标记数据质量。融合引擎是核心,它不仅要判断链路状态,还要检测数据一致性,防止错误数据进入BAS系统。双链路冗余不是简单的"两条路都走",而是有主备、有切换、有校验的完整高可用方案。

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

目录
  • 楼宇自控IoT网关:POE温湿度变送器接入,UDP和SNMP双链路冗余设计
    • 一、为什么需要双链路冗余
    • 二、架构设计
      • 2.1 网络路径隔离
    • 三、UDP链路实现
      • 3.1 UDP接收服务
      • 3.2 UDP链路健康度评估
    • 四、SNMP链路实现
      • 4.1 SNMP轮询服务
      • 4.2 SNMP链路健康度评估
    • 五、数据融合与链路切换引擎
      • 5.1 融合策略
      • 5.2 切换逻辑
    • 六、北向接口
      • 6.1 BACnet/IP对象更新
    • 七、典型坑
    • 八、小结
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档