字数:约3800字 | 阅读时间:12分钟
“数据丢了不可怕,可怕的是丢了还不知道”


问题背景:偏远风电场的网络困境

风电场的选址有一个铁律——风大的地方,通常也是人烟稀少的地方。对于华北平原上几十台风机组成的小型风场来说,4G/5G信号覆盖已经不是什么新鲜事,但”有信号”和”信号稳定”之间隔着一道巨大的鸿沟。

我们的监控系统架构是这样的:每台风机的PLC(可编程逻辑控制器)以50ms间隔采集振动、温度、转速等时序数据,通过边缘计算网关汇聚后,经MQTT协议上报到中控室的时序数据库。数据流设计上,每台风机每秒约产生20条记录,一个50台风机的风场,每秒就是1000条数据在管道中流动。

这套架构在城市里跑得好好的,但到了偏远风场就出问题了。

一个真实的故障场景: 2026年3月,某个风场连续一周出现数据缺口。值班人员查看监控画面,发现图表上每隔几小时就有一段空白——不是传感器故障,也不是设备断电,纯粹是网络抖动导致MQTT连接断开,数据在传输途中丢失了。更让人头疼的是,丢数据这件事本身也被”丢”了——监控面板上只显示空白,没有告警,值班人员根本不知道什么时候丢的、丢了多少。

这就是本文要解决的问题:在网络不稳定的工业场景下,如何保证时序数据的完整性?

时序数据写入的”窗口”机制

在讨论保数策略之前,先理解时序数据库的写入机制。以我们使用的IoTDB为例,它的写入流程大致是:

1
传感器 → 边缘网关 → MQTT Broker → 写入服务 → IoTDB Storage

IoTDB的写入有一个关键特性:批量写入窗口。单条写入虽然也能成功,但性能极差,TPS大概在几千到一万之间。如果开启批量写入,把数据攒到一定数量或时间窗口再flush,TPS可以提升到几十万。我们的配置是:

1
2
3
4
5
6
7
8
9
// IoTDB Session 批量写入配置
SessionPool sessionPool = new SessionPool(
"10.0.1.100", 6667, "root", "root",
10, // session数量
false, // 不启用 thrift 压缩
5000, // 批量写入阈值:5000条/次
1000, // flush间隔:1000ms
3 // 重试次数
);

这个配置在稳定网络下工作得很好,但在网络抖动场景下,问题就来了。当MQTT连接断开时,Broker侧会缓存消息,但如果断开时间超过Broker的缓存上限(默认通常是1GB内存),老消息会被淘汰。而写入服务这边,如果在flush窗口内网络恢复, buffered 的数据可以继续写入;如果 flush 失败,IoTDB会按重试次数重试3次——重试全部失败后,这批数据就真的丢了。

踩坑记录:MQTT QoS级别的”静默丢失”

在最初设计保数策略时,我们犯了一个经典错误——选择了MQTT QoS 0(至多一次)

选择QoS 0的理由看似充分:

  • QoS 0没有ACK机制,传输延迟较低
  • 风场数据实时性要求高,宁可丢旧数据也不要延迟新数据
  • 减少网络开销(偏远地区带宽有限)

但实际运行后发现一个致命问题:QoS 0的消息丢失是静默的

Broker不会告诉你消息丢了,写入服务不会告诉你消息没收到,时序数据库更不会提醒你少了几条记录。唯一的发现方式是人工比对——把传感器本地存储的原始数据和中控室的数据做diff,这在大规模风场里几乎不可行。

排查过程:

我们花了两天时间才定位到这个问题。起初怀疑是传感器故障,因为某些风机的数据确实出现了间断。但更换传感器后问题依旧,而且丢数据的风机位置不固定,和时间也没有明显关联。

后来在一个偶然的机会下,我们在边缘网关上抓包分析MQTT流量,发现QoS 0的PUBLISH消息在网络抖动时确实被丢弃了,而Broker侧的入站计数和出站计数存在差值——这说明消息在Broker之前就丢了,根本没进入Broker的缓存。

