首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >独立配电柜集群监测:TCP/IP长连接,批量温湿度数据上传服务器

独立配电柜集群监测:TCP/IP长连接,批量温湿度数据上传服务器

原创
作者头像
盛世宏博小可
发布于 2026-09-28 10:34:00
发布于 2026-09-28 10:34:00
690
举报

独立配电柜集群监测:TCP/IP长连接,批量温湿度数据上传服务器

独立配电柜 #配电柜集群 #TCP/IP长连接 #批量数据上传 #温湿度监测 #KeepAlive #Nagle #数据聚合 #配电室 #动环监控 #工业以太网

先说清场景:独立配电柜集群指的是分散在厂区各处的配电室/配电间里,成排独立运行的低压配电柜。每个柜内装一台以太网温湿度变送器,监测柜内温湿度——这是判断柜内是否有过热隐患、是否凝露、是否有电气接点异常升温的关键指标。配电柜的特点是数量多、点位分散、单点数据量小但采集频率要求不高,最适合用"长连接 + 批量上传"的模式。

先说一个真实案例:

某汽车零部件厂,4 个配电室共 96 台独立配电柜,每台柜内装 1 台以太网温湿度变送器。 初始方案:采集服务器对 96 台设备各开一个 TCP 连接,每 10s 轮询一次,每次读 2 个寄存器。 问题:

  • 96 个并发 TCP 连接,服务器端口和文件描述符占用高
  • 每次轮询 96 个独立请求,网络小包泛滥(每个请求 12 字节 + 响应 13 字节)
  • 采集周期 10s 但完成一轮需要 8-12s(取决于响应时间),实际周期不稳定
  • 交换机端口出现大量小包,广播域流量占比上升 优化方案:改为传感器主动长连接 + 批量上报模式
  • 每台传感器启动后主动连接采集服务器,维持长连接(KeepAlive)
  • 传感器本地缓存,每 30s 打包一次(温度+湿度+状态+时标),批量上报
  • 服务器单端口监听,统一接收所有传感器的数据 结果:服务器连接数从 96 降到 96(长连接复用,但每个设备仍是独立连接), 关键优化是减少了请求-响应交互次数:从每秒 9.6 次交互降到每 30 秒 1 次,网络包数减少 97%。

一、架构设计

1. 整体架构

代码语言:javascript
复制
┌─────────────────────────────────────────────────────────────────────┐
│                        采集服务器(中心侧)                          │
│  ┌────────────────────────────────────────────────────────────────┐ │
│  │              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                 │
└─────────────────────────────────────────────────────────────────────┘

2. 两种工作模式对比

维度

服务器轮询模式(旧)

传感器主动上报模式(新)

连接方向

服务器 → 传感器(Server 端被动等待请求)

传感器 → 服务器(Client 端主动发起)

连接数量

96 个独立连接

96 个长连接(但连接一旦建立不释放)

请求方式

请求-响应模型

传感器单向推送(或请求-响应混合)

数据实时性

取决于轮询周期

采集即上报,延迟更低

网络开销

每次 2 个小包(请求+响应)

批量打包,一个包含多台设备数据

服务器压力

需管理轮询调度、超时、重试

只需维护连接表、解析数据

适用场景

设备数量少、需频繁控制

设备数量多、数据量小、单向采集

核心转变:从"服务器问,传感器答"变成"传感器采集完,主动告诉服务器"。


二、传感器端实现(TCP Client + 批量上报)

1. 固件架构

代码语言:javascript
复制
┌─────────────────────────────────────────────────────────────┐
│                    传感器固件架构                             │
│                                                             │
│  ┌────────────┐  ┌────────────┐  ┌───────────────────────┐ │
│  │ 采样任务    │→│ 数据缓存    │→│ TCP Client 上报任务    │ │
│  │ 2s 采集一次 │  │ Ring Buffer│  │ 长连接管理             │ │
│  │ 温湿度+I/O │  │ 容量: 128  │  │ 批量打包               │ │
│  └────────────┘  └────────────┘  │ 心跳/重连              │ │
│                                  └───────────────────────┘ │
│                                             │               │
│                                  ┌──────────┴────────────┐ │
│                                  │   TCP/IP 协议栈        │ │
│                                  │   LwIP + FreeRTOS     │ │
│                                  └────────────────────────┘ │
└─────────────────────────────────────────────────────────────┘

