
电厂配电室环境数据不是"看看就行"的辅助信息。在事故分析、设备状态评估、运行规程执行中,环境数据经常作为关键佐证材料被调取——比如某次开关柜异常发热事件,事后需要回溯当时的温度曲线、湿度变化、与设备动作时间线的对应关系。
这意味着两个硬约束:

而电厂环境的现实是:
风险场景 | 发生频率 | 影响 |
|---|---|---|
网络分区(交换机堆叠分裂) | 偶发,季度级 | 持续数分钟到数小时 |
光纤链路中断(施工、外力) | 低,但每次持续长 | 数小时到数天 |
采集服务器维护/重启 | 每周或每两周一次 | 数分钟到数小时 |
POE交换机端口故障 | 偶发 | 单点设备中断 |
传感器到交换机网线故障 | 偶发 | 单点设备中断 |
这些问题在普通机房可以容忍,在电厂意味着数据空白。因此,本地缓存 + TCP断点续传不是"锦上添花",而是系统设计的必选项。
┌─────────────────────────────────────────────────────────────────────┐
│ 配电室现场(生产控制大区) │
│ │
│ ┌────────────┐ ┌────────────┐ ┌────────────┐ │
│ │ 温湿度 │ │ 温湿度 │ │ 温湿度 │ │
│ │ 记录仪 │ │ 记录仪 │ │ 记录仪 │ │
│ │ (本地存储) │ │ (本地存储) │ │ (本地存储) │ │
│ └─────┬──────┘ └─────┬──────┘ └─────┬──────┘ │
│ │ │ │ │
│ └────────────────┼────────────────┘ │
│ │ TCP 502 (Modbus TCP) │
│ ┌────▼─────┐ │
│ │ POE交换机 │ │
│ └────┬─────┘ │
└──────────────────────────┼──────────────────────────────────────────┘
│
┌────────────▼────────────┐
│ 正向隔离装置 │
│ (单向,生产区→管理区) │
└────────────┬────────────┘
│
┌──────────────────────────┼──────────────────────────────────────────┐
│ 管理信息大区 │ │
│ │ │
│ ┌────────────────────────▼────────────────────────────────────┐ │
│ │ 采集服务器 │ │
│ │ ┌──────────────────────────────────────────────────────┐ │ │
│ │ │ 断点续传管理器 │ │ │
│ │ │ - 维护每台设备的同步状态(最后同步时间戳) │ │ │
│ │ │ - 网络恢复后自动发起补传请求 │ │ │
│ │ │ - 补传限速,避免影响实时采集 │ │ │
│ │ └──────────────────────────────────────────────────────┘ │ │
│ │ │ │
│ │ ┌──────────────────────────────────────────────────────┐ │ │
│ │ │ 实时采集 + 缓存双写 │ │ │
│ │ │ - 正常时:写数据库 + 更新同步状态 │ │ │
│ │ │ - 异常时:只写本地缓存 │ │ │
│ │ └──────────────────────────────────────────────────────┘ │ │
│ └──────────────────────────────────────────────────────────────┘ │
│ │ │
│ ┌───────────▼───────────┐ │
│ │ 时序数据库 │ │
│ │ TDengine/InfluxDB │ │
│ └───────────────────────┘ │
└─────────────────────────────────────────────────────────────────────┘核心设计原则:
记录仪内置存储(Flash或EEPROM),采用循环存储方式:
存储布局(假设512KB Flash):
┌─────────────────────────────────────────────┐
│ 配置区 (4KB) │
│ - 采样间隔、报警阈值、设备ID │
├─────────────────────────────────────────────┤
│ 索引区 (8KB) │
│ - 每段数据的起始时间戳、偏移量、长度 │
├─────────────────────────────────────────────┤
│ 数据区 (496KB) │
│ - 循环存储,新数据覆盖最旧数据 │
│ - 每条记录固定32字节 │
└─────────────────────────────────────────────┘数据记录格式(32字节):
┌──────────┬──────────┬─────────┬─────────┬─────────┬────────┐
│ 时间戳 │ 温度 │ 湿度 │ 状态字 │ 序列号 │ CRC16 │
│ (4字节) │ (4字节) │ (4字节) │ (2字节) │ (4字节) │(2字节) │
│ Unix时间戳│ float ℃│ float %RH│ │ │ │
└──────────┴──────────┴─────────┴─────────┴─────────┴────────┘存储容量计算:
单条记录:32字节
可用空间:496KB ≈ 496 × 1024 / 32 = 15872 条
采样间隔:60秒
容量:15872 × 60秒 = 952320秒 ≈ 11天即记录仪本地可存储约11天的历史数据。对于网络中断场景,这个时间窗口足够覆盖绝大多数情况。
标准温湿度变送器通常只提供实时数据读取(输入寄存器),记录仪需要额外提供历史数据读取功能。在私有寄存器中扩展:
读取历史数据命令(自定义功能码或Holding Register写入):
步骤1:设置查询参数(写入Holding Registers)
寄存器 0x0100: 起始时间戳(高16位)
寄存器 0x0101: 起始时间戳(低16位)
寄存器 0x0102: 结束时间戳(高16位)
寄存器 0x0103: 结束时间戳(低16位)
寄存器 0x0104: 最大返回记录数(默认100)
寄存器 0x0105: 查询触发标志(写入1触发查询)
步骤2:读取查询结果(读取Holding Registers)
寄存器 0x0200: 查询结果状态
0x0000 = 查询成功
0x0001 = 查询中
0x0002 = 无数据
0x0003 = 参数错误
寄存器 0x0201: 返回记录数
寄存器 0x0202: 下一批起始时间戳(高16位)
寄存器 0x0203: 下一批起始时间戳(低16位)
步骤3:读取数据块(读取Input Registers)
寄存器 0x0300 开始,每16个寄存器(32字节)为一条记录
共最多100条记录(1600个寄存器)#define MAX_RECORDS_PER_QUERY 100
#define RECORD_SIZE 32
#define DATA_FLASH_START_ADDR 0x08010000 // Flash起始地址
#define DATA_FLASH_SIZE 0x7C000 // 496KB
typedef struct {
uint32_t timestamp;
float temp;
float humi;
uint16_t status;
uint16_t seq;
uint16_t crc;
} EnvRecord;
static uint32_t write_offset = 0; // 当前写入偏移
static uint32_t oldest_ts = 0; // 最旧记录时间戳
static uint32_t newest_ts = 0; // 最新记录时间戳
void env_record_init(void) {
// 初始化时扫描Flash,找到有效数据范围
EnvRecord rec;
uint32_t offset = 0;
while (offset < DATA_FLASH_SIZE) {
flash_read(DATA_FLASH_START_ADDR + offset, &rec, sizeof(rec));
if (rec.crc == calc_crc16((uint8_t*)&rec, RECORD_SIZE - 2)) {
if (oldest_ts == 0 || rec.timestamp < oldest_ts)
oldest_ts = rec.timestamp;
if (rec.timestamp > newest_ts)
newest_ts = rec.timestamp;
write_offset = offset + RECORD_SIZE;
} else {
break; // 遇到无效记录,停止扫描
}
offset += RECORD_SIZE;
}
if (write_offset >= DATA_FLASH_SIZE) {
write_offset = 0; // 循环覆盖
}
}
void env_record_write(float temp, float humi, uint16_t status) {
EnvRecord rec;
rec.timestamp = get_unix_time();
rec.temp = temp;
rec.humi = humi;
rec.status = status;
rec.seq = get_seq();
rec.crc = calc_crc16((uint8_t*)&rec, RECORD_SIZE - 2);
// 写入Flash
flash_write(DATA_FLASH_START_ADDR + write_offset, &rec, sizeof(rec));
// 更新指针
write_offset += RECORD_SIZE;
if (write_offset >= DATA_FLASH_SIZE) {
write_offset = 0;
// 覆盖最旧数据,更新oldest_ts
EnvRecord next;
flash_read(DATA_FLASH_START_ADDR + write_offset, &next, sizeof(next));
oldest_ts = next.timestamp;
}
newest_ts = rec.timestamp;
}package main
import (
"database/sql"
"encoding/binary"
"fmt"
"log"
"net"
"sync"
"time"
_ "github.com/mattn/go-sqlite3"
"github.com/go-modbus/modbus"
)
// 设备同步状态
type SyncState struct {
DeviceID string
LastSyncTS time.Time // 上次成功同步的时间戳
OldestLocalTS time.Time // 记录仪中最旧数据的时间戳
NewestLocalTS time.Time // 记录仪中最新数据的时间戳
Syncing bool // 是否正在补传
LastError error
}
var (
syncStates = make(map[string]*SyncState)
statesMu sync.RWMutex
db *sql.DB
)
func initSyncState(deviceID string) *SyncState {
statesMu.Lock()
defer statesMu.Unlock()
if s, ok := syncStates[deviceID]; ok {
return s
}
// 从数据库加载上次同步状态
var lastSyncUnix int64
row := db.QueryRow(
"SELECT COALESCE(MAX(ts), 0) FROM env_data WHERE device_id = ?",
deviceID)
row.Scan(&lastSyncUnix)
s := &SyncState{
DeviceID: deviceID,
LastSyncTS: time.Unix(lastSyncUnix, 0),
}
syncStates[deviceID] = s
return s
}func pollDeviceWithCache(dev DeviceConfig) {
ticker := time.NewTicker(10 * time.Second)
defer ticker.Stop()
// 本地缓存(服务器侧,用于采集服务自身异常时)
cacheFile := fmt.Sprintf("/var/cache/envd/%s.db", dev.ID)
for range ticker.C {
data, err := readRealtime(dev)
if err != nil {
log.Printf("[%s] 实时读取失败: %v", dev.ID, err)
// 标记设备离线,但不影响补传逻辑
markDeviceOffline(dev.ID)
continue
}
// 尝试写入数据库
if err := writeToDB(data); err != nil {
// 数据库不可用时,写入本地缓存
log.Printf("[%s] 数据库写入失败,缓存: %v", dev.ID, err)
cacheToFile(cacheFile, data)
} else {
// 数据库写入成功,更新同步状态
updateSyncState(dev.ID, data.Timestamp)
// 尝试将本地缓存flush到数据库
flushCacheToFile(cacheFile)
}
// 检查是否需要补传
go checkAndBackfill(dev)
}
}
func checkAndBackfill(dev DeviceConfig) {
state := initSyncState(dev.ID)
statesMu.Lock()
if state.Syncing {
statesMu.Unlock()
return // 已经在补传中
}
state.Syncing = true
statesMu.Unlock()
defer func() {
statesMu.Lock()
state.Syncing = false
statesMu.Unlock()
}()
// 查询记录仪的时间范围
oldest, newest, err := queryRecorderRange(dev)
if err != nil {
log.Printf("[%s] 查询记录仪范围失败: %v", dev.ID, err)
return
}
state.OldestLocalTS = oldest
state.NewestLocalTS = newest
// 判断是否需要补传
if state.LastSyncTS.Before(newest.Add(-5 * time.Second)) {
// 有数据需要补传
log.Printf("[%s] 开始补传: %v -> %v",
dev.ID, state.LastSyncTS, newest)
backfill(dev, state.LastSyncTS, newest)
}
}func backfill(dev DeviceConfig, start, end time.Time) error {
const batchSize = 100
const rateLimit = 10 * time.Millisecond // 每批间隔,限速
current := start
for current.Before(end) {
// 设置查询参数
if err := setQueryRange(dev, current, end, batchSize); err != nil {
return fmt.Errorf("set query range: %w", err)
}
// 等待记录仪准备数据
time.Sleep(100 * time.Millisecond)
// 读取查询结果状态
status, count, nextTS, err := getQueryResult(dev)
if err != nil {
return fmt.Errorf("get query result: %w", err)
}
switch status {
case 0x0000: // 成功
if count == 0 {
current = nextTS
continue
}
// 读取数据块
records, err := readRecordBatch(dev, count)
if err != nil {
return fmt.Errorf("read batch: %w", err)
}
// 写入数据库
if err := batchWriteToDB(dev.ID, records); err != nil {
return fmt.Errorf("batch write: %w", err)
}
log.Printf("[%s] 补传 %d 条 (%v -> %v)",
dev.ID, count, records[0].Timestamp, records[len(records)-1].Timestamp)
current = nextTS
case 0x0001: // 查询中,等待
time.Sleep(200 * time.Millisecond)
continue
case 0x0002: // 无数据
current = end // 跳过
case 0x0003: // 参数错误
return fmt.Errorf("invalid query parameters")
default:
return fmt.Errorf("unknown status: 0x%04X", status)
}
// 限速,避免影响实时采集
time.Sleep(rateLimit)
}
log.Printf("[%s] 补传完成", dev.ID)
return nil
}
func setQueryRange(dev DeviceConfig, start, end time.Time, maxCount uint16) error {
handler := modbus.NewTCPClientHandler(fmt.Sprintf("%s:%d", dev.IP, dev.Port))
handler.SlaveId = dev.SlaveID
handler.Timeout = 3 * time.Second
if err := handler.Connect(); err != nil {
return err
}
defer handler.Close()
client := modbus.NewClient(handler)
// 写入起始时间戳
startUnix := uint32(start.Unix())
if _, err := client.WriteMultipleRegisters(0x0100, 2,
uint16ToBytes(startUnix)); err != nil {
return err
}
// 写入结束时间戳
endUnix := uint32(end.Unix())
if _, err := client.WriteMultipleRegisters(0x0102, 2,
uint16ToBytes(endUnix)); err != nil {
return err
}
// 写入最大返回记录数
if _, err := client.WriteSingleRegister(0x0104, maxCount); err != nil {
return err
}
// 触发查询
if _, err := client.WriteSingleRegister(0x0105, 1); err != nil {
return err
}
return nil
}
func readRecordBatch(dev DeviceConfig, count uint16) ([]EnvRecord, error) {
handler := modbus.NewTCPClientHandler(fmt.Sprintf("%s:%d", dev.IP, dev.Port))
handler.SlaveId = dev.SlaveID
handler.Timeout = 5 * time.Second
if err := handler.Connect(); err != nil {
return nil, err
}
defer handler.Close()
client := modbus.NewClient(handler)
// 读取数据块:每条记录32字节=16个寄存器,最多100条=1600寄存器
regCount := count * 16
results, err := client.ReadInputRegisters(0x0300, regCount)
if err != nil {
return nil, err
}
records := make([]EnvRecord, count)
for i := uint16(0); i < count; i++ {
offset := i * 32
records[i].Timestamp = binary.BigEndian.Uint32(results[offset : offset+4])
records[i].Temp = math.Float32frombits(
binary.BigEndian.Uint32(results[offset+4 : offset+8]))
records[i].Humi = math.Float32frombits(
binary.BigEndian.Uint32(results[offset+8 : offset+12]))
records[i].Status = binary.BigEndian.Uint16(results[offset+12 : offset+14])
records[i].Seq = binary.BigEndian.Uint32(results[offset+14 : offset+18])
records[i].CRC = binary.BigEndian.Uint16(results[offset+18 : offset+20])
// 校验CRC
if records[i].CRC != crc16(results[offset:offset+18]) {
return nil, fmt.Errorf("CRC mismatch at record %d", i)
}
}
return records, nil
}// 补传调度器:错峰、限速、优先级
type BackfillScheduler struct {
queue chan BackfillTask
limiter chan struct{} // 信号量,限制并发补传数
}
type BackfillTask struct {
DeviceID string
Priority int // 1=高(数据缺口大), 2=中, 3=低
Start time.Time
End time.Time
}
func (s *BackfillScheduler) Start() {
// 最多同时3个补传任务,避免占满采集资源
s.limiter = make(chan struct{}, 3)
for task := range s.queue {
s.limiter <- struct{}{} // 获取令牌
go func(t BackfillTask) {
defer func() { <-s.limiter }()
dev := getDevice(t.DeviceID)
if err := backfill(dev, t.Start, t.End); err != nil {
log.Printf("[%s] 补传失败: %v, 将重试", t.DeviceID, err)
// 重新入队,延迟重试
time.AfterFunc(30*time.Second, func() {
s.queue <- t
})
}
}(task)
}
}
func (s *BackfillScheduler) Submit(t BackfillTask) {
// 高优先级插队到队列前端
if t.Priority == 1 {
go func() { s.queue <- t }()
} else {
s.queue <- t
}
}func detectGaps(deviceID string, start, end time.Time) []TimeGap {
// 查询数据库中该设备的时间序列
query := `
SELECT timestamp FROM env_data
WHERE device_id = ? AND timestamp BETWEEN ? AND ?
ORDER BY timestamp
`
rows, _ := db.Query(query, deviceID, start.Unix(), end.Unix())
defer rows.Close()
var timestamps []time.Time
for rows.Next() {
var ts int64
rows.Scan(&ts)
timestamps = append(timestamps, time.Unix(ts, 0))
}
// 检测缺口
var gaps []TimeGap
expectedInterval := 10 * time.Second // 采样间隔
for i := 1; i < len(timestamps); i++ {
gap := timestamps[i].Sub(timestamps[i-1])
if gap > expectedInterval*2 { // 超过2倍间隔视为缺口
gaps = append(gaps, TimeGap{
Start: timestamps[i-1].Add(expectedInterval),
End: timestamps[i].Add(-expectedInterval),
})
}
}
return gaps
}
type TimeGap struct {
Start time.Time
End time.Time
}
func (g TimeGap) Duration() time.Duration {
return g.End.Sub(g.Start)
}┌─────────────────────────────────────────────────────────────┐
│ 数据完整性报告 - 2025-01-15 │
├─────────────────────────────────────────────────────────────┤
│ │
│ 设备: T-6kVA-01 (6kV配电室-A列#3开关柜) │
│ │
│ 统计周期: 2025-01-14 00:00 ~ 2025-01-15 00:00 │
│ 预期数据点: 8640 (每10秒) │
│ 实际数据点: 8635 │
│ 完整率: 99.94% │
│ │
│ 数据缺口: │
│ [2025-01-14 03:15:00 ~ 03:18:20] 缺失20条 (网络中断) │
│ → 已自动补传 (记录仪缓存) │
│ │
│ [2025-01-14 14:22:00 ~ 14:22:10] 缺失1条 (设备重启) │
│ → 记录仪重启期间无法记录,无法补传 │
│ │
│ 补传统计: │
│ 自动补传次数: 3 │
│ 成功补传数据点: 45 │
│ 平均补传耗时: 1.2秒 │
│ │
└─────────────────────────────────────────────────────────────┘时间线:
T0: 网络正常,实时采集+入库
T1: 网络中断(交换机端口down)
→ 记录仪继续本地存储(不受影响)
→ 服务器侧采集失败,记录设备离线
T2: 网络恢复(端口up)
→ 服务器检测到设备在线
→ 触发补传流程
→ 查询记录仪 [T1, T2] 区间数据
→ 分批读取并写入数据库
T3: 补传完成
→ 更新同步状态为T2
→ 恢复正常实时采集场景:网络中断超过11天(记录仪存储容量上限),最旧数据被覆盖
处理:
1. 补传时发现记录仪 oldest_ts > 数据库 last_sync_ts
→ 说明中间有一段数据已被覆盖,无法补传
2. 记录数据缺口,生成告警
3. 从 oldest_ts 开始补传剩余数据
4. 完整性报告中明确标注"不可恢复缺口"处理:
1. 每批数据写入数据库成功后,立即更新 last_sync_ts
2. 网络中断时,补传中止
3. 网络恢复后,从 last_sync_ts 继续补传
4. 不会重复写入(数据库唯一索引:device_id + timestamp)单台记录仪:
Flash: 512KB
可用数据空间: 496KB
单条记录: 32字节
存储条数: ~15872条
采样间隔: 60秒
存储时长: ~11天
200台记录仪总存储: 200 × 512KB = 100MB (Flash,分布式)单台记录仪补传:
100条/批 × 32字节 = 3200字节/批
批间隔: 10ms (限速)
补传速率: ~320KB/s (理论)
实际: 受Modbus TCP报文开销影响,约50-100KB/s
11天数据补传耗时: 15872条 ÷ 100条/批 × 10ms = 1.6秒 (理论上)
实际: 约5-10秒实时采集: 200台 × 每10秒1次 = 20 QPS
补传: 3并发 × 50 QPS = 150 QPS
数据库写入: 实时20 + 补传150 = 170 QPS (峰值)
CPU占用: < 5% (单核)测试项 | 方法 | 预期结果 |
|---|---|---|
正常采集 | 运行24小时 | 数据完整率≥99.9% |
网络中断30分钟 | 断开交换机端口 | 记录仪本地存储不丢数据 |
网络恢复自动补传 | 恢复端口 | 30秒内开始补传,5分钟内完成 |
补传数据正确性 | 比对记录仪原始数据和数据库 | 完全一致 |
补传限速 | 监控采集延迟 | 实时采集不受影响 |
重复补传 | 模拟补传中断再恢复 | 不重复写入数据库 |
存储覆盖 | 模拟中断超过11天 | 正确报告不可恢复缺口 |
多设备并发补传 | 同时断开10台设备 | 并发控制正常,3台同时补传 |
服务器重启 | 重启采集服务 | 重启后从断点继续补传 |
□ 连续运行7天,数据完整率≥99.9%
□ 网络中断后恢复,100%自动补传成功
□ 补传过程中实时采集不受影响(延迟<2倍采样间隔)
□ 补传数据零丢失、零重复
□ 存储覆盖场景正确处理,告警明确
□ 所有操作有日志,可审计关键词:电厂配电室,环境系统,以太网温湿度记录仪,本地缓存,TCP断点续传,数据完整性,增量同步,Modbus TCP,时序数据库,容灾设计
标签:#电厂 #配电室 #温湿度记录仪 #本地缓存 #断点续传 #数据完整性 #ModbusTCP #容灾 #环境监测 #时序数据库
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。