解决思路: 立即切换到QoS 1(至少一次),虽然增加了ACK开销,但确保每条消息至少被Broker接收一次。对于偶尔出现的重复消息(QoS 1在网络抖动时可能产生重复),我们在写入服务侧做了去重处理(基于消息ID)。

1
2
3
4
5
6
7
8
9
// MQTT客户端配置 - 修正后
MqttConnectOptions options = new MqttConnectOptions();
options.setAutomaticReconnect(true);
options.setConnectionTimeout(10);
options.setKeepAliveInterval(60);
options.setMqttVersion(MqttConnectOptions.MQTT_VERSION_3_1_1);

// ⚠️ 关键:使用QoS 1而不是QoS 0
int qos = 1;

QoS 2(精确一次)我们没用,因为在时序数据场景下,少量重复数据的代价远小于”丢失后需要去重”的复杂度。

三层保数方案:边缘缓冲 + 断点续传 + 去重校验

解决了QoS的问题后,我们开始设计系统性的保数方案。最终采用了三层架构:

第一层:边缘端本地缓冲

在每台边缘网关上部署本地缓冲队列,使用Disruptor作为高性能环形缓冲。当MQTT连接正常时,数据直接上报;当连接异常时,数据写入本地缓冲队列,等待连接恢复后批量续传。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
// 边缘网关 - 本地缓冲设计
public class EdgeBufferManager {
// 使用Disruptor作为高性能环形缓冲
private final RingBuffer<DataEvent> ringBuffer;
// 本地持久化目录
private final String persistPath = "/opt/edge/buffer/";
// 缓冲队列容量:50万条(约25分钟的数据量)
private static final int BUFFER_SIZE = 500_000;
// 持久化阈值:队列超过10万条时写入磁盘
private static final int FLUSH_THRESHOLD = 100_000;

public void onDataReceived(SensorData data) {
if (mqttClient.isConnected()) {
// 正常上报
publishWithQoS1(data);
} else {
// 写入本地缓冲
ringBuffer.publishEvent((event, sequence) -> {
event.setData(data);
event.setTimestamp(System.currentTimeMillis());
});
}
}

// 定时检查缓冲区,溢出时持久化到磁盘
@Scheduled(fixedRate = 5000)
public void checkAndPersist() {
long remaining = ringBuffer.remainingCapacity();
if (remaining < BUFFER_SIZE - FLUSH_THRESHOLD) {
List<SensorData> batch = drainBatch(FLUSH_THRESHOLD);
persistToDisk(batch);
}
}

// 连接恢复后批量续传
public void onConnectionRestored() {
List<SensorData> buffered = drainAll();
for (SensorData data : buffered) {
publishWithQoS1(data);
}
}
}

这里有一个设计考量:为什么用内存环形缓冲 + 磁盘持久化的混合方案,而不是纯磁盘队列?

原因有两点:

  1. 纯磁盘队列(如基于Kafka的本地队列)在边缘网关有限的存储资源下,随机写入性能不够。Disruptor的单线程吞吐量可以达到百万级别,足以应对每秒20条/风机的数据量。
  2. 大多数网络抖动在几秒到几分钟内恢复,内存缓冲足以覆盖绝大部分场景。磁盘持久化是兜底,防止长时间断网导致内存溢出。

磁盘持久化采用追加写入的WAL(Write-Ahead Log)格式,恢复时顺序读取,开销很小:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
// WAL持久化格式
// [timestamp:8bytes][deviceId:32bytes][dataLength:4bytes][data:N bytes][checksum:4bytes]

