
关键词:算力机房动环改造、UDP并发采集、以太网温湿度变送器、多网卡负载均衡、内核协议栈优化、异步I/O、高可用采集、数据完整性 标签:#物联网 #Modbus #TCP/IP #UDP #POE供电 #腾讯云 #Wireshark #Python #InfluxDB #以太网温湿度传感器 #网口温湿度变送器 #机房监控 #边缘计算
前几篇覆盖了中小型机房和通用服务器机房的采集方案。本篇聚焦算力机房——AI 集群、GPU 服务器、智算中心这类场景。这类机房对动环系统提出了截然不同的要求:
维度 | 传统机房 | 算力机房 |
|---|---|---|
功率密度 | 3~8 kW/机柜 | 20~50 kW/机柜(液冷可达 100+ kW) |
温湿度梯度 | 相对均匀 | 热点明显,冷热通道温差大 |
采集点密度 | 每机柜 1~2 点 | 每机柜 4~8 点(进风/出风/冷通道/热通道) |
规模 | 数十~数百点 | 数千点 |
告警延迟要求 | 分钟级 | 秒级(GPU 降频阈值窄) |
网络环境 | 独立管理网 | 与计算网络可能共用 Spine-Leaf,QoS 复杂 |
核心矛盾:算力机房需要高密度、低延迟、高可靠的采集,但 UDP 协议本身不保证送达。如何在数千台 UDP 变送器的规模下,既利用 UDP 的低开销优势,又保证数据完整性?

