极简架构在IoT平台中的项目复盘:设备接入层的高并发设计经验
一、项目场景与约束
某工业物联网项目:5000台设备,每台每5秒上报一次数据(温度、湿度、振动频率)。计算得出约1000 QPS的数据写入。设备运行在3G/4G网络下,网络不稳定,需要支持数据补传。
技术约束与之前类似:团队3人,运维能力有限,不能引入Kafka或Kubernetes。需要在一台4C8G的云服务器上完成设备接入、数据存储和API服务。
二、极简设计的四个核心决策
决策一:自定义二进制协议替代MQTT Broker。
标准做法是搭建Mosquitto或EMQX做MQTT Broker。但对5000台设备而言,引入一个独立需要运维的消息中间件性价比不高。自定义了一个极简的二进制协议:
| Header(2B) | Length(2B) | Type(1B) | DeviceID(8B) | Timestamp(8B) | Payload(NB) | CRC(2B) |单个上报消息约30-40字节。对比MQTT的MQTT CONNECT + PUBLISH流程(约200+字节),自定义协议节省了约80%的网络开销。对3G/4G网络下的设备而言,低带宽意味着更快的上报速度和更低的流量费用。
决策二:连接复用而不是短连接。
初始版本使用HTTP POST上报数据。5000台设备每5秒一次就是每小时360万次TCP握手。TCP的SYN+ACK+ACK三次握手在高频场景下占用了大量CPU。
改为TCP长连接 + 连接池管理:
type DeviceConnPool struct { mu sync.RWMutex conns map[string]*DeviceConn // deviceID -> conn maxConns int timeout time.Duration } type DeviceConn struct { conn net.Conn lastSeen time.Time writeLock sync.Mutex // 每个连接独立锁,避免全局锁竞争 deviceID string }每台设备建立连接后保持。空闲超过5分钟才关闭。单服务器可维护5000个并发连接,Go的goroutine模型天然适合这种场景。内存占用约80MB。
决策三:批量写入减少数据库压力。
1000 QPS的数据写入如果每条单独INSERT,PostgreSQL很难支撑。采用三级缓冲:
- 内存RingBuffer:累计100条或100ms后触发一次批量INSERT
- Redis缓冲:设备断线时的数据暂存,重连后批量补传
- PostgreSQL时序分区:按设备ID Hash分16个分区,避免写入热点
type BatchWriter struct { buffer chan DataPoint batchSize int ticker *time.Ticker db *sql.DB } func (w *BatchWriter) Run(ctx context.Context) { batch := make([]DataPoint, 0, w.batchSize) for { select { case dp := <-w.buffer: batch = append(batch, dp) if len(batch) >= w.batchSize { w.flush(batch) batch = batch[:0] } case <-w.ticker.C: if len(batch) > 0 { w.flush(batch) batch = batch[:0] } case <-ctx.Done(): if len(batch) > 0 { w.flush(batch) // 优雅关闭 } return } } } func (w *BatchWriter) flush(batch []DataPoint) { // PostgreSQL COPY协议批量插入,比INSERT快10倍 tx, _ := w.db.Begin() stmt, _ := tx.Prepare(pq.CopyIn( "device_data", "device_id", "timestamp", "metric_type", "value", )) for _, dp := range batch { stmt.Exec(dp.DeviceID, dp.Timestamp, dp.MetricType, dp.Value) } stmt.Close() tx.Commit() }优化后写入延迟:单条INSERT方案P99约120ms → 批量方案P99约15ms。
决策四:Redis做热数据缓存。
最新的数据(最近1小时)放在Redis中。查询最近数据的API从Redis读,毫秒级响应。历史数据走PostgreSQL。Redis使用Sorted Set,member是deviceID:metricType,score是时间戳,方便范围查询。
三、数据补传机制
设备在3G/4G网络下频繁断线。断线期间的数据需要缓存和补传:
func (s *Server) handleReconnect(deviceID string, conn net.Conn) { // 设备重连后,检查断线期间是否有未上报数据 lastReportTime := s.getLastReportTime(deviceID) gap := time.Since(lastReportTime) if gap > 30*time.Second { // 发送补传指令 s.sendCommand(conn, Command{ Type: CmdBatchRetransmit, DeviceID: deviceID, Payload: marshalTimeRange(lastReportTime, time.Now()), }) } // 接收补传数据 s.receiveRetransmit(conn, deviceID) } func (s *Server) receiveRetransmit(conn net.Conn, deviceID string) { // 超时保护:补传不超过60秒 deadline := time.Now().Add(60 * time.Second) conn.SetReadDeadline(deadline) for { packet, err := s.readPacket(conn) if err != nil { break } if packet.Type == PacketRetransmitEnd { break } s.batchWriter.buffer <- packet.ToDataPoint() } }四、性能边界与可扩展性
当前系统性能:
- 5000台设备稳定在线,1000 QPS数据写入
- CPU使用率约35%,内存约1.2GB
- 写入P99延迟约15ms,读取P99延迟约5ms(Redis热数据)
已知边界:
- 单服务器连接上限约2万台设备(TCP文件描述符限制 + 内存),超出需要水平扩展
- 批量写入的100ms缓冲窗口意味着数据有100ms的延迟,对秒级实时性有影响
- PostgreSQL时序分区方案在数据量超过10亿条后需要归档策略
当设备数突破1万台时的扩展路径:
- 先加服务器做水平扩展——通过一致性哈希将设备分流
- 引入TSDB(如TimescaleDB)替代原生PostgreSQL,提升时序查询性能
- 但不需要引入Kafka——当前的自定义TCP + 内存批处理方案在10万设备级别仍然有效
五、总结
IoT平台极简架构的核心原则:根据真实负载选型,不提前为"可能的扩展"引入复杂度。
关键决策:
- 5000设备的规模不需要MQTT Broker——自定义二进制协议+TCP长连接足够
- 批量写入是性价比最高的优化,PostgreSQL的COPY协议比单条INSERT快10倍
- Redis做热数据分层,兼顾实时查询和存储成本
- Go的goroutine天然适合高并发连接场景,5000连接仅80MB内存
最大的教训:初期尝试了标准的MQTT + InfluxDB组合,但运维复杂度和资源消耗远超预期。回滚到自定义协议+PostgreSQL的方案后,开发时间反而缩短。这再次验证了那句话:在可用范围内,简单方案永远比"标准方案"更可靠。