private void persistToDisk(List<SensorData> batch) {
try (FileChannel channel = FileChannel.open(
Paths.get(persistPath, "buffer.wal"),
StandardOpenOption.CREATE, StandardOpenOption.WRITE,
StandardOpenOption.APPEND)) {
ByteBuffer buffer = ByteBuffer.allocate(1024 * batch.size());
for (SensorData data : batch) {
buffer.putLong(data.getTimestamp());
buffer.put(data.getDeviceId().getBytes(StandardCharsets.UTF_8));
byte[] payload = data.toJson().getBytes(StandardCharsets.UTF_8);
buffer.putInt(payload.length);
buffer.put(payload);
int checksum = crc32(buffer.array());
buffer.putInt(checksum);
}
buffer.flip();
channel.write(buffer);
}
}

第二层:断点续传机制

连接恢复后的续传需要考虑几个问题:

  • 续传顺序: 按时间戳升序续传,保证时序数据库的写入顺序正确
  • 续传速率控制: 避免续传数据和新数据同时涌入导致写入服务过载
  • 续传超时处理: 如果续传过程中再次断开,未完成的批次重新回到缓冲队列
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
// 断点续传 - 速率控制
public class ResumeTransferManager {
private final MqttAsyncClient mqttClient;
private final ExecutorService transferExecutor;
// 续传速率:每批200条,间隔100ms
private static final int BATCH_SIZE = 200;
private static final int BATCH_INTERVAL_MS = 100;

public void resumeTransfer(List<SensorData> backlog) {
// 按时间戳排序
backlog.sort(Comparator.comparingLong(SensorData::getTimestamp));

transferExecutor.submit(() -> {
int offset = 0;
while (offset < backlog.size()) {
if (!mqttClient.isConnected()) {
// 传输中再次断开,剩余数据回写缓冲
writeBackToBuffer(backlog.subList(offset, backlog.size()));
return;
}

int end = Math.min(offset + BATCH_SIZE, backlog.size());
List<SensorData> batch = backlog.subList(offset, end);
publishBatch(batch);

offset = end;
try {
Thread.sleep(BATCH_INTERVAL_MS);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
writeBackToBuffer(backlog.subList(offset, backlog.size()));
}
}
});
}
}

速率控制的设计依据:每批200条、间隔100ms,续传速率约2000条/秒,大约是正常数据产生速率的2倍。这样既能较快补齐历史数据,又不会对写入服务造成太大压力。如果风场规模更大(比如200台风机),可以适当提高批次大小或缩短间隔。

第三层:去重校验

QoS 1的ACK重传机制和网络抖动可能导致重复消息到达。我们在写入服务侧实现基于消息ID的去重:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
// 写入服务 - 去重设计
public class DedupWriteService {
// 布隆过滤器用于快速判断是否见过该消息
private final BloomFilter<String> recentMsgFilter;
// 最近消息ID缓存,用于精确去重
private final Cache<String, Boolean> msgIdCache;

// 布隆过滤器容量:100万条,误判率0.01%
// 过期时间:30分钟(超过30分钟的消息ID不再缓存)
@PostConstruct
public void init() {
recentMsgFilter = BloomFilter.create(
Funnels.stringFunnel(StandardCharsets.UTF_8),
1_000_000, 0.01);
msgIdCache = Caffeine.newBuilder()
.maximumSize(1_000_000)
.expireAfterWrite(30, TimeUnit.MINUTES)
.build();
}

public WriteResult write(SensorData data) {
String msgId = data.getDeviceId() + ":" + data.getTimestamp();

// 第一层:布隆过滤器快速判断
if (recentMsgFilter.mightContain(msgId)) {
// 可能重复,精确检查缓存
if (msgIdCache.getIfPresent(msgId) != null) {
return WriteResult.DUPLICATED;
}
}

// 执行写入
boolean success = sessionPool.insertRecord(
data.getDeviceId(),
data.getTimestamp(),
data.getMeasurements()
);

if (success) {
recentMsgFilter.put(msgId);
msgIdCache.put(msgId, true);
return WriteResult.SUCCESS;
}
return WriteResult.FAILED;
}
}

