

近期在调试一套工业环境监控系统时,遇到了关于Modbus连接保持的问题。现场部署了一批网口温湿度变送器,PoE取电、Modbus TCP上云,同时利旧接入存量RS485探头,并开启了SNMP和UDP Trap服务。算力服务器机房改造项目里,节点规模从百级跳到千级,传统"一问一答"的Modbus TCP轮询模式在扇出和时延上开始吃力:400节点×5s轮询 ≈ 80 RPS,勉强够用;1200节点×3s轮询 ≈ 400 RPS,单采集服务已经到瓶颈。这篇讲一种工程上已经验证的替代路径——UDP多播批量采集,把"服务端逐个轮询"变成"设备主动批量上报、服务端一组收",在算力机房高密度场景下显著降低采集面负载。
维度 | 传统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和内存占用下降一个数量级。
┌─────────────────────────────┐
│ 采集服务端 │
│ UDP多播监听 (239.255.1.100) │
│ ↓ │
│ 解析 → 去重 → 时序库 │
└──────────────┬──────────────┘
│ 多播组 239.255.1.100:9000
┌────────────────────────┼────────────────────────┐
│ │ │
[列1: 变送器×48] [列2: 变送器×48] [列3: 变送器×48]
定时组播推送 定时组播推送 定时组播推送
(每3秒) (每3秒) (每3秒)多播组 | 用途 | TTL |
|---|---|---|
239.255.1.100 | 周期遥测数据 | 1(本子网) |
239.255.1.101 | 事件Trap(越限/掉电) | 1 |
239.255.1.102 | 管理面心跳(设备在线检测) | 1 |
239.255.1.200 | 服务端→设备广播指令(可选) | 1 |
地址选择原则:

// 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;
}
}// 多播批量上报帧:一帧携带多个节点的数据
// 用于高密度机柜场景,一台设备可接多个探头或自身多传感器
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))// 推送周期可配置: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);
}// 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 // 不需要绑定特定接口#!/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}")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! 华为交换机示例
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关键配置:
! 跨子网时需要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-modeclass 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# 当丢包率超过阈值时,触发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}")指标 | 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兜底补读,在规模化场景下实现低开销、高实时、可降级的环境数据采集。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。