┌─────────────────────────────────────────────────────────────────────┐
│ 算力机房动环采集架构 │
├─────────────────────────────────────────────────────────────────────┤
│ │
│ ┌─────────────────────────────────────────────────────────────┐ │
│ │ 采集接入层(多节点,每节点负责一个 Pod/Row) │ │
│ │ ┌─────────┐ ┌─────────┐ ┌─────────┐ ┌─────────┐ │ │
│ │ │ 节点 A │ │ 节点 B │ │ 节点 C │ │ 节点 D │ ... │ │
│ │ │ 128台 │ │ 128台 │ │ 128台 │ │ 128台 │ │ │
│ │ │ 4×25G │ │ 4×25G │ │ 4×25G │ │ 4×25G │ │ │
│ │ └────┬────┘ └────┬────┘ └────┬────┘ └────┬────┘ │ │
│ └───────┼────────────┼────────────┼────────────┼────────────┘ │
│ │ │ │ │ │
│ ┌───────▼────────────▼────────────▼────────────▼────────────┐ │
│ │ 汇聚层(Kafka / Pulsar) │ │
│ │ 消息队列解耦采集与存储,削峰填谷 │ │
│ └───────────────────────┬────────────────────────────────────┘ │
│ │ │
│ ┌───────────────────────▼────────────────────────────────────┐ │
│ │ 处理层(流式计算 + 规则引擎) │ │
│ │ - 实时聚合(每机柜平均/最大/最小) │ │
│ │ - 异常检测(梯度突变、连续越限) │ │
│ │ - 数据补全(插值、标记缺失) │ │
│ └──────────┬────────────────────────┬────────────────────────┘ │
│ │ │ │
│ ┌──────────▼──────┐ ┌─────────▼────────┐ │
│ │ 时序数据库 │ │ 告警/联动系统 │ │
│ │ InfluxDB集群 │ │ 大屏/Webhook │ │
│ └─────────────────┘ └──────────────────┘ │
└─────────────────────────────────────────────────────────────────────┘设计要点:
以单节点负责 128 台变送器、上报周期 5s 为例:
指标 | 数值 |
|---|---|
设备数 | 128 |
上报周期 | 5s |
理论 PPS | 25.6 pps(极低) |
实际突发 PPS | ~50 pps(微突发) |
单包大小 | 32~64 字节 |
带宽 | < 50 Kbps |
CPU 占用 | < 5%(优化后) |
结论:UDP 并发采集的瓶颈不在网络带宽,而在:
# 1. 增大 UDP 接收缓冲区
sysctl -w net.core.rmem_max=268435456 # 256MB
sysctl -w net.core.rmem_default=16777216 # 16MB
sysctl -w net.ipv4.udp_mem="262144 327680 393216" # min/pressure/max (pages)
# 2. 增大套接字接收队列长度
sysctl -w net.core.netdev_max_backlog=16384
sysctl -w net.core.somaxconn=8192
# 3. 减少 softirq 延迟
sysctl -w net.core.dev_weight=600
sysctl -w net.core.busy_poll=50
sysctl -w net.core.busy_read=50
# 4. 启用 RSS(接收侧缩放)和 RPS(软件 RPS)
# 查看网卡队列数
ls /sys/class/net/eth0/queues/
# 配置 RPS(将软中断分散到多个 CPU)
echo ffff > /sys/class/net/eth0/queues/rx-0/rps_cpus
echo 32768 > /sys/class/net/eth0/queues/rx-0/rps_flow_cnt
# 5. 启用 XPS(发送侧)
echo ffff > /sys/class/net/eth0/queues/tx-0/xps_cpus
# 6. 中断亲和性绑定(irqbalance 或手动)
# 查看中断号
cat /proc/interrupts | grep eth0
# 绑定到特定 CPU(示例)
echo 2 > /proc/irq/123/smp_affinity # 绑定到 CPU 1
方案对比:
方案 | 优点 | 缺点 | 推荐度 |
|---|---|---|---|
单进程 asyncio | 简单、无锁 | GIL 瓶颈、单核 | ⭐⭐ |
多进程(SO_REUSEPORT) | 真并行、无 GIL | 进程管理复杂 | ⭐⭐⭐⭐⭐ |
多线程 | 共享内存方便 | GIL 仍有影响 | ⭐⭐⭐ |
DPDK/用户态协议栈 | 极致性能 | 开发量大、维护难 | ⭐(过度) |
推荐:多进程 + SO_REUSEPORT,每个进程绑定一个 CPU 核,内核自动做四元组哈希分发。
"""
collector_node.py - 多进程 UDP 采集节点
"""
import asyncio
import socket
import struct
import multiprocessing as mp
import signal
import os
from dataclasses import dataclass
from kafka import KafkaProducer
import json
NUM_WORKERS = mp.cpu_count() # 通常 8~32 核
UDP_PORT = 9000
KAFKA_SERVERS = ["kafka-1:9092", "kafka-2:9092", "kafka-3:9092"]
@dataclass
class RawReading:
dev_id: str
seq: int
temperature: float
humidity: float
timestamp: float
recv_ts: float
worker_pid: int
class Worker:
def __init__(self, worker_id: int, cpu_affinity: int):
self.worker_id = worker_id
self.cpu_affinity = cpu_affinity
self.producer = None
def setup_kafka(self):
"""每个 worker 独立的 Kafka producer(线程安全)"""
self.producer = KafkaProducer(
bootstrap_servers=KAFKA_SERVERS,
value_serializer=lambda v: json.dumps(v).encode("utf-8"),
acks=1,
linger_ms=50, # 批量发送
batch_size=65536,
compression_type="snappy",
retries=3,
max_in_flight_requests_per_connection=5
)
def parse_packet(self, data: bytes, addr: tuple, recv_ts: float) -> RawReading | None:
if len(data) < 32:
return None
try:
header, dev_id, seq, ts, temp_raw, hum_raw, status, crc = \
struct.unpack_from(">HHHIHHBH", data, 0)
if header != 0xAA55:
return None
return RawReading(
dev_id=f"{addr[0]}:{dev_id}",
seq=seq,
temperature=temp_raw / 10.0 if temp_raw < 0x8000 else (temp_raw - 0x10000) / 10.0,
humidity=hum_raw / 10.0,
timestamp=ts,
recv_ts=recv_ts,
worker_pid=os.getpid()
)
except Exception:
return None
async def run(self):
"""Worker 主循环"""
# CPU 亲和性绑定
os.sched_setaffinity(0, {self.cpu_affinity})
self.setup_kafka()
# 创建 UDP socket,启用 SO_REUSEPORT
sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEPORT, 1)
sock.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, 16 * 1024 * 1024) # 16MB
sock.bind(("0.0.0.0", UDP_PORT))
sock.setblocking(False)
loop = asyncio.get_event_loop()
print(f"Worker {self.worker_id} (PID={os.getpid()}) bound to CPU {self.cpu_affinity}, port {UDP_PORT}")
while True:
try:
data, addr = await loop.sock_recvfrom(sock, 1024)
recv_ts = loop.time()
reading = self.parse_packet(data, addr, recv_ts)
if reading:
# 发送到 Kafka
self.producer.send(
"env_raw",
key=reading.dev_id.encode(),
value={
"dev_id": reading.dev_id,
"seq": reading.seq,
"temperature": reading.temperature,
"humidity": reading.humidity,
"timestamp": reading.timestamp,
"recv_ts": reading.recv_ts,
"worker_pid": reading.worker_pid
}
)
except asyncio.CancelledError:
break
except Exception as e:
print(f"Worker {self.worker_id} error: {e}")
sock.close()
self.producer.close()
def start_worker(worker_id: int, cpu_affinity: int):
"""子进程入口"""
worker = Worker(worker_id, cpu_affinity)
loop = asyncio.new_event_loop()
asyncio.set_event_loop(loop)
# 信号处理
def shutdown():
for task in asyncio.all_tasks(loop):
task.cancel()
loop.stop()
loop.add_signal_handler(signal.SIGTERM, shutdown)
loop.add_signal_handler(signal.SIGINT, shutdown)
try:
loop.run_until_complete(worker.run())
finally:
loop.close()
def main():
"""主进程:启动所有 worker"""
print(f"Starting {NUM_WORKERS} workers for UDP port {UDP_PORT}")
processes = []
for i in range(NUM_WORKERS):
p = mp.Process(target=start_worker, args=(i, i % mp.cpu_count()))
p.start()
processes.append(p)
print(f"Spawned worker {i} (PID={p.pid})")
def shutdown_all(sig, frame):
print("Shutting down...")
for p in processes:
p.terminate()
for p in processes:
p.join(timeout=5)
if p.is_alive():
p.kill()
exit(0)
signal.signal(signal.SIGTERM, shutdown_all)
signal.signal(signal.SIGINT, shutdown_all)
# 主进程等待
for p in processes:
p.join()
if __name__ == "__main__":
main()# 查看每个 worker 的 CPU 使用情况
top -H -p $(pgrep -f collector_node | head -1)
# 查看 UDP 接收统计
watch -n 1 'cat /proc/net/udp | wc -l'
nstat -az | grep UdpRcv
# 查看软中断分布
cat /proc/softirqs | grep NET_RX
# 查看套接字接收队列
ss -lunpem | grep 9000
# 压测:用 pktgen 模拟 UDP 流量
# 单节点模拟 512 台 × 2s 周期
python pktgen_sim.py --count 512 --interval 2 --port 9000 --target 10.20.1.100每个 worker 维护 per-device 的状态表:
"""
packet_tracker.py - 丢包检测
"""
from collections import defaultdict
from dataclasses import dataclass, field
@dataclass
class DeviceState:
last_seq: int = -1
lost_count: int = 0
received_count: int = 0
last_recv_ts: float = 0.0
gap_history: list = field(default_factory=list) # 最近 N 次丢包 burst 大小
class PacketTracker:
def __init__(self, gap_threshold=5, window_size=100):
self.devices: dict[str, DeviceState] = {}
self.gap_threshold = gap_threshold # 序号跳变超过此值视为回绕或严重丢包
self.window_size = window_size
def on_packet(self, dev_id: str, seq: int, recv_ts: float) -> dict:
"""返回丢包信息"""
state = self.devices.get(dev_id)
if state is None:
state = DeviceState()
self.devices[dev_id] = state
state.last_seq = seq
state.received_count = 1
state.last_recv_ts = recv_ts
return {"dev_id": dev_id, "lost": 0, "event": "first_packet"}
state.received_count += 1
state.last_recv_ts = recv_ts
if seq == state.last_seq:
return {"dev_id": dev_id, "lost": 0, "event": "duplicate"}
# 处理序号回绕(16 位序号)
expected = (state.last_seq + 1) & 0xFFFF
if seq == expected:
state.last_seq = seq
return {"dev_id": dev_id, "lost": 0, "event": "ok"}
# 计算跳变
if seq > state.last_seq:
gap = seq - state.last_seq - 1
else:
# 回绕
gap = (0x10000 - state.last_seq) + seq - 1
if gap > self.gap_threshold:
# 序号跳变过大,可能是设备重启或严重丢包
state.last_seq = seq
return {"dev_id": dev_id, "lost": -1, "event": "seq_reset", "gap": gap}
state.lost_count += gap
state.gap_history.append(gap)
if len(state.gap_history) > self.window_size:
state.gap_history.pop(0)
state.last_seq = seq
return {"dev_id": dev_id, "lost": gap, "event": "gap", "total_lost": state.lost_count}"""
stream_processor.py - Kafka 消费者 + 数据质量处理
"""
from kafka import KafkaConsumer, KafkaProducer
import json
from datetime import datetime, timedelta
from collections import defaultdict
consumer = KafkaConsumer(
"env_raw",
bootstrap_servers=["kafka-1:9092"],
value_deserializer=lambda v: json.loads(v.decode("utf-8")),
group_id="env_processor",
auto_offset_reset="latest"
)
producer = KafkaProducer(
bootstrap_servers=["kafka-1:9092"],
value_serializer=lambda v: json.dumps(v).encode("utf-8")
)
# 设备状态缓存
device_cache = defaultdict(dict)
# 最近有效值(用于插值)
last_valid = {}
for msg in consumer:
reading = msg.value
dev_id = reading["dev_id"]
# 1. 数据质量标记
quality = "good"
if reading["recv_ts"] - reading["timestamp"] > 10:
quality = "clock_drift"
# 2. 连续丢失检测(基于 recv_ts 间隔)
last = last_valid.get(dev_id)
if last:
interval = reading["recv_ts"] - last["recv_ts"]
expected_interval = 5.0 # 设备上报周期
if interval > expected_interval * 2.5:
# 中间有丢包
missing_count = int(interval / expected_interval) - 1
# 发送插值标记
for i in range(1, missing_count + 1):
interp_ts = last["recv_ts"] + i * expected_interval
producer.send("env_interpolated", value={
"dev_id": dev_id,
"timestamp": interp_ts,
"temperature": None, # 标记为插值缺失
"humidity": None,
"quality": "missing",
"missing_seq": i,
"total_missing": missing_count
})
# 3. 更新缓存
last_valid[dev_id] = reading
# 4. 发送到处理后的 topic
producer.send("env_processed", value={
"dev_id": dev_id,
"seq": reading["seq"],
"temperature": reading["temperature"],
"humidity": reading["humidity"],
"timestamp": reading["timestamp"],
"recv_ts": reading["recv_ts"],
"quality": quality
})接入交换机(每 ToR 48 端口):
- 启用 IGMP Snooping(如果设备支持组播)
- 配置广播抑制:broadcast-suppression 10%
- 配置风暴控制:storm-control broadcast level 10
- PoE 功率预算:每端口 15.4W,总预算 370W(24 端口满配)
- 启用 LLDP:便于设备发现和拓扑
- QoS:为管理流量(UDP 9000)打 DSCP AF11VLAN 规划:
- VLAN 10:计算网络(GPU 训练流量)
- VLAN 20:存储网络(NVMe over Fabrics)
- VLAN 30:动环管理(温湿度变送器)
- VLAN 40:IPMI/BMC
流量隔离:
- 动环 VLAN 独立,不与计算网络混跑
- 汇聚交换机上联带宽 ≥ 10Gbps(即使动环流量只有几 Mbps,也要避免拥塞)
- 如果 Spine-Leaf 架构,动环流量走独立 Leaf 对参数 | 推荐值 | 理由 |
|---|---|---|
上报周期 | 5s(常规)/ 2s(热点区域) | 平衡实时性和网络负载 |
目标端口 | 9000(固定) | 便于防火墙和 ACL |
目标 IP | 采集节点 VIP(Keepalived) | 高可用 |
本地存储 | 启用(≥10 万条) | 断网补传 |
心跳间隔 | 与上报周期一致 | 简化逻辑 |
┌─────────────┐
│ Keepalived │
│ VIP: x.x.x.x│
└──────┬──────┘
│
┌─────────┼─────────┐
│ │ │
┌────▼───┐ ┌──▼────┐ ┌──▼────┐
│ 主节点 │ │ 备节点 │ │ 备节点 │
│ (Active)│ │(Standby)│ │(Standby)│
└────────┘ └───────┘ └───────┘
VIP 漂移条件:
- 主节点进程挂掉
- 主节点网络不可达
- 主节点 UDP 接收丢包率 > 1%
设备侧:
- 设备配置主备双目标 IP(部分高端设备支持)
- 或:设备只发一个目标,VIP 漂移后自动切换三层保障:
1. 设备本地存储:断网期间数据不丢
2. 采集节点 Kafka:消息持久化,多副本
3. 时序数据库:InfluxDB 多节点集群,数据冗余
恢复流程:
- 网络恢复 → 设备主动补传(如果支持)
- 或:采集节点从设备拉取历史(Modbus TCP 读取 Flash)
- 或:人工导出 CSV 补录指标 | 类型 | 告警阈值 |
|---|---|---|
udp_packets_received | Counter | — |
udp_packets_dropped_kernel | Counter | > 0 持续 1min |
udp_receive_buffer_errors | Counter | > 10/s |
per_device_packet_loss_rate | Gauge | > 0.1% 持续 5min |
per_device_last_seen_age | Gauge | > 30s |
kafka_producer_queue_size | Gauge | > 10000 |
worker_cpu_usage | Gauge | > 80% |
worker_memory_rss | Gauge | > 500MB |
max_block_ms 或降级到本地文件。 recv_ts 为准。 算力机房动环改造中,UDP 多并发采集的核心不是"能不能收到",而是收到多少、丢了多少、丢了怎么办。通过多进程 SO_REUSEPORT 模型横向扩展接收能力,内核参数调优消除软中断瓶颈,Kafka 解耦采集与存储,流式处理做数据质量标记和插值补全,三层保障(设备本地存储 + 消息队列 + 时序库)确保数据完整性。最终交付的是一个可观测、可扩展、可容错的采集系统,而不是一个"尽力而为"的 UDP 接收器。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。