布隆过滤器 + Caffeine缓存的组合,在保证低内存占用的同时实现了高效去重。布隆过滤器的误判率控制在0.01%,即便有误判,也只是让一个不重复的消息被跳过(可以通过定期对账补偿),不会产生重复数据。

数据完整性校验

保数策略的最后一环是完整性校验。我们在三个层面进行校验:

1. 边缘端自校验

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
// 边缘网关 - 每小时对账
@Scheduled(cron = "0 0 * * * ?")
public void hourlyReconciliation() {
// 统计本地产生和已上报的数据量
long produced = sensorCollector.getTotalCount();
long published = mqttReporter.getPublishedCount();
long buffered = bufferManager.getBufferedCount();

if (produced != published + buffered) {
// 数据量不匹配,可能存在丢失
alertService.sendAlert(
"EDGE_DATA_MISMATCH",
String.format("产生=%d, 上报=%d, 缓冲=%d, 缺口=%d",
produced, published, buffered, produced - published - buffered)
);
}
}

2. 中控端校验

中控室侧,我们通过检查时间戳连续性来判断数据完整性:

1
2
3
4
5
-- IoTDB: 检查某台风机最近1小时的数据连续性
-- 正常情况下,50ms间隔意味着每小时72000条记录
SELECT COUNT(*) FROM wind_turbine_01.root.sg.d001
WHERE time >= 2026-06-30 00:00:00
AND time < 2026-06-30 01:00:00;

如果某台风机的记录数远低于预期值(比如低于60000条),系统会自动标记为”数据不完整”并触发告警。

3. 定期全量对账

每周进行一次全量对账:边缘网关导出一周的传感器原始数据摘要(按时间窗口的计数和校验和),与中控室的数据做对比。这个过程在凌晨低峰期执行,对业务几乎无影响。

从单机到集群的延伸

单台风机层面的保数方案相对简单,但当系统扩展到几十个风场、几千台风机时,挑战就升级了。

写入瓶颈: 几千台风机同时上报,中控室的单点IoTDB可能扛不住。我们的解决方案是引入集群拓扑——每个风场部署一个IoTDB DataNode,通过IoTDB的集群模式(Consensus Protocol)同步到中控室的元数据节点。

缓冲溢出: 如果某个风场断网时间过长(比如超过24小时),边缘网关的磁盘缓冲可能写满。我们在设计中加入了”降级采样”策略:当缓冲队列超过容量的80%时,自动降低采样频率(从50ms降到500ms),优先保留近期完整数据,丢弃低频细节数据,同时持续告警提醒运维人员介入。

1
2
3
4
5
6
7
8
9
10
11
12
13
// 降级采样策略
public void checkBufferHealth() {
double usage = bufferManager.getUsagePercent();
if (usage > 80) {
// 降级:降低采样频率
sensorCollector.setInterval(500); // 从50ms降到500ms
alertService.sendAlert("BUFFER_OVERFLOW_WARNING",
"缓冲使用率" + usage + "%,已降级采样");
} else if (usage < 50) {
// 恢复:恢复正常采样
sensorCollector.setInterval(50);
}
}

这套方案上线后,经过两个月的运行观察,数据完整率从之前的约96%提升到了99.7%。剩下0.3%的缺失主要来自极端断网场景(比如光缆被挖断,断网超过48小时),这种情况下降级采样的数据虽然不完整,但核心趋势数据仍然保留,足以支撑事后分析。


总结来说,工业时序数据的保数策略核心就三件事:先保证数据不被静默丢弃(QoS 1 + ACK),然后设计合理的缓冲机制应对断网(内存 + 磁盘混合缓冲),最后在接收端做去重校验确保数据质量。 在风电场这种网络条件有限的场景下,这套方案经过了实际验证,确实把数据完整率提升到了可接受的水平。


本文由AI辅助生成框架,技术细节来自真实项目经验。