首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >算力服务器机房改造:UDP多播方式,批量采集RJ45温湿度变送器数据

算力服务器机房改造:UDP多播方式,批量采集RJ45温湿度变送器数据

原创
作者头像
盛世宏博科技
发布于 2026-09-24 17:08:47
发布于 2026-09-24 17:08:47
750
举报

算力服务器机房改造:UDP多播方式,批量采集RJ45温湿度变送器数据

近期在调试一套工业环境监控系统时,遇到了关于Modbus连接保持的问题。现场部署了一批网口温湿度变送器,PoE取电、Modbus TCP上云,同时利旧接入存量RS485探头,并开启了SNMP和UDP Trap服务。算力服务器机房改造项目里,节点规模从百级跳到千级,传统"一问一答"的Modbus TCP轮询模式在扇出和时延上开始吃力:400节点×5s轮询 ≈ 80 RPS,勉强够用;1200节点×3s轮询 ≈ 400 RPS,单采集服务已经到瓶颈。这篇讲一种工程上已经验证的替代路径——UDP多播批量采集,把"服务端逐个轮询"变成"设备主动批量上报、服务端一组收",在算力机房高密度场景下显著降低采集面负载。


一、为什么算力机房需要UDP多播

1.1 规模带来的矛盾

维度

传统Modbus TCP轮询

UDP多播批量采集

采集模式

服务端主动轮询,一问一答

设备主动推送,一组多收

网络负载

N节点 × M寄存器 × 轮询频率

设备侧定时组播,单次发包覆盖所有订阅者

TCP连接数

每节点1-N个PCB

无连接,零PCB

服务端扇出

O(N) 并发连接

O(1) 监听一个多播组

设备侧资源

TCP状态机 + pbuf池

仅pbuf发送,无状态

数据实时性

取决于轮询周期

设备侧定时推送,周期可控

可靠性

TCP保证

不保证,需应用层去重+Modbus兜底

核心收益:服务端从"N个TCP连接 × 轮询频率"降为"1个UDP socket监听",设备侧从"维持TCP PCB"降为"定时组播发包",在千节点规模下,采集服务的CPU和内存占用下降一个数量级。

1.2 适用边界

  • 主数据面:UDP多播作为高频轻量遥测通道,Modbus TCP仍作为可靠兜底和配置下发通道
  • 网络要求:局域网内多播路由可达(PIM-SM或IGMP Snooping),跨子网需路由支持
  • 节点密度:百节点以上规模才有明显收益,小规模场景轮询足够

二、多播架构设计

2.1 整体拓扑

代码语言:javascript
复制
┌─────────────────────────────┐
                    │      采集服务端              │
                    │  UDP多播监听 (239.255.1.100) │
                    │       ↓                     │
                    │  解析 → 去重 → 时序库        │
                    └──────────────┬──────────────┘
                                   │ 多播组 239.255.1.100:9000
          ┌────────────────────────┼────────────────────────┐
          │                        │                        │
   [列1: 变送器×48]        [列2: 变送器×48]        [列3: 变送器×48]
   定时组播推送           定时组播推送           定时组播推送
   (每3秒)               (每3秒)               (每3秒)

2.2 多播地址规划

多播组

用途

TTL

239.255.1.100

周期遥测数据

1(本子网)

239.255.1.101

事件Trap(越限/掉电)

1

239.255.1.102

管理面心跳(设备在线检测)

1

239.255.1.200

服务端→设备广播指令(可选)

1

地址选择原则:

  • 使用239.0.0.0/8(管理范围多播),不污染公网
  • TTL=1确保不出子网,跨子网时由PIM路由器按需转发
  • 不同业务用不同组,接收端按组订阅,避免无关数据唤醒

三、设备侧多播推送实现

3.1 LWIP多播初始化

代码语言:javascript
复制
// multicast.c - 基于LWIP raw API
#include "lwip/udp.h"
#include "lwip/igmp.h"
#include "lwip/ip_addr.h"
#include "lwip/netif.h"
#include <string.h>

static struct udp_pcb *g_mcast_pcb = NULL;
static ip_addr_t g_mcast_addr;
static u16_t g_mcast_port = 9000;
static uint16_t g_seq = 0;

#define MULTICAST_TTL  1   // 不出子网

