上一篇讲了多线程并行处理 IoT 数据的方案。但"处理顺序对了"不等于"数据时序对了"。设备时钟漂移、网络延迟、乱序到达,都会让数据库里的数据看起来"时间倒流"。本文记录如何保证 IoT 数据的时序正确性。
一、两种时间,两个问题
IoT 数据涉及两种时间:
| 时间 | 含义 | 谁决定 |
|---|---|---|
| 设备时间(TM) | 数据产生的真实时刻 | 设备本地时钟 |
| 服务器时间(receiveTime) | 数据到达服务器的时刻 | 服务器时钟 |
它们不一定一致,甚至经常不一致:
设备时间线: T1 → T2 → T3 → T4(正确顺序)
网络延迟后: T1 到达 → T3 到达 → T2 到达 → T4 到达(乱序!)
服务器记录: receiveTime: 10:00:01 → 10:00:02 → 10:00:03 → 10:00:04
如果按 receiveTime 排序,T2 和 T3 的顺序是反的。
二、我的 tcp-server 项目中的实际处理
行车监控项目中,设备发来的数据格式:
ID:{clientId},{status},TM:{datetime},X:{x},Y:{y},Z:{z},T:{weight},V:{extra};
TM 格式是 yyyy-M-d-H-m-s(如 2026-6-25-16-14-13),这是设备时间。
当前做法
项目里 ClientData 实体有 receiveTime(@PrePersist 自动填充服务器时间),但没有单独存设备时间。查询时按 receiveTime 排序。
问题
- 设备时钟漂移 → TM 和 receiveTime 差几秒甚至几分钟
- 项目已经检测了时间漂移(
timeDiffSeconds),超过阈值标记为warning - 但数据仍然按 receiveTime 存储和查询
改进方案
@Entity
public class ClientData {
@Id
private String clientId;
private LocalDateTime receiveTime; // 服务器接收时间(已有)
private LocalDateTime deviceTime; // 设备时间(新增,从 TM 解析)
@PrePersist
void prePersist() {
this.receiveTime = LocalDateTime.now();
}
}
解析 TM 字段:
// TcpChannelHandler 中
private LocalDateTime parseDeviceTime(String tm) {
// TM 格式: yyyy-M-d-H-m-s
DateTimeFormatter fmt = DateTimeFormatter.ofPattern("yyyy-M-d-H-m-s");
return LocalDateTime.parse(tm, fmt);
}
查询时按设备时间排序:
SELECT * FROM client_data
WHERE client_id = ?
ORDER BY device_time ASC; -- 而不是 receive_time
三、乱序到达的处理策略
网络抖动导致数据乱序到达是常态,有几种处理策略:
策略一:容忍乱序,查询时排序(最简单)
// 不管到达顺序,全部入库
// 查询时 ORDER BY device_time
适用:数据展示、历史回放。大多数监控场景用这个就够了。
策略二:水位线(Watermark)等待
借鉴 Flink 的思路:不急着处理,等一小段时间让乱序数据到齐。
// 每台设备维护一个缓冲区
Map<String, PriorityQueue<DeviceData>> buffers = new ConcurrentHashMap<>();
void onData(DeviceData data) {
buffers.computeIfAbsent(data.getDeviceId(), k -> new PriorityQueue<>(
Comparator.comparing(DeviceData::getDeviceTime)
)).offer(data);
}
// 定时刷新:只处理"足够老"的数据
@Scheduled(fixedRate = 2000)
void flush() {
long watermark = System.currentTimeMillis() - 5000; // 5秒水位线
for (var entry : buffers.entrySet()) {
PriorityQueue<DeviceData> pq = entry.getValue();
while (!pq.isEmpty() && pq.peek().getReceiveTimestamp() < watermark) {
DeviceData data = pq.poll();
processInOrder(data); // 按设备时间顺序处理
}
}
}
代价:增加 5 秒延迟。 适用:需要严格按序处理的场景(如状态机判断"先吊起后放下")。
策略三:序列号(最可靠)
让设备在每条数据里带一个递增序列号:
ID:crane-001,1,SEQ:1042,TM:2026-6-25-16-14-13,...
服务端检查序列号连续性:
Map<String, Long> lastSeq = new ConcurrentHashMap<>();
void onData(DeviceData data) {
long expected = lastSeq.getOrDefault(data.getDeviceId(), 0L) + 1;
if (data.getSeq() == expected) {
process(data);
lastSeq.put(data.getDeviceId(), data.getSeq());
} else if (data.getSeq() > expected) {
// 有数据还没到,暂存等待
pendingBuffer.offer(data);
}
// data.getSeq() <= lastSeq → 重复数据,丢弃
}
代价:需要设备端配合改造协议。 适用:金融交易、工业控制等不能容忍任何乱序的场景。
四、批量写入时的顺序问题
我的 tcp-server 用 LinkedBlockingQueue + 单线程批量 saveAll。这里有个微妙的问题:
// 16 个 Netty 线程同时往队列里放数据
// 队列是 FIFO,但多个线程 offer 的顺序是不确定的!
线程1: offer(设备A-T3) ← 先拿到 CPU
线程2: offer(设备A-T2) ← 后拿到 CPU
队列顺序: [A-T3, A-T2] ← 反了!
解决方案
方案 A:入库前按设备时间排序
@Scheduled(fixedRate = 500)
void flush() {
List<DeviceData> batch = new ArrayList<>(50);
buffer.drainTo(batch, 50);
// 按设备时间排序后再写入
batch.sort(Comparator.comparing(DeviceData::getDeviceTime));
repository.saveAll(batch);
}
方案 B:用上一篇的哈希分桶,同设备数据进同一个队列,单线程消费天然有序。
方案 C:数据库层面用 device_time 做排序键,不依赖写入顺序。
推荐 A + C 组合:写入前排序 + 查询时按设备时间排。双保险。
五、时间漂移检测与处理
设备时钟不准是 IoT 的常态。我的项目已经在做:
// DataController 中计算 timeDiffSeconds
long diff = serverTime - deviceTime; // 秒
if (diff > 300 || diff < -180) { // 超前5分钟或落后3分钟
status = "warning"; // 标记异常
errorMessage += "设备时间异常\n";
}
更完善的处理
// 1. 记录漂移量,但不丢弃数据
data.setTimeDiffSeconds(diff);
data.setDeviceTime(parsedDeviceTime); // 仍然存设备时间
// 2. 如果漂移过大,用服务器时间修正
if (Math.abs(diff) > 3600) { // 差 1 小时以上,设备时钟肯定坏了
data.setDeviceTime(data.getReceiveTime()); // 降级用服务器时间
data.setTimeCorrected(true);
}
// 3. 定期校准(如果设备支持 NTP 或下发校时命令)
六、Kafka 中的时序保证
如果用 Kafka 做数据管道(我的项目里 MQTT → TCP 就是类似思路):
设备 → Kafka (partition by deviceId) → Consumer → 数据库
Kafka 保证同一 partition 内消息有序。但要注意:
| 配置 | 影响 |
|---|---|
acks=1 | Leader 写入就确认,可能丢数据 |
acks=all | 所有副本写入才确认,不丢但慢 |
retries > 0 + max.in.flight=1 | 保证重试不亂序 |
enable.idempotence=true | 幂等生产者,自动保证顺序 |
关键配置:
# 生产者:保证顺序
enable.idempotence=true
max.in.flight.requests.per.connection=1
acks=all
retries=3
七、完整方案总结
设备数据到达
│
├─ 解析设备时间(TM)
├─ 检测时间漂移(对比服务器时间)
├─ 哈希分桶 → 同设备同线程(保证处理顺序)
├─ 批量写入前按 device_time 排序
├─ 数据库存储 device_time + receive_time 双字段
└─ 查询时 ORDER BY device_time
| 层面 | 保证手段 |
|---|---|
| 接收层 | Netty 多 worker 并行接收 |
| 处理层 | 哈希分桶,同设备串行 |
| 存储层 | 批量写入前排序 |
| 查询层 | ORDER BY device_time |
| 异常处理 | 时间漂移检测 + 降级 |
八、一句话总结
处理顺序靠线程模型(哈希分桶),数据时序靠设备时间(TM 字段)。 两者缺一不可:线程模型保证"同一设备的数据按到达顺序处理",设备时间保证"即使乱序到达,查询结果仍然是正确的时间线"。
附录:Flink 是什么?
上文提到"借鉴 Flink 的思路",这里补充介绍一下 Flink。
Flink 是 Apache 开源的分布式流处理引擎,专门处理"永不停止的数据流"。
传统数据库:数据存好了,你去查(批处理)
Flink:数据像水龙头一样一直流过来,你边流边处理(流处理)
核心概念
数据流(DataStream):
DataStream<DeviceData> stream = env
.addSource(new MqttSource()) // 数据源
.keyBy(data -> data.deviceId) // 按设备分组(同设备串行!)
.window(TumblingEventTimeWindows.of(Time.seconds(10)))
.aggregate(new AvgAggregator()); // 聚合计算
keyBy(deviceId) 和上一篇讲的哈希分桶是同一个思想。
事件时间 vs 处理时间:
| 时间语义 | 含义 |
|---|---|
| Event Time | 数据产生的真实时间(TM 字段) |
| Processing Time | 数据到达服务器的时间 |
Flink 推荐用 Event Time(设备时间),正是本文说的"查询按 device_time 排序"。
水位线(Watermark):处理乱序的核心机制,本文"策略二"的来源:
数据到达顺序:T1, T3, T2, T5, T4(乱序!)
水位线 = 最大事件时间 - 允许延迟(5秒)
水位线推进到 T5 时 → T1~T5 都"到齐了",安全输出结果
代价:延迟 5 秒;收益:顺序正确
WatermarkStrategy
.<DeviceData>forBoundedOutOfOrderness(Duration.ofSeconds(5))
.withTimestampAssigner((data, ts) -> data.getDeviceTimeMs())
窗口(Window):把无限数据流切成有限的块:
滚动窗口:[0-10s] [10-20s] [20-30s] → 每10秒统计一次
滑动窗口:[0-10s] [5-15s] [10-20s] → 每5秒更新最近10秒
会话窗口:设备活跃期一个窗口,沉默30秒关闭
和 tcp-server 项目的对应关系
| Flink 概念 | tcp-server 对应 |
|---|---|
| Source | Netty TCP + MQTT 接入 |
| keyBy | 按 clientId 处理 |
| Event Time | TM 字段 |
| Watermark | timeDiffSeconds 漂移检测(手动版) |
| Window + Aggregate | 每 500ms 批量 saveAll |
| Sink | MySQL 入库 |
tcp-server 本质上是一个手写单节点版流处理系统。
什么时候该用 Flink
| 场景 | 需要吗 |
|---|---|
| 行车监控,千级设备 | 不需要 |
| 实时告警:连续 3 次超标就报警 | 可选(Flink CEP 很方便) |
| 十万级设备实时聚合 | 需要 |
| 实时大屏 | 需要 |
本文借鉴的就是 Flink 的事件时间和水位线两个思想。这两个思想不需要引入 Flink,在任何 Java 项目里都能实现。
水位线详解:用一个具体例子走一遍
场景
行车设备每秒发一条数据,数据里带设备时间 TM。但网络有延迟,数据不一定按顺序到。
没有水位线会怎样
假设要统计"每 10 秒的平均重量":
窗口 [10:00:00 ~ 10:00:10) 的数据:
实际到达顺序:
10:00:11 收到 TM=10:00:03 重量=5吨
10:00:12 收到 TM=10:00:01 重量=3吨
10:00:12 收到 TM=10:00:05 重量=7吨
10:00:13 收到 TM=10:00:02 重量=4吨 ← 迟到了!
10:00:13 收到 TM=10:00:08 重量=6吨
如果你不等:10:00:10 一到就关闭窗口算平均 → 只收到 3 条 → 结果错了。
如果你永远等:结果永远出不来。
水位线怎么解决
水位线 = “我估计 TM <= X 的数据应该都到齐了”
// 设置:允许 5 秒延迟
WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(5))
意思:水位线 = 目前见过的最大 TM - 5 秒
走一遍时间线:
时刻 收到数据 最大TM 水位线(最大TM-5s) 能关哪个窗口?
─────────────────────────────────────────────────────────────────────────
10:00:11 TM=10:00:03 10:00:03 09:59:58 不能关任何窗口
10:00:12 TM=10:00:01 10:00:03 09:59:58 不能(没推进)
10:00:12 TM=10:00:05 10:00:05 10:00:00 不能(刚到边界)
10:00:13 TM=10:00:02 10:00:05 10:00:00 不能(没推进)
10:00:13 TM=10:00:08 10:00:08 10:00:03 不能
10:00:15 TM=10:00:12 10:00:12 10:00:07 不能
10:00:16 TM=10:00:14 10:00:14 10:00:09 不能
10:00:17 TM=10:00:15 10:00:15 10:00:10 ← 到10了! 关闭[00~10)窗口!
水位线推进到 10:00:10 的那一刻,Flink 说:“我估计 TM < 10:00:10 的数据都到了”,于是关闭 [10:00:00, 10:00:10) 窗口,计算平均值,输出结果。
关键点
窗口 [10:00:00 ~ 10:00:10) 什么时候关?
→ 不是 10:00:10 关
→ 是水位线到达 10:00:10 时关
→ 水位线 = 最大TM - 5秒
→ 所以要等到出现 TM=10:00:15 的数据
→ 实际关闭时间 ≈ 10:00:17(晚了 7 秒)
5 秒的延迟容忍,换来的是:TM=10:00:02 那条迟到的数据赶上了窗口。
如果数据迟到超过 5 秒呢
10:00:25 收到 TM=10:00:04 ← 迟到了 21 秒!超过 5 秒容忍
此时窗口 [00~10) 早已关闭。这条数据被丢弃(默认行为),或者进入"侧输出流"单独处理:
OutputTag<DeviceData> lateTag = new OutputTag<>("late-data") {};
stream.keyBy(...)
.window(...)
.allowedLateness(Time.seconds(30)) // 额外再等 30 秒
.sideOutputLateData(lateTag) // 超过 30 秒的走这里
.aggregate(...);
// 单独处理迟到数据
DataStream<DeviceData> lateStream = stream.getSideOutput(lateTag);
lateStream.addSink(new LateDataHandler());
和 tcp-server 的类比
tcp-server 有一个"手动版水位线":
if (timeDiffSeconds > 300 || timeDiffSeconds < -180) {
status = "warning"; // 设备时间偏差太大 → 标记异常
}
区别:
| Flink 水位线 | tcp-server | |
|---|---|---|
| 怎么处理迟到 | 等 5 秒,窗口内都能收 | 不等,到了就存 |
| 怎么判断异常 | 超过水位线 → 侧输出 | 超过 300 秒 → warning |
| 排序方式 | 窗口内按事件时间聚合 | 查询时 ORDER BY device_time |
一句话
水位线就是一个"截止时间估算器":它根据已经看到的数据,猜测"某个时间点之前的数据应该都到了",然后触发计算。允许多大延迟由你设定。延迟设大了,结果更准但更慢;设小了,更快但可能丢迟到数据。