IoT设备数据的多线程并行处理

星期五, 7月 31, 2026 | 3分钟阅读 | 更新于 星期五, 7月 31, 2026

@

物联网场景下,成百上千台设备每秒各发一条数据。单线程处理不过来,多线程又可能打乱顺序。本文记录如何用多线程并行处理 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
设备量 < 10004~8
设备量 > 1000016~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);
        }
    });
}

能用,但问题很多:

  1. 锁对象越来越多:设备几万台,就要维护几万个锁对象,内存浪费
  2. 设备下线后锁不清理:内存泄漏
  3. 锁竞争:如果 processData 里有共享资源(比如数据库连接池),不同设备的线程还是在争抢
  4. 调试困难:死锁、锁等待,出问题很难排查
  5. 线程池 + 锁的组合本身就容易写出 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 分区高(分布式)同分区有序无上限分布式/微服务

七、实践建议

  1. 设备 < 1000 台:单队列 + 批量写入就够了,别过度设计
  2. 设备 1000~50000 台:哈希分桶,线程数 = CPU 核数 × 2
  3. 分布式多节点:Kafka 按 deviceId 分区
  4. 不要为了并行而并行:如果瓶颈在数据库,加线程没用,先优化批量写入

记住那个原则:同设备串行,跨设备并行。所有方案都是这个思想的不同实现。

© 2026 My Blog

🌱 Powered by Hugo with theme Dream.