int multicast_init(const char *mcast_ip, u16_t port) {
    g_mcast_pcb = udp_new();
    if (!g_mcast_pcb) return -1;
    
    ipaddr_aton(mcast_ip, &g_mcast_addr);
    g_mcast_port = port;
    
    // 绑定本地端口(0=自动分配)
    udp_bind(g_mcast_pcb, IP_ADDR_ANY, 0);
    
    // 加入多播组
    struct netif *netif = netif_default;
    if (!netif) return -2;
    
    err_t err = igmp_joingroup(&netif->ip_addr, &g_mcast_addr);
    if (err != ERR_OK) return -3;
    
    // 设置TTL
    udp_set_multicast_ttl(g_mcast_pcb, MULTICAST_TTL);
    
    return 0;
}

void multicast_deinit(void) {
    if (g_mcast_pcb) {
        igmp_leavegroup(IP_ADDR_ANY, &g_mcast_addr);
        udp_remove(g_mcast_pcb);
        g_mcast_pcb = NULL;
    }
}

3.2 批量数据帧格式

代码语言:javascript
复制
// 多播批量上报帧:一帧携带多个节点的数据
// 用于高密度机柜场景,一台设备可接多个探头或自身多传感器
typedef struct __attribute__((packed)) {
    uint16_t header;       // 0x55BB 多播帧标识
    uint8_t  dev_addr;     // 设备地址
    uint8_t  sensor_count; // 本帧包含的传感器数量(1-8)
    uint16_t seq;          // 序列号
    uint32_t batch_ts;     // 批次时间戳
    uint16_t interval_ms;  // 推送间隔(毫秒),用于服务端计算丢包
} mcast_frame_hdr_t;       // 12字节

typedef struct __attribute__((packed)) {
    uint8_t  sensor_id;    // 传感器编号(0=主传感器)
    int16_t  temperature;  // 温度 ×100
    uint16_t humidity;     // 湿度 ×10
    uint16_t voltage;      // 供电电压 ×100
    uint8_t  status;       // 状态位
} mcast_sensor_t;          // 8字节

// 最大帧大小:12 + 8×8 = 76字节(远小于以太网MTU 1500)
#define MAX_SENSORS_PER_FRAME  8
#define MCAST_FRAME_MAX_SIZE   (sizeof(mcast_frame_hdr_t) + MAX_SENSORS_PER_FRAME * sizeof(mcast_sensor_t))

3.3 定时推送任务

代码语言:javascript
复制
// 推送周期可配置:3s / 5s / 10s
static uint32_t g_push_interval_ms = 3000;

void mcast_push_task(void) {
    static uint32_t last_push = 0;
    uint32_t now = sys_now();
    
    if (now - last_push < g_push_interval_ms) return;
    last_push = now;
    
    if (!network_is_up() || !g_mcast_pcb) return;
    
    uint8_t buf[MCAST_FRAME_MAX_SIZE];
    mcast_frame_hdr_t *hdr = (mcast_frame_hdr_t *)buf;
    
    memset(hdr, 0, sizeof(mcast_frame_hdr_t));
    hdr->header = 0x55BB;
    hdr->dev_addr = g_device_addr;
    hdr->seq = g_seq++;
    hdr->batch_ts = get_corrected_timestamp();
    hdr->interval_ms = g_push_interval_ms;
    
    // 填充传感器数据
    mcast_sensor_t *sensors = (mcast_sensor_t *)(buf + sizeof(mcast_frame_hdr_t));
    uint8_t count = 0;
    
    // 主传感器(自身)
    sensors[count].sensor_id = 0;
    sensors[count].temperature = (int16_t)(read_temperature() * 100);
    sensors[count].humidity = (uint16_t)(read_humidity() * 10);
    sensors[count].voltage = (uint16_t)(read_voltage() * 100);
    sensors[count].status = get_status_bits();
    count++;
    
    // 扩展传感器(如外接RS485探头,最多7个)
    for (int i = 1; i < MAX_SENSORS_PER_FRAME && i <= g_ext_sensor_count; i++) {
        sensors[count].sensor_id = i;
        sensors[count].temperature = (int16_t)(read_ext_temperature(i) * 100);
        sensors[count].humidity = (uint16_t)(read_ext_humidity(i) * 10);
        sensors[count].voltage = 0;  // 外接探头无独立电压
        sensors[count].status = get_ext_status(i);
        count++;
    }
    
    hdr->sensor_count = count;
    
    // CRC覆盖整个帧
    uint16_t frame_len = sizeof(mcast_frame_hdr_t) + count * sizeof(mcast_sensor_t);
    uint16_t crc = crc16_modbus(buf, frame_len);
    memcpy(buf + frame_len, &crc, 2);
    frame_len += 2;
    
    // 多播发送
    udp_sendto(g_mcast_pcb, 
               pbuf_alloc_ref(buf, frame_len),  // 简化:实际用pbuf_alloc + memcpy
               &g_mcast_addr, g_mcast_port);
}

