
关键词:机房环境数据采集、温湿度传感器、TCP/UDP双模式切换、记录仪存储、断网续传、本地缓存、边缘计算、Modbus TCP、私有UDP上报、数据完整性 标签:#物联网 #Modbus #TCP/IP #UDP #POE供电 #InfluxDB #以太网温湿度传感器 #网口温湿度变送器 #机房监控 #边缘计算
前一篇讲了 UDP 采集与丢包排查。实际机房项目中,只跑一种协议往往不够——
纯 TCP(Modbus TCP)的问题:
纯 UDP 的问题:
双模式的核心思路:TCP 做"可靠基线",UDP 做"高频快采",两者互补。同时,传感器本地记录仪提供"断网兜底"——网络全断时数据不丢,恢复后补传。
模式 | 协议 | 方向 | 用途 | 数据特征 |
|---|---|---|---|---|
Mode A:TCP 轮询 | Modbus TCP | 采集端→设备,请求-响应 | 可靠采集、配置下发、校准 | 周期 10~30s,每次读 4~8 寄存器 |
Mode B:UDP 主动上报 | 私有 UDP | 设备→采集端,单向推送 | 高频采样、实时告警触发 | 周期 2~5s,带序号和时标 |
Mode C:混合模式 | TCP + UDP 并行 | 双向 | 正常运行态 | TCP 做校验基线,UDP 做高频流 |
Mode D:离线记录 | 本地 Flash/SD | 设备内部 | 断网时本地存储 | 按采样周期存,网络恢复后补传 |
┌─────────────────────────────────────────────────────────┐
│ 模式决策状态机 │
├─────────────────────────────────────────────────────────┤
│ │
│ 上电启动 │
│ │ │
│ ▼ │
│ ┌──────────┐ TCP 连接成功 ┌──────────────────┐ │
│ │ TCP探测 │──────────────────▶│ Mode C: 混合运行 │ │
│ └──────────┘ └──────────────────┘ │
│ │ │ │ │
│ │ TCP 失败 │ │ UDP 连续丢包 │
│ ▼ │ ▼ │
│ ┌──────────┐ UDP 可达 ┌──────────────────┐ │
│ │ UDP探测 │──────────────────▶│ Mode B: UDP单模 │ │
│ └──────────┘ └──────────────────┘ │
│ │ │ │
│ │ 均失败 │ 持续失败超时 │
│ ▼ ▼ │
│ ┌──────────────────────────────────────────────┐ │
│ │ Mode D: 离线记录模式(本地存储,等待网络恢复)│ │
│ └──────────────────────────────────────────────┘ │
│ │ │
│ │ 网络恢复检测到连通 │
│ ▼ │
│ 重新进入 TCP探测 → 状态机循环 │
└─────────────────────────────────────────────────────────┘切换判定参数(建议值,可按环境调整):
参数 | 含义 | 默认值 |
|---|---|---|
tcp_probe_timeout | TCP 连接尝试超时 | 5s |
tcp_fail_threshold | 连续 TCP 失败次数触发降级 | 3 次 |
udp_loss_threshold | UDP 连续丢包数触发告警/降级 | 10 个周期 |
udp_loss_window | 丢包统计滑动窗口 | 30s |
offline_threshold | 双模均失败持续时间,进入离线模式 | 60s |
recovery_probe_interval | 离线模式下探测网络恢复间隔 | 30s |
传感器 MCU 需要同时维护:
/* sensor_dual_mode.c - 传感器端双模式核心逻辑(伪代码/参考实现) */
#include "lwip/tcp.h"
#include "lwip/udp.h"
#include "flash_storage.h"
/* 模式状态 */
typedef enum {
MODE_TCP_ONLY = 0,
MODE_UDP_ONLY = 1,
MODE_HYBRID = 2,
MODE_OFFLINE = 3
} op_mode_t;
static op_mode_t current_mode = MODE_HYBRID;
static uint16_t udp_seq = 0;
static uint32_t tcp_fail_count = 0;
static uint32_t udp_loss_count = 0;
static uint32_t last_udp_ack_ts = 0;
/* Modbus TCP 服务端回调 */
static err_t modbus_tcp_accept(void *arg, struct tcp_pcb *newpcb, errr err) {
/* 接受连接,注册 recv 回调 */
tcp_recv(newpcb, modbus_tcp_recv);
tcp_err(newpcb, modbus_tcp_err);
tcp_fail_count = 0; /* 连接成功,重置失败计数 */
return ERR_OK;
}
static void modbus_tcp_err(void *arg, err_t err) {
/* TCP 连接异常断开 */
tcp_fail_count++;
if (tcp_fail_count >= TCP_FAIL_THRESHOLD) {
/* 降级:TCP 不可靠,切到 UDP 单模或混合 */
evaluate_mode_downgrade();
}
}
/* UDP 上报定时器回调(每 5s 触发) */
static void udp_report_timer(void *arg) {
if (current_mode == MODE_OFFLINE) {
/* 离线模式:只写本地存储 */
record_to_local_storage();
probe_network_recovery();
return;
}
/* 构造 UDP 上报帧 */
uint8_t frame[32];
build_udp_frame(frame, sizeof(frame), udp_seq++);
/* 发送 */
udp_sendto(udp_pcb, frame, remote_ip, REMOTE_UDP_PORT);
/* 如果混合模式,TCP 侧仍在响应请求(由 lwIP 栈自动处理) */
}
/* 网络恢复探测 */
static void probe_network_recovery(void) {
/* 尝试建立 TCP 连接到采集端 */
struct tcp_pcb *tpcb = tcp_new();
err_t err = tcp_connect(tpcb, &collector_ip, 502, on_recovery_connected);
if (err != ERR_OK) {
tcp_abort(tpcb);
}
}
static err_t on_recovery_connected(void *arg, struct tcp_pcb *tpcb, err_t err) {
if (err == ERR_OK) {
/* 网络恢复,切回混合模式 */
current_mode = MODE_HYBRID;
tcp_fail_count = 0;
/* 触发补传 */
trigger_backfill();
}
tcp_close(tpcb);
return ERR_OK;
}传感器本地存储的核心约束:
存储布局(以 SPI Flash 4MB 为例):
├── 0x000000 ~ 0x000FFF 元数据区(256KB)
│ ├── 设备ID、固件版本
│ ├── 环形缓冲头指针/尾指针
│ ├── 采样配置(周期、量程)
│ └── 补传状态位图
│
├── 0x001000 ~ 0x3FFFFF 数据区(约 4MB - 256KB)
│ ├── 每条记录 16 字节:
│ │ ├── timestamp (4B, Unix epoch)
│ │ ├── temperature (2B, int16, ×0.1℃)
│ │ ├── humidity (2B, int16, ×0.1%RH)
│ │ ├── quality (1B, bit0=传感器OK, bit1=校准有效)
│ │ ├── seq (2B, 序号)
│ │ └── reserved (5B)
│ └── 可存约 262000 条
│ @ 5s 采样 = 连续存储 15 天
│ @ 30s 采样 = 连续存储 90 天
│
└── 写策略:
├── 顺序写,到末尾回绕到开头
├── 尾指针追上头指针时覆盖最旧数据
├── 每次写入后 flush + 更新尾指针
└── 补传后更新"已同步"位图/* flash_storage.c - 本地存储核心 */
#define RECORD_SIZE 16
#define DATA_START_ADDR 0x001000
#define DATA_END_ADDR 0x3FFFFF
#define MAX_RECORDS ((DATA_END_ADDR - DATA_START_ADDR) / RECORD_SIZE)
typedef struct {
uint32_t timestamp;
int16_t temperature; /* ×0.1℃ */
int16_t humidity; /* ×0.1%RH */
uint8_t quality;
uint16_t seq;
uint8_t reserved[5];
} env_record_t;
typedef struct {
uint32_t head; /* 最旧已同步数据的位置 */
uint32_t tail; /* 下一次写入位置 */
uint32_t count; /* 当前存储的记录数 */
uint16_t last_seq;
} storage_meta_t;
static storage_meta_t meta;
void record_sample(float temp, float hum) {
env_record_t rec;
rec.timestamp = get_unix_time();
rec.temperature = (int16_t)(temp * 10);
rec.humidity = (int16_t)(hum * 10);
rec.quality = 0x03; /* 传感器OK + 校准有效 */
rec.seq = meta.last_seq++;
/* 写入 Flash */
flash_write(DATA_START_ADDR + meta.tail * RECORD_SIZE, &rec, RECORD_SIZE);
/* 更新尾指针 */
meta.tail = (meta.tail + 1) % MAX_RECORDS;
meta.count = (meta.count < MAX_RECORDS) ? meta.count + 1 : MAX_RECORDS;
/* 环形覆盖:tail 追上 head 时 head 前进 */
if (meta.tail == meta.head && meta.count == MAX_RECORDS) {
meta.head = (meta.head + 1) % MAX_RECORDS;
}
/* 持久化元数据 */
flash_write_meta(&meta);
}┌─────────────────────────────────────────────────────────────┐
│ 采集端进程架构 │
├─────────────────────────────────────────────────────────────┤
│ │
│ ┌─────────────┐ ┌─────────────┐ ┌───────────────┐ │
│ │ TCP Poller │ │ UDP Receiver│ │ Mode Manager │ │
│ │ (Modbus) │ │ │ │ (状态机) │ │
│ └──────┬──────┘ └──────┬──────┘ └───────┬───────┘ │
│ │ │ │ │
│ ▼ ▼ ▼ │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ Data Fusion Layer(数据融合层) │ │
│ │ - 时间戳对齐(最近邻/插值) │ │
│ │ - 质量位合并(TCP 校验 vs UDP 序号连续性) │ │
│ │ - 去重(UDP 重传/重复检测) │ │
│ │ - 缺失判定(双通道都无数据 → 标记 stale) │ │
│ └───────────────────────────┬─────────────────────────┘ │
│ ▼ │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ Write Pipeline(写入管道) │ │
│ │ - InfluxDB 批量写入 │ │
│ │ - 本地 WAL(Write-Ahead Log)用于采集端自身断网 │ │
│ │ - 指标暴露(Prometheus /metrics) │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ Backfill Client(补传客户端) │ │
│ │ - 设备重连后,读取设备侧未同步数据 │ │
│ │ - 通过 TCP Modbus 扩展寄存器或专用协议拉取 │ │
│ │ - 按 seq 排序后批量写入,避免时序错乱 │ │
│ └─────────────────────────────────────────────────────┘ │
└─────────────────────────────────────────────────────────────┘"""
data_fusion.py - TCP/UDP 双通道数据融合
"""
import asyncio
import time
from collections import defaultdict, deque
from dataclasses import dataclass, field
from typing import Optional
@dataclass
class SensorReading:
dev_id: str
temperature: float
humidity: float
seq: int
timestamp: float
source: str # "tcp" or "udp"
quality: int = 1
@dataclass
class DeviceState:
last_tcp_reading: Optional[SensorReading] = None
last_udp_reading: Optional[SensorReading] = None
tcp_seq: int = 0
udp_seq: int = 0
udp_loss_streak: int = 0
last_fusion_output: float = 0.0
fusion_interval: float = 5.0 # 每 5s 输出一次融合结果
class DataFusion:
def __init__(self):
self.devices: dict[str, DeviceState] = {}
self.output_queue = asyncio.Queue(maxsize=10000)
def on_tcp_reading(self, reading: SensorReading):
"""TCP 轮询到数据"""
dev = self.devices.setdefault(reading.dev_id, DeviceState())
dev.last_tcp_reading = reading
# TCP 成功意味着通道健康
dev.udp_loss_streak = 0 # 重置 UDP 丢包计数(因为 TCP 兜底)
def on_udp_reading(self, reading: SensorReading):
"""UDP 收到上报"""
dev = self.devices.setdefault(reading.dev_id, DeviceState())
# 序号连续性检查
expected = dev.udp_seq + 1
if reading.seq != expected and dev.udp_seq > 0:
gap = reading.seq - dev.udp_seq
if gap > 1:
dev.udp_loss_streak += gap - 1
else:
dev.udp_loss_streak = max(0, dev.udp_loss_streak - 1)
dev.last_udp_reading = reading
dev.udp_seq = reading.seq
async def fusion_loop(self):
"""融合输出循环:每 fusion_interval 输出一条"""
while True:
await asyncio.sleep(self.fusion_interval)
now = time.time()
for dev_id, dev in self.devices.items():
# 选择策略:
# 1. UDP 有新鲜数据(< 2×interval)→ 优先用 UDP(高频)
# 2. UDP 超时或无数据 → 用 TCP 数据
# 3. 两者都有 → 交叉校验,偏差大则标 quality=0
udp_fresh = (
dev.last_udp_reading and
(now - dev.last_udp_reading.timestamp) < self.fusion_interval * 2
)
tcp_fresh = (
dev.last_tcp_reading and
(now - dev.last_tcp_reading.timestamp) < 30.0 # TCP 周期长
)
if udp_fresh:
result = dev.last_udp_reading
if tcp_fresh:
# 交叉校验
t_diff = abs(dev.last_udp_reading.temperature -
dev.last_tcp_reading.temperature)
h_diff = abs(dev.last_udp_reading.humidity -
dev.last_tcp_reading.humidity)
if t_diff > 1.0 or h_diff > 5.0:
result.quality = 0 # 通道不一致,标记需关注
elif tcp_fresh:
result = dev.last_tcp_reading
else:
# 双通道都无数据
result = SensorReading(
dev_id=dev_id, temperature=0, humidity=0,
seq=0, timestamp=now, source="none", quality=0
)
await self.output_queue.put(result)
dev.last_fusion_output = now"""
mode_manager_v2.py - 双模式切换管理
"""
import asyncio
import time
from enum import Enum
class CollectorMode(Enum):
HYBRID = "hybrid" # TCP + UDP 并行
TCP_ONLY = "tcp_only" # UDP 丢包严重,降级到 TCP
UDP_ONLY = "udp_only" # TCP 不可达,UDP 单模
DEGRADED = "degraded" # 双通道都异常,标记设备 stale
@dataclass
class DeviceHealth:
dev_id: str
tcp_success_count: int = 0
tcp_fail_count: int = 0
udp_packets_received: int = 0
udp_packets_expected: int = 0
last_tcp_ts: float = 0.0
last_udp_ts: float = 0.0
mode: CollectorMode = CollectorMode.HYBRID
offline_since: Optional[float] = None
class ModeManager:
def __init__(self, tcp_fail_threshold=3, udp_loss_ratio_threshold=0.3,
evaluation_window=60):
self.devices: dict[str, DeviceHealth] = {}
self.tcp_fail_threshold = tcp_fail_threshold
self.udp_loss_ratio_threshold = udp_loss_ratio_threshold
self.evaluation_window = evaluation_window
def report_tcp_result(self, dev_id: str, success: bool):
h = self.devices.setdefault(dev_id, DeviceHealth(dev_id=dev_id))
if success:
h.tcp_fail_count = 0
h.tcp_success_count += 1
h.last_tcp_ts = time.time()
else:
h.tcp_fail_count += 1
self._reevaluate(h)
def report_udp_packet(self, dev_id: str, received: bool):
h = self.devices.setdefault(dev_id, DeviceHealth(dev_id=dev_id))
h.udp_packets_expected += 1
if received:
h.udp_packets_received += 1
h.last_udp_ts = time.time()
# 滑动窗口:每 evaluation_window 秒重置计数
# 实际实现用环形缓冲或时间桶
def _reevaluate(self, h: DeviceHealth):
now = time.time()
tcp_ok = h.tcp_fail_count < self.tcp_fail_threshold
udp_loss_ratio = 1.0 - (h.udp_packets_received / max(1, h.udp_packets_expected))
udp_ok = udp_loss_ratio < self.udp_loss_ratio_threshold
if tcp_ok and udp_ok:
h.mode = CollectorMode.HYBRID
h.offline_since = None
elif tcp_ok and not udp_ok:
h.mode = CollectorMode.TCP_ONLY
elif not tcp_ok and udp_ok:
h.mode = CollectorMode.UDP_ONLY
else:
if h.offline_since is None:
h.offline_since = now
h.mode = CollectorMode.DEGRADED
def get_mode(self, dev_id: str) -> CollectorMode:
return self.devices.get(dev_id, DeviceHealth(dev_id)).mode设备重连(TCP 建立成功)
│
▼
采集端发送 "BACKFILL_REQ" 命令
│ 参数:start_seq, end_seq(或 start_time, end_time)
▼
设备端从 Flash 读取对应范围记录
│ 按页读取,每页 16 条记录(256B)
▼
设备端通过 TCP 逐页回传
│ 每页带 CRC16 校验
▼
采集端验证 CRC,写入 InfluxDB
│ 成功后发送 "ACK_PAGE" 带页号
▼
设备端标记该页已同步(更新位图)
│
▼
所有页传输完成 → 发送 "BACKFILL_DONE"保持寄存器区(0x2000 起始):
0x2000: backfill_status (RO) - bit0=in_progress, bit1=complete, bit2=error
0x2001: total_records (RO) - Flash 中总记录数
0x2002: unsynced_count (RO) - 未同步记录数
0x2003: oldest_unsynced_seq (RO) - 最旧未同步序号
0x2004: newest_stored_seq (RO) - 最新存储序号
0x2005: page_size (RO) - 每页记录数(固定 16)
0x2006: backfill_start_seq (RW) - 补传起始序号(采集端写入)
0x2007: backfill_end_seq (RW) - 补传结束序号(采集端写入)
0x2008: backfill_page_index (RW) - 请求第几页(采集端写入)
0x2009: backfill_control (RW) - bit0=start, bit1=abort, bit2=clear_synced
0x2100~0x211F: page_data[16] (RO) - 一页数据(16条×16B=256B,分 128 个寄存器)"""
backfill_client.py - 补传客户端
"""
import struct
import crc16
from pymodbus.client import ModbusTcpClient
class BackfillClient:
BASE = 0x2000
PAGE_DATA_BASE = 0x2100
PAGE_RECORDS = 16
def __init__(self, client: ModbusTcpClient):
self.c = client
def get_status(self):
rr = self.c.read_holding_registers(self.BASE, 5)
return {
"status": rr.registers[0],
"total": rr.registers[1],
"unsynced": rr.registers[2],
"oldest_seq": rr.registers[3],
"newest_seq": rr.registers[4],
}
def request_page(self, start_seq, end_seq, page_index):
# 设置补传参数
self.c.write_registers(self.BASE + 6, [start_seq, end_seq])
# 设置页索引并触发
self.c.write_registers(self.BASE + 8, [page_index, 1]) # control=start
# 读取页数据(128 个寄存器 = 256 字节)
rr = self.c.read_holding_registers(self.PAGE_DATA_BASE, 128)
raw = struct.pack('>128H', *rr.registers)
# 解析记录
records = []
for i in range(self.PAGE_RECORDS):
offset = i * 16
rec = struct.unpack_from('>IHhHBH5s', raw, offset)
records.append({
"seq": rec[0],
"timestamp": rec[1],
"temperature": rec[2] / 10.0,
"humidity": rec[3] / 10.0,
"quality": rec[4],
})
return records
def run_backfill(self, dev_id: str):
status = self.get_status()
if status["unsynced"] == 0:
return []
total_pages = (status["unsynced"] + self.PAGE_RECORDS - 1) // self.PAGE_RECORDS
all_records = []
for page in range(total_pages):
records = self.request_page(
status["oldest_seq"], status["newest_seq"], page
)
all_records.extend(records)
# 写入 InfluxDB(按原始时间戳,保持时序)
self.write_to_influx(dev_id, records)
# 清除已同步标记
self.c.write_registers(self.BASE + 8, [0, 4]) # control=clear_synced
return all_records场景 | 现象 | 根因 | 处理 |
|---|---|---|---|
TCP 连接数爆满 | 新传感器无法接入 | fd 耗尽或内存不足 | 采集端连接池复用,或改 UDP 为主 |
UDP 大量丢包 | 融合层频繁降级 TCP | 交换机缓冲溢出或主机 rmem 不足 | 调大 rmem、降低 UDP 发送频率、或增加采集端实例 |
补传数据时间戳错乱 | InfluxDB 曲线断层后突然回补一大段 | 补传写入未用原始时间戳 | 补传时强制使用设备端采集时标 |
Flash 写满后数据丢失 | 离线 7 天后恢复,只看到最近 3 天 | 环形缓冲覆盖 | 增大存储容量或缩短采样周期 |
模式频繁抖动 | 设备在线/离线状态反复变化 | 网络不稳定 + 切换阈值太敏感 | 加滞环和延时锁定 |
补传阻塞正常采集 | 补传期间实时数据延迟 | 单线程同步处理 | 补传走独立线程/协程,限速 |
双模式切换 + 本地记录仪,本质是在"实时性、可靠性、成本"三角中找平衡。TCP 保证不丢,UDP 保证快,本地存储保证断网不丢。工程落地的关键不在协议本身,而在模式切换的判定逻辑要稳(别抖)、融合策略要清晰(别糊)、补传要有序(别乱)。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。