物联网场景下,成百上千台设备每秒各发一条数据。单线程处理不过来,多线程又可能打乱顺序。本文记录如何用多线程并行处理 IoT 数据,同时保证同一设备的数据不被乱序。
一、问题本质
100 台行车设备,每秒各发 1 条数据 = 每秒 100 条。如果每条处理耗时 10ms(解析 + 校验 + 入库),单线程每秒只能处理 100 条,刚好卡在极限。设备一多就崩。
但直接开多线程有个致命问题:
设备A的第1条 → 线程1处理(耗时10ms)
设备A的第2条 → 线程2处理(耗时2ms,先完成)
→ 第2条比第1条先入库,顺序错了!
核心原则:同一设备的数据必须串行,不同设备的数据可以并行。
二、方案一:哈希分桶(最推荐)
思路:按设备 ID 哈希,把数据固定分配到 N 个队列,每个队列一个专属线程消费。
设备数据 → hash(deviceId) % N → 第K个队列 → 第K个线程
同一设备永远哈希到同一个队列 → 同一个线程 → FIFO 顺序天然保证。
完整实现
public class DeviceDispatcher {
private final int threadCount;
private final LinkedBlockingQueue<DeviceData>[] queues;
private final Thread[] consumers;
@SuppressWarnings("unchecked")
public DeviceDispatcher(int threadCount) {
this.threadCount = threadCount;
this.queues = new LinkedBlockingQueue[threadCount];
this.consumers = new Thread[threadCount];
for (int i = 0; i < threadCount; i++) {
queues[i] = new LinkedBlockingQueue<>(10000);
int idx = i;
consumers[i] = new Thread(() -> {
while (!Thread.currentThread().isInterrupted()) {
try {
DeviceData data = queues[idx].take(); // 阻塞等待
processDeviceData(data);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
}
}
}, "device-consumer-" + i);
consumers[i].start();
}
}
/** 生产端:Netty handler 收到数据后调用 */
public void dispatch(DeviceData data) {
int slot = Math.abs(data.getDeviceId().hashCode() % threadCount);
boolean ok = queues[slot].offer(data);
if (!ok) {
log.warn("队列已满,丢弃设备 {} 的数据", data.getDeviceId());
}
}
private void processDeviceData(DeviceData data) {
// 解析 → 校验 → 入库,同一设备永远在这个线程里串行执行
}
}
线程数怎么选
| 场景 | 建议线程数 |
|---|---|
| CPU 密集(复杂计算) | CPU 核数 |
| IO 密集(数据库写入) | CPU 核数 × 2 |
| 设备量 < 1000 | 4~8 |
| 设备量 > 10000 | 16~32 |
为什么不用线程池 + synchronized?
这是一种直觉上"看起来能行"的做法:用线程池并行处理,给每个设备加一把锁保证顺序。
// 线程池:10 个线程并行处理
ExecutorService pool = Executors.newFixedThreadPool(10);
// 每个设备一把锁
Map<String, Object> deviceLocks = new ConcurrentHashMap<>();
void onData(DeviceData data) {
pool.submit(() -> {
// 拿到这台设备的锁
Object lock = deviceLocks.computeIfAbsent(data.getDeviceId(), k -> new Object());
synchronized (lock) {
// 同一设备:锁住了,必须排队进来 → 顺序保证
// 不同设备:不同的锁对象,互不影响 → 并行
processData(data);
}
});
}
能用,但问题很多:
- 锁对象越来越多:设备几万台,就要维护几万个锁对象,内存浪费
- 设备下线后锁不清理:内存泄漏
- 锁竞争:如果
processData里有共享资源(比如数据库连接池),不同设备的线程还是在争抢 - 调试困难:死锁、锁等待,出问题很难排查
- 线程池 + 锁的组合本身就容易写出 bug
对比哈希分桶:
| 线程池 + synchronized | 哈希分桶 + N 队列 | |
|---|---|---|
| 顺序保证 | 靠锁 | 靠队列 FIFO(天然有序) |
| 锁 | 每设备一把,几万把 | 无锁 |
| 复杂度 | 锁管理、清理、死锁风险 | 队列 + 线程,简单 |
| 性能 | 锁竞争开销 | 无竞争,更快 |
| 内存 | 锁对象随设备增长 | 固定 N 个队列 |
一句话:能用队列解决的问题,不要用锁。 队列天然 FIFO,不需要额外同步机制。哈希分桶就是把"加锁排队"变成了"进队列排队",更简单也更快。
三、方案二:单队列 + 批量消费(tcp-server 的做法)
我的行车监控项目用的就是这种模式:
Netty Worker 线程(16个)接收数据
→ DedupService(SHA-1 去重,60s TTL)
→ StatusGuard(状态转换校验)
→ LinkedBlockingQueue<DeviceData>(10000)
→ 单线程每 500ms 取 50 条 → JPA saveAll 批量入库
// DataService 核心逻辑
private final LinkedBlockingQueue<DeviceData> buffer = new LinkedBlockingQueue<>(10000);
@Scheduled(fixedRate = 500)
public void flush() {
List<DeviceData> batch = new ArrayList<>(50);
buffer.drainTo(batch, 50); // 一次最多取 50 条
if (!batch.isEmpty()) {
repository.saveAll(batch); // 批量写入
}
}
优点:
- 顺序绝对保证(单线程 FIFO)
- 数据库压力极小(批量写)
- 代码简单,不容易出 bug
缺点:
- 写入是单线程瓶颈
- 适合每秒几千条,不适合几十万条
四、方案三:Disruptor 环形缓冲区(极致性能)
LMAX Disruptor 是金融级高性能队列,单机每秒处理数百万事件。
// 定义事件
public class DeviceEvent {
public String deviceId;
public String rawData;
public long timestamp;
}
// 创建 Disruptor
Disruptor<DeviceEvent> disruptor = new Disruptor<>(
DeviceEvent::new, // 事件工厂
65536, // 环形缓冲区大小(必须 2 的幂)
DaemonThreadFactory.INSTANCE,
ProducerType.MULTI, // 多生产者(多个 Netty 线程)
new YieldingWaitStrategy() // 等待策略
);
// 多消费者并行处理
disruptor.handleEventsWithWorkerPool(
new DeviceHandler(), new DeviceHandler(),
new DeviceHandler(), new DeviceHandler()
);
disruptor.start();
// 生产端
RingBuffer<DeviceEvent> ringBuffer = disruptor.getRingBuffer();
long seq = ringBuffer.next();
try {
DeviceEvent event = ringBuffer.get(seq);
event.deviceId = "crane-001";
event.rawData = rawLine;
} finally {
ringBuffer.publish(seq);
}
WorkerPool 模式保证:每个事件只被一个消费者处理。但注意——它不保证同一设备的事件被同一个消费者处理。如果需要严格的同设备顺序,要配合哈希路由:
// 按设备 ID 路由到不同 Disruptor 实例
Disruptor<DeviceEvent>[] disruptors = new Disruptor[N];
int slot = Math.abs(deviceId.hashCode() % N);
disruptors[slot].getRingBuffer().publishEvent(...);
五、方案四:Kafka 分区(分布式场景)
Kafka 的 partition 机制本质就是哈希分桶的分布式版本:
// 生产者:指定 key = deviceId
ProducerRecord<String, String> record = new ProducerRecord<>(
"device-data",
deviceId, // key → 决定 partition
jsonData // value
);
producer.send(record);
key="crane-001" → hash % partitions → partition 3 → consumer-3 按序消费
key="crane-002" → hash % partitions → partition 7 → consumer-7 按序消费
同一 partition 内严格有序,不同 partition 并行消费。这就是分布式版的哈希分桶。
六、方案对比
| 方案 | 并行度 | 顺序保证 | 吞吐量 | 复杂度 | 适用场景 |
|---|---|---|---|---|---|
| 哈希分桶 + N 队列 | 高 | 同设备串行 | 10万+/s | 中 | 单机万级设备 |
| 单队列 + 批量(tcp-server) | 低 | 全局有序 | 几千/s | 低 | 单机千级设备 |
| Disruptor | 极高 | 需配合路由 | 百万+/s | 高 | 金融/高频交易 |
| Kafka 分区 | 高(分布式) | 同分区有序 | 无上限 | 中 | 分布式/微服务 |
七、实践建议
- 设备 < 1000 台:单队列 + 批量写入就够了,别过度设计
- 设备 1000~50000 台:哈希分桶,线程数 = CPU 核数 × 2
- 分布式多节点:Kafka 按 deviceId 分区
- 不要为了并行而并行:如果瓶颈在数据库,加线程没用,先优化批量写入
记住那个原则:同设备串行,跨设备并行。所有方案都是这个思想的不同实现。