3.4 LWIP配置注意

代码语言:javascript
复制
// lwipopts.h 多播相关
#define LWIP_IGMP               1       // 使能IGMP
#define MEMP_NUM_IGMP_GROUP     8       // 多播组数量
#define LWIP_MULTICAST_TX_OPTIONS 1     // 允许设置多播TTL
#define LWIP_SO_BINDTODEVICE    0       // 不需要绑定特定接口

四、服务端多播接收

4.1 Python异步多播接收

代码语言:javascript
复制
#!/usr/bin/env python3
"""
UDP多播接收服务 - 算力机房批量采集
"""
import asyncio
import struct
import socket
import crc16
from dataclasses import dataclass, field
from datetime import datetime, timezone
from typing import Dict, List, Optional


@dataclass
class SensorData:
    sensor_id: int
    temperature: float
    humidity: float
    voltage: float
    status: int


@dataclass
class MulticastReading:
    dev_addr: int
    seq: int
    batch_ts: int
    interval_ms: int
    sensors: List[SensorData] = field(default_factory=list)
    source_ip: str = ""
    recv_time: float = 0.0


class MulticastParser:
    """多播帧解析器"""
    
    HEADER_MAGIC = 0x55BB
    FRAME_MIN_SIZE = 14  # 12字节头 + 2字节CRC
    
    @classmethod
    def parse(cls, data: bytes, addr: tuple) -> Optional[MulticastReading]:
        if len(data) < cls.FRAME_MIN_SIZE:
            return None
        
        # 解析帧头
        try:
            header, dev_addr, sensor_count, seq, batch_ts, interval_ms = \
                struct.unpack(">HBBHI20xH", data[:12])  # 简化,实际按结构解
        except struct.error:
            return None
        
        if header != cls.HEADER_MAGIC:
            return None
        
        if sensor_count > 8 or sensor_count == 0:
            return None
        
        # CRC校验
        expected_len = 12 + sensor_count * 8 + 2
        if len(data) < expected_len:
            return None
        
        crc_calc = crc16.modbus(data[:expected_len - 2])
        crc_recv = struct.unpack("<H", data[expected_len - 2:expected_len])[0]
        if crc_calc != crc_recv:
            return None
        
        # 解析传感器数据
        sensors = []
        offset = 12
        for _ in range(sensor_count):
            sid, temp_raw, hum_raw, volt_raw, status = \
                struct.unpack(">BhHHB", data[offset:offset + 8])
            offset += 8
            
            sensors.append(SensorData(
                sensor_id=sid,
                temperature=temp_raw / 100.0,
                humidity=hum_raw / 10.0,
                voltage=volt_raw / 100.0,
                status=status,
            ))
        
        return MulticastReading(
            dev_addr=dev_addr,
            seq=seq,
            batch_ts=batch_ts,
            interval_ms=interval_ms,
            sensors=sensors,
            source_ip=addr[0],
        )


class MulticastReceiver:
    """多播接收器"""
    
    def __init__(self, mcast_group: str, mcast_port: int, handler: callable):
        self.mcast_group = mcast_group
        self.mcast_port = mcast_port
        self.handler = handler
        self._running = False
        self._transport = None
        self._protocol = None
    
    async def start(self):
        loop = asyncio.get_running_loop()
        
        # 创建UDP socket并加入多播组
        sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM, socket.IPPROTO_UDP)
        sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
        sock.bind(('', self.mcast_port))
        
        mreq = struct.pack(
            "4sl",
            socket.inet_aton(self.mcast_group),
            socket.INADDR_ANY
        )
        sock.setsockopt(socket.IPPROTO_IP, socket.IP_ADD_MEMBERSHIP, mreq)
        
        # 增大接收缓冲区
        sock.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, 4 * 1024 * 1024)
        
        self._transport, self._protocol = await loop.create_datagram_endpoint(
            lambda: MulticastProtocol(self.handler),
            sock=sock,
        )
        self._running = True
        print(f"多播接收器启动: {self.mcast_group}:{self.mcast_port}")
    
    async def stop(self):
        self._running = False
        if self._transport:
            self._transport.close()


