

先说清场景:独立配电柜集群指的是分散在厂区各处的配电室/配电间里,成排独立运行的低压配电柜。每个柜内装一台以太网温湿度变送器,监测柜内温湿度——这是判断柜内是否有过热隐患、是否凝露、是否有电气接点异常升温的关键指标。配电柜的特点是数量多、点位分散、单点数据量小但采集频率要求不高,最适合用"长连接 + 批量上传"的模式。
先说一个真实案例:
某汽车零部件厂,4 个配电室共 96 台独立配电柜,每台柜内装 1 台以太网温湿度变送器。 初始方案:采集服务器对 96 台设备各开一个 TCP 连接,每 10s 轮询一次,每次读 2 个寄存器。 问题:

┌─────────────────────────────────────────────────────────────────────┐
│ 采集服务器(中心侧) │
│ ┌────────────────────────────────────────────────────────────────┐ │
│ │ TCP 长连接服务端(单端口 5020) │ │
│ │ ┌────────────┐ ┌────────────┐ ┌────────────┐ │ │
│ │ │ 连接管理器 │ │ 协议解析器 │ │ 数据分发器 │ │ │
│ │ │ 96 连接 │→│ 批量解包 │→│ 写时序库 │ │ │
│ │ │ KeepAlive │ │ CRC 校验 │ │ 告警判断 │ │ │
│ │ │ 心跳检测 │ │ 序列号检测 │ │ 数据转发 │ │ │
│ │ └────────────┘ └────────────┘ └────────────┘ │ │
│ └────────────────────────────────────────────────────────────────┘ │
│ │ │
│ ┌───────────────┼───────────────┐ │
│ ▼ ▼ ▼ │
│ ┌──────────────┐ ┌──────────────┐ ┌──────────────┐ │
│ │ 时序数据库 │ │ 告警引擎 │ │ BA 平台接口 │ │
│ │ InfluxDB │ │ 越限判断 │ │ BACnet/MQTT │ │
│ └──────────────┘ └──────────────┘ └──────────────┘ │
└───────────────────────┬─────────────────────────────────────────────┘
│
▼ (监测 VLAN)
┌─────────────────────────────────────────────────────────────────────┐
│ 配电柜集群(传感器侧) │
│ │
│ ┌─────────┐ ┌─────────┐ ┌─────────┐ ┌─────────┐ │
│ │ 柜01 │ │ 柜02 │ │ 柜03 │ ... │ 柜96 │ │
│ │ 传感器 │ │ 传感器 │ │ 传感器 │ │ 传感器 │ │
│ │ TCP Client│ │ TCP Client│ │ TCP Client│ │ TCP Client│ │
│ │ 本地缓存 │ │ 本地缓存 │ │ 本地缓存 │ │ 本地缓存 │ │
│ └────┬────┘ └────┬────┘ └────┬────┘ └────┬────┘ │
│ └────────────┬┴────────────┬┴──────────────┘ │
│ ▼ │
│ ┌───────────────────┐ │
│ │ 接入交换机(PoE) │ │
│ │ VLAN 100 │ │
│ └───────────────────┘ │
│ │
│ 供电:每台配电柜内 PoE 交换机或 DC 24V 供电 │
│ 布线:Cat6 S/FTP,沿桥架敷设,与动力电缆隔离 ≥30cm │
└─────────────────────────────────────────────────────────────────────┘维度 | 服务器轮询模式(旧) | 传感器主动上报模式(新) |
|---|---|---|
连接方向 | 服务器 → 传感器(Server 端被动等待请求) | 传感器 → 服务器(Client 端主动发起) |
连接数量 | 96 个独立连接 | 96 个长连接(但连接一旦建立不释放) |
请求方式 | 请求-响应模型 | 传感器单向推送(或请求-响应混合) |
数据实时性 | 取决于轮询周期 | 采集即上报,延迟更低 |
网络开销 | 每次 2 个小包(请求+响应) | 批量打包,一个包含多台设备数据 |
服务器压力 | 需管理轮询调度、超时、重试 | 只需维护连接表、解析数据 |
适用场景 | 设备数量少、需频繁控制 | 设备数量多、数据量小、单向采集 |
核心转变:从"服务器问,传感器答"变成"传感器采集完,主动告诉服务器"。
┌─────────────────────────────────────────────────────────────┐
│ 传感器固件架构 │
│ │
│ ┌────────────┐ ┌────────────┐ ┌───────────────────────┐ │
│ │ 采样任务 │→│ 数据缓存 │→│ TCP Client 上报任务 │ │
│ │ 2s 采集一次 │ │ Ring Buffer│ │ 长连接管理 │ │
│ │ 温湿度+I/O │ │ 容量: 128 │ │ 批量打包 │ │
│ └────────────┘ └────────────┘ │ 心跳/重连 │ │
│ └───────────────────────┘ │
│ │ │
│ ┌──────────┴────────────┐ │
│ │ TCP/IP 协议栈 │ │
│ │ LwIP + FreeRTOS │ │
│ └────────────────────────┘ │
└─────────────────────────────────────────────────────────────┘参数 | 值 | 说明 |
|---|---|---|
采集周期 | 2s | 温湿度采样间隔 |
上报周期 | 30s | 批量上报间隔 |
每批数据点 | 15 个(30s / 2s) | 每个包包含 15 组温湿度数据 |
上报缓冲 | 128 组 | Ring Buffer 容量,断线时缓存 |
TCP 端口 | 5020 | 服务器监听端口(非标准 502) |
KeepAlive | 30s/5s/3 | 空闲 30s 开始探测,间隔 5s,3 次失败断开 |
重连退避 | 1/2/4/8/16/32/60s | 指数退避 + 随机抖动 |
心跳 | 每 30s 发送一次 | 保持连接活跃 |
报文序列号 | 32 位递增 | 检测丢包 |
断线缓存 | 最多 128 组 × 32 台 | 超容量后覆盖最旧数据 |
/* 配电柜温湿度传感器 - TCP Client 批量上报 */
#include "lwip/tcp.h"
#include "lwip/sockets.h"
#include "freertos/FreeRTOS.h"
#include "freertos/task.h"
#include "freertos/queue.h"
#include <string.h>
#define SERVER_IP "192.168.100.200"
#define SERVER_PORT 5020
#define SAMPLE_INTERVAL 2000 // 2s 采样
#define REPORT_INTERVAL 30000 // 30s 上报
#define BUFFER_SIZE 128 // 本地缓存 128 组
#define KEEPALIVE_IDLE 30
#define KEEPALIVE_INTVL 5
#define KEEPALIVE_PROBES 3
/* 单组采样数据 */
typedef struct {
uint32_t timestamp; // 采样时间戳
int16_t temperature; // 温度 ×10(如 256 = 25.6℃)
int16_t humidity; // 湿度 ×10(如 523 = 52.3%RH)
uint8_t status; // 状态字
uint8_t reserved; // 保留
} __attribute__((packed)) SampleData;
/* 批量上报报文 */
typedef struct {
uint16_t header; // 帧头 0xAA55
uint16_t length; // 报文总长度
uint32_t device_id; // 设备 ID
uint32_t sequence; // 序列号
uint32_t timestamp; // 报文生成时间戳
uint8_t count; // 数据点数量
uint8_t reserved[3];
SampleData samples[0]; // 变长数组
} __attribute__((packed)) BatchReport;
static int g_sock = -1;
static bool g_connected = false;
static uint32_t g_sequence = 0;
static QueueHandle_t g_sample_queue = NULL; // 采样数据队列
static SemaphoreHandle_t g_buffer_mutex = NULL;
static SampleData g_buffer[BUFFER_SIZE];
static uint16_t g_buffer_head = 0;
static uint16_t g_buffer_tail = 0;
static uint16_t g_buffer_count = 0;
/* 设置 TCP KeepAlive */
static void set_keepalive(int sockfd)
{
int optval = 1;
socklen_t optlen = sizeof(optval);
setsockopt(sockfd, SOL_SOCKET, SO_KEEPALIVE, &optval, optlen);
optval = KEEPALIVE_IDLE;
setsockopt(sockfd, IPPROTO_TCP, TCP_KEEPIDLE, &optval, optlen);
optval = KEEPALIVE_INTVL;
setsockopt(sockfd, IPPROTO_TCP, TCP_KEEPINTVL, &optval, optlen);
optval = KEEPALIVE_PROBES;
setsockopt(sockfd, IPPROTO_TCP, TCP_KEEPCNT, &optval, optlen);
}
/* 连接到服务器 */
static int connect_to_server(void)
{
struct sockaddr_in server_addr;
int sock;
sock = socket(AF_INET, SOCK_STREAM, 0);
if (sock < 0) {
LWIP_LOG("socket failed");
return -1;
}
// 设置超时
struct timeval tv = {3, 0};
setsockopt(sock, SOL_SOCKET, SO_SNDTIMEO, &tv, sizeof(tv));
setsockopt(sock, SOL_SOCKET, SO_RCVTIMEO, &tv, sizeof(tv));
set_keepalive(sock);
// 禁用 Nagle(小包实时性优先,但批量上报时启用)
int nodelay = 0; // 批量上报时启用 Nagle,减少小包
setsockopt(sock, IPPROTO_TCP, TCP_NODELAY, &nodelay, sizeof(nodelay));
memset(&server_addr, 0, sizeof(server_addr));
server_addr.sin_family = AF_INET;
server_addr.sin_port = htons(SERVER_PORT);
server_addr.sin_addr.s_addr = inet_addr(SERVER_IP);
if (connect(sock, (struct sockaddr *)&server_addr, sizeof(server_addr)) < 0) {
close(sock);
return -1;
}
g_sock = sock;
g_connected = true;
LWIP_LOG("connected to %s:%d", SERVER_IP, SERVER_PORT);
return sock;
}
/* 采样任务 */
void sample_task(void *pvParameters)
{
TickType_t last_wake = xTaskGetTickCount();
while (1) {
// 读取温湿度
int16_t temp = read_temperature() * 10; // ×10 存储
int16_t hum = read_humidity() * 10;
uint8_t status = get_status();
SampleData sample = {
.timestamp = get_timestamp(),
.temperature = temp,
.humidity = hum,
.status = status,
.reserved = 0
};
// 写入环形缓冲区
xSemaphoreTake(g_buffer_mutex, portMAX_DELAY);
g_buffer[g_buffer_head] = sample;
g_buffer_head = (g_buffer_head + 1) % BUFFER_SIZE;
if (g_buffer_count < BUFFER_SIZE) {
g_buffer_count++;
} else {
// 缓冲区满,覆盖最旧数据
g_buffer_tail = (g_buffer_tail + 1) % BUFFER_SIZE;
}
xSemaphoreGive(g_buffer_mutex);
vTaskDelayUntil(&last_wake, pdMS_TO_TICKS(SAMPLE_INTERVAL));
}
}
/* 构建批量上报报文 */
static int build_batch_report(uint8_t *buf, int buf_size)
{
BatchReport *report = (BatchReport *)buf;
xSemaphoreTake(g_buffer_mutex, portMAX_DELAY);
// 确定本次上报的数据点数量
uint8_t count = g_buffer_count;
if (count > 15) count = 15; // 最多 15 个点
if (count == 0) {
xSemaphoreGive(g_buffer_mutex);
return 0;
}
// 填充报文头
report->header = 0x55AA;
report->length = sizeof(BatchReport) + count * sizeof(SampleData);
report->device_id = get_device_id();
report->sequence = ++g_sequence;
report->timestamp = get_timestamp();
report->count = count;
report->reserved[0] = 0;
report->reserved[1] = 0;
report->reserved[2] = 0;
// 复制数据
for (int i = 0; i < count; i++) {
report->samples[i] = g_buffer[g_buffer_tail];
g_buffer_tail = (g_buffer_tail + 1) % BUFFER_SIZE;
g_buffer_count--;
}
xSemaphoreGive(g_buffer_mutex);
return report->length;
}
/* 上报任务 */
void report_task(void *pvParameters)
{
uint8_t buffer[sizeof(BatchReport) + 15 * sizeof(SampleData) + 4]; // +4 for CRC
TickType_t last_report = xTaskGetTickCount();
// 初始连接
while (connect_to_server() < 0) {
LWIP_LOG("connect failed, retry in 5s");
vTaskDelay(pdMS_TO_TICKS(5000));
}
while (1) {
TickType_t now = xTaskGetTickCount();
// 检查上报周期
if (now - last_report >= pdMS_TO_TICKS(REPORT_INTERVAL)) {
if (!g_connected) {
// 重连
static uint8_t retry = 0;
uint32_t delay = 1000 * (1 << (retry > 6 ? 6 : retry));
delay = delay * (0.7f + ((float)rand() / RAND_MAX) * 0.6f);
retry++;
vTaskDelay(pdMS_TO_TICKS(delay));
connect_to_server();
if (g_connected) retry = 0;
continue;
}
// 构建并发送批量报文
int len = build_batch_report(buffer, sizeof(buffer));
if (len > 0) {
// 计算 CRC
uint16_t crc = calc_crc16(buffer, len);
memcpy(buffer + len, &crc, 2);
len += 2;
if (send(g_sock, buffer, len, 0) < 0) {
LWIP_LOG("send failed, errno=%d", errno);
g_connected = false;
close(g_sock);
g_sock = -1;
}
}
last_report = now;
}
// 心跳(每 30s)
static TickType_t last_heartbeat = 0;
if (now - last_heartbeat >= pdMS_TO_TICKS(30000)) {
if (g_connected) {
uint8_t hb[8] = {0xAA, 0x55, 0x00, 0x08, 0x00, 0x00, 0x00, 0x00};
// 填充设备 ID 和序列号
memcpy(hb + 2, &g_device_id, 4);
memcpy(hb + 6, &g_sequence, 2);
send(g_sock, hb, 8, 0);
}
last_heartbeat = now;
}
vTaskDelay(pdMS_TO_TICKS(1000));
}
}
/* 主入口 */
void app_main(void)
{
g_sample_queue = xQueueCreate(32, sizeof(SampleData));
g_buffer_mutex = xSemaphoreCreateMutex();
xTaskCreate(sample_task, "sample", 2048, NULL, 4, NULL);
xTaskCreate(report_task, "report", 4096, NULL, 3, NULL);
}import asyncio
import struct
import time
from typing import Dict, List, Optional, Tuple
from dataclasses import dataclass, field
from datetime import datetime
import logging
logging.basicConfig(level=logging.INFO, format='%(asctime)s [%(levelname)s] %(message)s')
logger = logging.getLogger(__name__)
@dataclass
class SensorSession:
"""传感器会话"""
device_id: int
reader: asyncio.StreamReader
writer: asyncio.StreamWriter
cabinet_id: str = "" # 配电柜编号
location: str = "" # 位置描述
connected_time: float = 0.0
last_heartbeat: float = 0.0
last_report: float = 0.0
total_reports: int = 0
total_samples: int = 0
total_bytes: int = 0
sequence: int = 0
last_seq: int = 0
packet_loss: int = 0
status: str = "CONNECTED"
buffer: List[Dict] = field(default_factory=list) # 数据缓冲
@dataclass
class SampleData:
timestamp: int
temperature: float
humidity: float
status: int
class PowerCabinetServer:
"""配电柜集群监测服务器"""
def __init__(self, host: str = "0.0.0.0", port: int = 5020,
max_connections: int = 512,
heartbeat_timeout: float = 90.0,
data_callback=None,
alarm_callback=None):
self.host = host
self.port = port
self.max_connections = max_connections
self.heartbeat_timeout = heartbeat_timeout
self.data_callback = data_callback
self.alarm_callback = alarm_callback
self.sessions: Dict[int, SensorSession] = {}
self.cabinet_map: Dict[str, int] = {} # cabinet_id -> device_id
self.running = False
async def _handle_client(self, reader: asyncio.StreamReader,
writer: asyncio.StreamWriter):
"""处理单个客户端连接"""
peer = writer.get_extra_info('peername')
device_id = None
try:
# 握手:等待设备注册
# 第一个报文应为注册报文(或直接在批量报文中解析 device_id)
data = await asyncio.wait_for(reader.read(1024), timeout=10.0)
if len(data) < 12:
logger.warning(f"Invalid handshake from {peer}")
writer.close()
return
# 解析注册信息(简化:从第一个批量报文中获取 device_id)
if data[0] == 0xAA and data[1] == 0x55:
device_id = struct.unpack_from('<I', data, 4)[0]
else:
logger.warning(f"Unknown packet from {peer}")
writer.close()
return
# 创建会话
if device_id in self.sessions:
# 设备重连,关闭旧连接
old = self.sessions[device_id]
logger.info(f"Device {device_id} reconnected, closing old session")
try:
old.writer.close()
except:
pass
session = SensorSession(
device_id=device_id,
reader=reader,
writer=writer,
connected_time=time.time(),
last_heartbeat=time.time(),
last_report=time.time()
)
self.sessions[device_id] = session
logger.info(f"Device {device_id} connected from {peer}, "
f"total sessions: {len(self.sessions)}")
# 持续读取数据
await self._read_loop(session)
except asyncio.TimeoutError:
logger.warning(f"Handshake timeout from {peer}")
except (ConnectionResetError, BrokenPipeError) as e:
logger.info(f"Device {device_id} disconnected: {e}")
except Exception as e:
logger.error(f"Session error for device {device_id}: {e}")
finally:
if device_id and device_id in self.sessions:
session = self.sessions[device_id]
session.status = "DISCONNECTED"
# 保留会话信息但关闭连接
try:
writer.close()
await writer.wait_closed()
except:
pass
logger.info(f"Device {device_id} session closed, "
f"total reports: {session.total_reports}, "
f"total samples: {session.total_samples}")
async def _read_loop(self, session: SensorSession):
"""读取循环"""
while True:
# 读取报文头(至少 16 字节)
header = await reader_read_exact(session.reader, 16)
if not header:
break
if header[0] != 0xAA or header[1] != 0x55:
logger.warning(f"Device {session.device_id}: invalid header "
f"{header[0]:02x}{header[1]:02x}")
continue
# 解析报文头
length, device_id, sequence, timestamp, count = \
struct.unpack_from('<HIHIB', header, 2)
# 读取剩余数据
remaining = length - 16
if remaining < 0 or remaining > 1024:
logger.warning(f"Device {session.device_id}: invalid length {length}")
break
payload = await reader_read_exact(session.reader, remaining)
if not payload:
break
# 更新会话信息
session.last_report = time.time()
session.sequence = sequence
session.total_reports += 1
# 解析数据
await self._parse_batch_report(session, timestamp, count, payload)
async def _parse_batch_report(self, session: SensorSession,
timestamp: int, count: int,
payload: bytes):
"""解析批量上报报文"""
offset = 0
samples = []
for i in range(count):
if offset + 8 > len(payload):
break
ts, temp_raw, hum_raw, status = struct.unpack_from(
'<IhhBB', payload, offset
)
offset += 8
sample = {
'device_id': session.device_id,
'cabinet_id': session.cabinet_id,
'location': session.location,
'timestamp': ts,
'temperature': temp_raw * 0.1,
'humidity': hum_raw * 0.1,
'status': status,
'received_at': time.time()
}
samples.append(sample)
session.total_samples += 1
# 数据回调
if self.data_callback:
await self.data_callback(samples)
# 告警检测
for s in samples:
if self.alarm_callback:
await self.alarm_callback(s)
# 心跳更新
session.last_heartbeat = time.time()
# 序列号连续性检查
if session.last_seq > 0:
expected = (session.last_seq + 1) & 0xFFFFFFFF
if sequence != expected:
loss = (sequence - expected) & 0xFFFFFFFF
if loss < 100: # 防止回绕误判
session.packet_loss += loss
logger.debug(f"Device {session.device_id}: "
f"seq jump {session.last_seq} -> {sequence}, "
f"loss={loss}")
session.last_seq = sequence
async def _heartbeat_monitor(self):
"""心跳监控任务"""
while self.running:
await asyncio.sleep(10)
now = time.time()
disconnected = []
for device_id, session in self.sessions.items():
if session.status == "CONNECTED":
idle_time = now - session.last_heartbeat
if idle_time > self.heartbeat_timeout:
logger.warning(
f"Device {device_id} heartbeat timeout "
f"({idle_time:.0f}s), closing"
)
disconnected.append(device_id)
try:
session.writer.close()
except:
pass
for did in disconnected:
if did in self.sessions:
self.sessions[did].status = "TIMEOUT"
async def _stats_reporter(self):
"""统计报告"""
while self.running:
await asyncio.sleep(60)
connected = sum(1 for s in self.sessions.values()
if s.status == "CONNECTED")
total_samples = sum(s.total_samples for s in self.sessions.values())
total_reports = sum(s.total_reports for s in self.sessions.values())
total_loss = sum(s.packet_loss for s in self.sessions.values())
logger.info(
f"[STATS] Connected: {connected}/{len(self.sessions)} | "
f"Reports: {total_reports} | Samples: {total_samples} | "
f"Packet Loss: {total_loss}"
)
async def start(self):
"""启动服务器"""
self.running = True
# 启动心跳监控
asyncio.create_task(self._heartbeat_monitor())
# 启动统计报告
asyncio.create_task(self._stats_reporter())
server = await asyncio.start_server(
self._handle_client, self.host, self.port
)
addrs = ', '.join(str(sock.getsockname()) for sock in server.sockets)
logger.info(f"Power Cabinet Server listening on {addrs}, "
f"max connections: {self.max_connections}")
async with server:
await server.serve_forever()
def get_status(self) -> Dict:
"""获取所有会话状态"""
status = {}
for did, session in self.sessions.items():
status[did] = {
'cabinet_id': session.cabinet_id,
'location': session.location,
'status': session.status,
'connected_time': session.connected_time,
'last_heartbeat': session.last_heartbeat,
'total_reports': session.total_reports,
'total_samples': session.total_samples,
'packet_loss': session.packet_loss,
'idle_time': time.time() - session.last_heartbeat
}
return status
def stop(self):
self.running = False
async def reader_read_exact(reader: asyncio.StreamReader, n: int) -> Optional[bytes]:
"""精确读取 n 字节"""
try:
return await reader.readexactly(n)
except (asyncio.IncompleteReadError, ConnectionError):
return None
# 数据回调示例
async def on_data(samples: List[Dict]):
"""数据接收回调"""
for s in samples:
print(f"[{datetime.fromtimestamp(s['timestamp'])}] "
f"柜{s['cabinet_id']} T={s['temperature']:.1f}℃ "
f"H={s['humidity']:.1f}%RH")
async def on_alarm(sample: Dict):
"""告警回调"""
if sample['status'] & 0x01: # 高温告警
print(f"⚠️ ALARM: 柜{sample['cabinet_id']} 高温 "
f"{sample['temperature']:.1f}℃")
if sample['status'] & 0x04: # 高湿告警
print(f"⚠️ ALARM: 柜{sample['cabinet_id']} 高湿 "
f"{sample['humidity']:.1f}%RH")
# 主程序
async def main():
server = PowerCabinetServer(
host="0.0.0.0",
port=5020,
max_connections=512,
heartbeat_timeout=90.0,
data_callback=on_data,
alarm_callback=on_alarm
)
# 配置配电柜映射
server.cabinet_map = {
"PS-01": 1001, "PS-02": 1002, "PS-03": 1003,
# ...
}
try:
await server.start()
except KeyboardInterrupt:
server.stop()
print("Server stopped")
if __name__ == "__main__":
asyncio.run(main())┌────────────────────────────────────────────────────────────────┐
│ 批量上报报文格式 │
├──────────┬──────┬──────────┬──────────┬─────────┬─────┬───────┤
│ 帧头 │ 长度 │ 设备 ID │ 序列号 │ 时间戳 │ 数量│ 保留 │
│ 2B │ 2B │ 4B │ 4B │ 4B │ 1B │ 3B │
│ 0xAA55 │ N │ uint32 │ uint32 │ uint32 │ u8 │ - │
├──────────┴──────┴──────────┴──────────┴─────────┴─────┴───────┤
│ 数据点(变长) │
├──────────┬──────────┬──────────┬───────────────────────────────┤
│ 时间戳 │ 温度 │ 湿度 │ 状态 │
│ 4B │ 2B(int16)│ 2B(int16)│ 1B │
│ uint32 │ ×0.1℃ │ ×0.1%RH │ bit0:高温 bit1:低温 bit2:高湿 │
│ │ │ │ bit3:低湿 bit4:传感器故障 │
├──────────┼──────────┼──────────┼───────────────────────────────┤
│ ... │ ... │ ... │ ... │
├──────────┴──────────┴──────────┴───────────────────────────────┤
│ CRC16 (2B) │
└────────────────────────────────────────────────────────────────┘
示例(1 个数据点):
AA 55 1A 00 00 00 03 E9 01 00 00 00 5E 8A A2 60
01 3E 8F 00 02 0D 00 00 00 00 00 00 00 00 00 00
01 00 00 00 00 00 00 00 00 00 00 00 00 00 00 00
XX XX
→ 长度 0x001A = 26 字节
→ 设备 ID 0x000003E9 = 1001
→ 序列号 0x00000001
→ 时间戳 0x60A28A5E
→ 数量 1
→ 数据点:时间戳 3E8F0000(小端) = 0x00008F3E, 温度 0100(小端)=0x0001,
湿度 0D00(小端)=0x000D, 状态 00
→ CRC16┌──────────┬──────────┬──────────┐
│ 帧头 │ 设备 ID │ 序列号 │
│ 2B │ 4B │ 2B │
│ 0xAA55 │ uint32 │ uint16 │
└──────────┴──────────┴──────────┘
无 CRC,无长度字段(固定 8 字节)
用于保持连接活跃,防止中间设备/NAT 超时断开Bit | 含义 | 触发条件 |
|---|---|---|
0 | 高温告警 | 温度 > 设定上限 |
1 | 低温告警 | 温度 < 设定下限 |
2 | 高湿告警 | 湿度 > 设定上限 |
3 | 低湿告警 | 湿度 < 设定下限 |
4 | 传感器故障 | 探头读取失败 |
5 | 通信异常 | 网络中断 |
6 | 电压异常 | 供电电压超范围 |
7 | 保留 | - |
方案 A:每台柜内独立 PoE 交换机
优点:每台柜独立供电、独立网络,单点故障不扩散
缺点:柜内设备多(不止温湿度传感器),成本高
方案 B:区域汇聚 + 长线
优点:成本最低
缺点:超 100m 需中继,故障面大
方案 C:柜内微型 PoE 交换机(推荐)
每台配电柜内装一台 5-8 口工业 PoE 交换机
传感器 + 智能断路器 + 局放传感器 共用
上行汇聚到配电室接入交换机
成本与可靠性平衡最好接入交换机端口配置:
→ 端口安全:MAC 限制 1-2 个(柜内只有 1-2 台网络设备)
→ BPDU Guard:启用
→ Storm Control:广播抑制 5%
→ PoE 优先级:Critical(温湿度监测是安全相关)
→ IGMP Snooping:如果后续加视频巡检
→ Port Fast:启用(接入端口直接转发)
汇聚交换机:
→ VLAN 隔离:监测 VLAN 独立
→ DHCP Snooping:防止私接 DHCP
→ DAI(动态 ARP 检测):防止 ARP 欺骗
→ ACL:只允许监测 VLAN 到采集服务器 5020 端口并发 96 台传感器:
CPU:2 核(asyncio 单线程即可处理)
内存:512MB(每个会话约 2-4KB)
网络:10Mbps(绰绰有余)
文件描述符:至少 512(96 连接 + 系统预留)
并发 500+ 台:
CPU:4-8 核
内存:2-4GB
网络:100Mbps
文件描述criptor:2048+
建议分片部署(按配电室分区)from influxdb_client import InfluxDBClient, Point, WriteOptions
from influxdb_client.client.write_api import SYNCHRONOUS
class DataWriter:
def __init__(self, url, token, org, bucket):
self.client = InfluxDBClient(url=url, token=token, org=org)
self.write_api = self.client.write_api(write_options=SYNCHRONOUS)
self.bucket = bucket
self.org = org
async def write_batch(self, samples: List[Dict]):
"""批量写入"""
points = []
for s in samples:
point = Point("power_cabinet_temp_hum") \
.tag("device_id", str(s['device_id'])) \
.tag("cabinet_id", s['cabinet_id']) \
.tag("location", s['location']) \
.field("temperature", s['temperature']) \
.field("humidity", s['humidity']) \
.field("status", s['status']) \
.time(datetime.fromtimestamp(s['timestamp']))
points.append(point)
# 批量写入
self.write_api.write(bucket=self.bucket, org=self.org, record=points)ALARM_RULES = {
'high_temp': {
'threshold': 45.0, # ℃
'duration': 300, # 连续 5 分钟
'severity': 'CRITICAL',
'action': ['sound_alarm', 'send_sms', 'trigger_ac']
},
'high_humidity': {
'threshold': 80.0, # %RH
'duration': 600, # 连续 10 分钟
'severity': 'WARNING',
'action': ['send_sms', 'start_dehumidifier']
},
'sensor_offline': {
'threshold': 90, # 秒(心跳超时)
'duration': 0,
'severity': 'WARNING',
'action': ['send_sms']
}
}
class AlarmEngine:
def __init__(self, rules: Dict):
self.rules = rules
self.state = {} # 持续时长跟踪
async def check(self, sample: Dict):
device_id = sample['device_id']
# 温度告警
if sample['temperature'] > self.rules['high_temp']['threshold']:
self._increment(device_id, 'high_temp')
if self.state[device_id]['high_temp'] >= self.rules['high_temp']['duration']:
await self._trigger_alarm('high_temp', sample)
else:
self._reset(device_id, 'high_temp')
# 湿度告警
if sample['humidity'] > self.rules['high_humidity']['threshold']:
self._increment(device_id, 'high_humidity')
if self.state[device_id]['high_humidity'] >= self.rules['high_humidity']['duration']:
await self._trigger_alarm('high_humidity', sample)
else:
self._reset(device_id, 'high_humidity')问题 | 原因 | 解决方案 |
|---|---|---|
传感器连接不上服务器 | 防火墙拦截 5020 端口 | 放行端口,检查 ACL |
连接建立后很快断开 | KeepAlive 参数不一致 | 统一空闲时间、探测间隔、探测次数 |
数据延迟大 | 上报周期过长 | 缩短到 10-15s(配电柜场景可接受) |
批量报文粘包 | TCP 流无边界 | 报文头含长度字段,按长度读取 |
序列号不连续 | 网络丢包 | 检测丢包但不重传(UDP 语义) |
服务器内存增长 | 会话未清理 | 心跳超时主动关闭,定期清理 DISCONNECTED 会话 |
配电柜内交换机掉电 | PoE 预算不足 | 交换机纳入 UPS,或传感器改 DC 24V 供电 |
数据时间不对 | 传感器 RTC 未同步 | 服务器端用接收时间覆盖,或部署 NTP |
同柜多传感器 ID 冲突 | 配置错误 | 每台设备唯一 ID,配置文件版本管理 |
项目 | 标准 | 测试方法 |
|---|---|---|
连接数 | 96 台全部在线 | 查看服务器会话列表 |
数据到达率 | ≥99.9%(30 天) | 对比传感器本地缓存与服务器数据 |
上报延迟 | P99 < 35s(上报周期 30s) | 记录采样时间戳与接收时间差 |
心跳存活 | 90s 内必有心跳 | 监控 last_heartbeat 字段 |
序列号连续 | 丢包率 < 0.1% | 服务器统计 packet_loss |
断线恢复 | <30s 自动重连 | 拔网线后计时 |
数据完整性 | 本地缓存与服务器一致 | 导出对比 |
告警响应 | <5s(从越限到触发) | 模拟越限,记录时间差 |
独立配电柜集群监测用"长连接 + 批量上报",本质是把请求-响应模型改成发布-订阅模型。 传感器采集完自己打包、自己发送,服务器只负责收和存,双方都不用等对方。 关键在三个设计:报文带长度字段解决粘包、序列号检测丢包、KeepAlive 维持长连接。 96 台配电柜、30 秒上报一次,服务器一个端口、一个 asyncio 循环就搞定——比轮询模式省了 97% 的网络交互。
关键词:独立配电柜、配电柜集群、TCP/IP长连接、批量数据上传、温湿度监测、KeepAlive、批量上报、配电室监控、传感器网络、数据缓存、断线重连、心跳机制、序列号、粘包处理
标签:#独立配电柜 #配电柜集群 #TCPIP长连接 #批量数据上传 #温湿度监测 #KeepAlive #Nagle #数据聚合 #配电室 #动环监控 #工业以太网 #传感器网络
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。