2. 关键参数配置

参数

值

说明

采集周期

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 台

超容量后覆盖最旧数据

3. 传感器端核心代码(C/FreeRTOS)

代码语言:javascript
复制
/* 配电柜温湿度传感器 - 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);
}

三、服务器端实现(Python asyncio)

1. 服务器核心代码

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

四、协议设计要点

1. 报文格式

代码语言:javascript
复制
┌────────────────────────────────────────────────────────────────┐
│                    批量上报报文格式                             │
├──────────┬──────┬──────────┬──────────┬─────────┬─────┬───────┤
│ 帧头     │ 长度 │ 设备 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

2. 心跳报文

代码语言:javascript
复制
┌──────────┬──────────┬──────────┐
│ 帧头     │ 设备 ID  │ 序列号   │
│ 2B       │ 4B       │ 2B       │
│ 0xAA55   │ uint32   │ uint16   │
└──────────┴──────────┴──────────┘

无 CRC,无长度字段(固定 8 字节)
用于保持连接活跃,防止中间设备/NAT 超时断开

3. 状态字定义

Bit

含义

触发条件

0

高温告警

温度 > 设定上限

1

低温告警

温度 < 设定下限

2

高湿告警

湿度 > 设定上限

3

低湿告警

湿度 < 设定下限

4

传感器故障

探头读取失败

5

通信异常

网络中断

6

电压异常

供电电压超范围

7

保留

-


五、网络与部署优化

1. 配电柜内组网方式

代码语言:javascript
复制
方案 A:每台柜内独立 PoE 交换机
  优点:每台柜独立供电、独立网络,单点故障不扩散
  缺点:柜内设备多(不止温湿度传感器),成本高

方案 B:区域汇聚 + 长线
  优点:成本最低
  缺点:超 100m 需中继,故障面大

方案 C:柜内微型 PoE 交换机(推荐)
  每台配电柜内装一台 5-8 口工业 PoE 交换机
  传感器 + 智能断路器 + 局放传感器 共用
  上行汇聚到配电室接入交换机
  
  成本与可靠性平衡最好

2. 交换机配置

代码语言:javascript
复制
接入交换机端口配置:
  → 端口安全: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 端口

3. 服务器资源配置

代码语言:javascript
复制
并发 96 台传感器:
  CPU:2 核(asyncio 单线程即可处理)
  内存:512MB(每个会话约 2-4KB)
  网络:10Mbps(绰绰有余)
  文件描述符:至少 512(96 连接 + 系统预留)

并发 500+ 台:
  CPU:4-8 核
  内存:2-4GB
  网络:100Mbps
  文件描述criptor:2048+
  建议分片部署(按配电室分区)

六、数据持久化与告警

1. 写入时序数据库(InfluxDB 示例)

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

2. 告警规则

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

目录
  • 独立配电柜集群监测:TCP/IP长连接,批量温湿度数据上传服务器
    • 独立配电柜 #配电柜集群 #TCP/IP长连接 #批量数据上传 #温湿度监测 #KeepAlive #Nagle #数据聚合 #配电室 #动环监控 #工业以太网
    • 一、架构设计
      • 1. 整体架构
      • 2. 两种工作模式对比
    • 二、传感器端实现(TCP Client + 批量上报)
      • 1. 固件架构
      • 2. 关键参数配置
      • 3. 传感器端核心代码(C/FreeRTOS)
    • 三、服务器端实现(Python asyncio)
      • 1. 服务器核心代码
    • 四、协议设计要点
      • 1. 报文格式
      • 2. 心跳报文
      • 3. 状态字定义
    • 五、网络与部署优化
      • 1. 配电柜内组网方式
      • 2. 交换机配置
      • 3. 服务器资源配置
    • 六、数据持久化与告警
      • 1. 写入时序数据库(InfluxDB 示例)
      • 2. 告警规则
    • 七、常见问题与排查
    • 八、验收标准
    • 九、一句话总结
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档