class MulticastProtocol(asyncio.DatagramProtocol):
    def __init__(self, handler):
        self.handler = handler
        self.parser = MulticastParser()
    
    def datagram_received(self, data, addr):
        reading = self.parser.parse(data, addr)
        if reading:
            reading.recv_time = asyncio.get_event_loop().time()
            asyncio.create_task(self.handler(reading))
    
    def error_received(self, exc):
        print(f"多播接收错误: {exc}")

4.2 多组订阅(遥测 + Trap + 心跳)

代码语言:javascript
复制
async def main():
    # 同时订阅三个多播组
    receivers = [
        MulticastReceiver("239.255.1.100", 9000, handle_telemetry),
        MulticastReceiver("239.255.1.101", 9001, handle_trap),
        MulticastReceiver("239.255.1.102", 9002, handle_heartbeat),
    ]
    
    for r in receivers:
        await r.start()
    
    try:
        await asyncio.Event().wait()
    finally:
        for r in receivers:
            await r.stop()


async def handle_telemetry(reading: MulticastReading):
    """处理遥测数据"""
    for sensor in reading.sensors:
        print(f"[{reading.dev_addr:03d}-{sensor.sensor_id}] "
              f"T={sensor.temperature:.1f}℃ "
              f"H={sensor.humidity:.1f}% "
              f"V={sensor.voltage:.2f}V")


async def handle_trap(reading: MulticastReading):
    """处理事件Trap"""
    for sensor in reading.sensors:
        if sensor.status != 0:
            print(f"⚠ Trap 设备{reading.dev_addr} 状态=0x{sensor.status:02X}")


async def handle_heartbeat(reading: MulticastReading):
    """处理心跳"""
    # 更新节点在线状态
    pass

五、网络基础设施配置

5.1 交换机IGMP Snooping

代码语言:javascript
复制
! 华为交换机示例
vlan 100
 igmp-snooping enable
 igmp-snooping version 2
 igmp-snooping querier enable
 igmp-snooping querier ip 10.10.100.1

! 端口配置
interface GigabitEthernet0/0/1
 port link-type access
 port default vlan 100
 storm-control broadcast min-cir 1000
 storm-control multicast min-cir 2000

关键配置:

  • 启用IGMP Snooping,防止多播泛洪到所有端口
  • 配置Querier,定期发送IGMP General Query
  • 风暴抑制:限制多播流量上限,防止异常时网络拥塞

5.2 多播路由(跨子网场景)

代码语言:javascript
复制
! 跨子网时需要PIM-SM
router pim
 ip pim rp-address 10.10.0.1
 ip pim ssm-range 232.0.0.0/8

interface Vlan100
 ip pim sparse-mode

interface Vlan200
 ip pim sparse-mode

六、丢包检测与补传联动

6.1 序列号连续性检测

代码语言:javascript
复制
class LossDetector:
    """基于序列号和时间间隔的丢包检测"""
    
    def __init__(self):
        self._last_seq: Dict[int, int] = {}      # dev_addr -> last_seq
        self._last_ts: Dict[int, float] = {}     # dev_addr -> last_recv_time
        self._lost_count: Dict[int, int] = {}
        self._total_count: Dict[int, int] = {}
    
    def on_packet(self, reading: MulticastReading) -> Optional[int]:
        """返回丢失的包数,None表示无法判断"""
        dev = reading.dev_addr
        seq = reading.seq
        now = reading.recv_time
        
        self._total_count[dev] = self._total_count.get(dev, 0) + 1
        
        if dev not in self._last_seq:
            self._last_seq[dev] = seq
            self._last_ts[dev] = now
            return None
        
        expected = (self._last_seq[dev] + 1) % 65536
        lost = 0
        
        if seq != expected:
            # 计算丢失数量(处理回绕)
            lost = (seq - expected) % 65536
            self._lost_count[dev] = self._lost_count.get(dev, 0) + lost
            print(f"设备{dev:03d} 丢包: 期望{expected} 收到{seq} "
                  f"丢失{lost} 丢包率{self.loss_rate(dev):.1f}%")
        
        self._last_seq[dev] = seq
        self._last_ts[dev] = now
        return lost if lost > 0 else 0
    
    def loss_rate(self, dev: int) -> float:
        total = self._total_count.get(dev, 0)
        if total == 0: return 0.0
        lost = self._lost_count.get(dev, 0)
        return lost / total * 100.0

6.2 Modbus兜底触发

代码语言:javascript
复制
# 当丢包率超过阈值时,触发Modbus TCP补读
MODBUS_FALLBACK_THRESHOLD = 5.0  # 丢包率>5%触发

async def check_and_fallback(dev_addr: int, loss_rate: float, 
                             modbus_clients: dict):
    """丢包率过高时,通过Modbus TCP补读"""
    if loss_rate > MODBUS_FALLBACK_THRESHOLD:
        client = modbus_clients.get(dev_addr)
        if client:
            try:
                data = await client.read_holding_registers(0, 4)
                if data:
                    # 写入时序库,标记为modbus_fallback
                    await influx_write(dev_addr, data, source="modbus_fallback")
                    print(f"设备{dev_addr:03d} Modbus补读成功")
            except Exception as e:
                print(f"设备{dev_addr:03d} Modbus补读失败: {e}")

七、性能对比

7.1 千节点场景实测

指标

Modbus TCP轮询

UDP多播推送

采集周期

5s

3s

服务端CPU

~45% (8核)

~3% (8核)

服务端内存

~1.2GB (连接池)

~80MB

网络带宽(采集面)

~2.4Mbps (请求+响应)

~0.6Mbps (单向)

TCP连接数

1200

0

设备侧PCB占用

2-4 per node

0

数据到达延迟

取决于轮询顺序,最大5s

最大3s,均匀

丢包率(正常)

0% (TCP重传)

<0.1%

丢包率(拥塞)

超时增加

可能>5%,触发Modbus兜底


八、排障清单

现象

原因

排查

服务端收不到多播包

未加入多播组/防火墙拦截

tcpdump igmp看成员报告,iptables -L

收到但CRC错

内存损坏或字节序错误

抓包对比原始hex

部分节点无数据

IGMP Snooping未正确学习

交换机看mrouter端口和成员表

丢包率高

交换机多播风暴抑制

检查风暴抑制阈值

重复包

网络环路或IGMP重复加入

检查STP状态和端口状态

跨子网收不到

PIM路由未配置

show ip mroute看转发状态


九、关键词

算力服务器机房、机房改造、UDP多播、批量采集、RJ45温湿度变送器、IGMP Snooping、多播路由、LWIP多播、丢包检测、Modbus兜底、高密节点、采集架构


十、一句话总结

算力机房千节点改造中,UDP多播把采集模式从"服务端逐个轮询"翻转为"设备主动批量推送、一组多收",服务端负载从O(N)降到O(1),设备侧零TCP状态;配合IGMP Snooping控制泛洪、序列号检测丢包、Modbus TCP兜底补读,在规模化场景下实现低开销、高实时、可降级的环境数据采集。

物联网 #Modbus #TCP/IP #UDP #POE供电 #腾讯云 #Wireshark #Python #InfluxDB #以太网温湿度传感器 #网口温湿度变送器 #机房监控 #算力机房 #UDP多播 #IGMP #批量采集 #工业物联网

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

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

目录
  • 算力服务器机房改造:UDP多播方式,批量采集RJ45温湿度变送器数据
    • 一、为什么算力机房需要UDP多播
      • 1.1 规模带来的矛盾
      • 1.2 适用边界
    • 二、多播架构设计
      • 2.1 整体拓扑
      • 2.2 多播地址规划
    • 三、设备侧多播推送实现
      • 3.1 LWIP多播初始化
      • 3.2 批量数据帧格式
      • 3.3 定时推送任务
      • 3.4 LWIP配置注意
    • 四、服务端多播接收
      • 4.1 Python异步多播接收
      • 4.2 多组订阅(遥测 + Trap + 心跳)
    • 五、网络基础设施配置
      • 5.1 交换机IGMP Snooping
      • 5.2 多播路由(跨子网场景)
    • 六、丢包检测与补传联动
      • 6.1 序列号连续性检测
      • 6.2 Modbus兜底触发
    • 七、性能对比
      • 7.1 千节点场景实测
    • 八、排障清单
    • 九、关键词
    • 十、一句话总结
  • 物联网 #Modbus #TCP/IP #UDP #POE供电 #腾讯云 #Wireshark #Python #InfluxDB #以太网温湿度传感器 #网口温湿度变送器 #机房监控 #算力机房 #UDP多播 #IGMP #批量采集 #工业物联网
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档