我接手这个平台的时候,在线设备规模还不到十万,半年后冲上了八百万,再往后半年跨过了千万级。千万级物联网设备接入,听起来是个可以写进PPT里的漂亮数字,但放在后端眼里,真正的拷问是:从设备上云的第一毫秒开始,到数据落到存储、跑完计算、被业务消费,整条链路还能不能稳定扛住。这篇文章不聊云厂商的宣传物料,只讲我们在千万级设备接入这个规模下,怎么设计数据处理链路,以及真正踩过、帮你避开的那些坑。
如果你正准备做一个千万级设备接入的物联网平台,或者正在为现有平台性能发愁,这篇文章可以当成一份设计参考。即便你现在只有几万台设备,文中大部分容量估算方法和架构思路同样适用,因为“向后退一步做设计”永远比“上线后再推倒重来”划算。
1. 千万级接入到底在挑战什么:先算账再看瓶颈
很多人一听到“千万级设备接入”,条件反射是加机器、上中间件、堆资源。但加机器之前先要弄清楚,千万级带来的压力到底是什么。我见过的项目里,有相当一部分不是被平均流量打垮的,而是被设计阶段没算清楚的“量级差异”打垮的。
1.1 一笔账算下来:千万设备每天产生多少数据
先做一道简单的算术题。假设我们有1000万台设备,在线率按70%算,也就是700万台设备同时在线。每台设备每30秒上报一条消息,包含设备ID、时间戳和几个关键测点,整条消息序列化之后约300字节。
平均每秒消息量:700万 / 30 ≈ 23.3万条/秒。
物联网流量从来不是均匀的。白天高峰、整点批量上报、设备定时唤醒,都会把流量抬高到均值的3倍甚至更高。按3倍峰值保守估算,峰值QPS在70万条/秒左右,峰值写带宽约70万 × 300B ≈ 210MB/s。
再算存储。每天消息总量:23.3万条/秒 × 86400秒 ≈ 201亿条。原始数据量约6TB/天。注意,这是压缩前的量。时序数据库通常能压掉6到10倍,所以落到存储层大概是600GB到1TB/天,一年就是200TB到360TB。
这个数字意味着什么?单机肯定没戏,少量几台机器也扛不住。更重要的是,如果设计时对消息量、峰值倍数、压缩率没有任何估计,后面所有选型都是拍脑袋。
1.2 最容易先崩的四个瓶颈点
千万级接入的压力不是均匀分布在所有组件上的,它集中在四个点。
第一是接入层。千万连接对操作系统文件描述符、线程模型、内存都是考验。每个MQTT连接在Broker里占用的内存大约几十KB到两百KB不等,取决于是否持久会话、QoS级别和飞行窗口大小。一千万连接意味着一两百GB起步的内存开销,这还没算业务逻辑。
第二是消息中间件。当生产速率达到几十万QPS时,中间件的分区数、磁盘顺序写能力、副本同步开销都会成为瓶颈。很多团队在几十万设备时用单机Kafka凑合,等规模上来再迁移,代价非常大。
第三是存储层。时序数据的写入模式是持续追加,但物联网数据的查询往往是按设备、按时间段扫描。写入吞吐和查询延迟是一对矛盾,存储引擎设计得不好,数据量一大,写入放大和查询慢会同时出现。
第四是下游消费。告警计算、实时大屏、消息推送、第三方回调,这些逻辑如果和主链路耦合,一旦某个下游处理慢了,背压会一路传导到接入层,最终表现为“设备上报变慢”甚至“连接被断开”。
1.3 从指标反推架构:估量是方案设计的第一步
我习惯在设计架构前,先把一组关键指标写在文档最前面:设备总数、在线率、每设备消息频率、单条消息大小、峰值倍数、数据保留周期。这些数据决定了后续所有的技术选型。
举个例子。如果单设备消息频率是每秒一条而不是30秒一条,那么上面的QPS会从23万直接变成700万,存储量也会膨胀30倍。这时你可能需要边缘网关做数据汇聚,先把高频数据在边缘做聚合降频,再上云。如果设备上报本身就是低频率、大报文,比如图片或文件,那重点就不在消息吞吐而在对象存储和链路带宽。
不同场景对量级的敏感度完全不一样。先算账,再选型,是千万级方案设计的第一步,也是很多人跳过的一步。
2. 接入层设计:把“连得上”和“连得稳”分开做
接入层是整个链路里最敏感的一层。连接数一上去,很多平时看不到的问题就全出来了。我在设计接入层时有个原则:接入只负责连接管理,不要在上面挂业务逻辑。把“连得上”和“连得稳”分开,后面扩容和排障都能省一半力气。
2.1 多协议网关:设备侧从来不是只有MQTT
千万级设备接入,首先要想清楚设备端用的是什么协议。MQTT是物联网的事实标准,但不是全部。智能表计很多走CoAP或UDP私有协议,老旧工业设备走Modbus、OPC-UA,摄像头走GB28181或者RTSP,还有一些NB-IoT设备走LwM2M。期望所有设备都统一到MQTT不现实。
所以在接入层前面要放一层协议网关。它的职责是接受各种协议接入,统一解析数据,转成内部的标准消息结构,再发给后端的Broker集群。这样整个数据链路只认一种内部协议,后续的消息处理、存储、计算都不需要关心设备原始协议是什么。
这层网关必须是无状态设计。连接落在哪个网关节点上不产生业务依赖,这样前面挂负载均衡才能随便扩缩容。网关节点只做协议解析、格式转换、QoS兜底,不保存设备业务状态。设备状态放到独立的会话管理服务里。
2.2 Broker集群:选型、节点数与千万连接的隐性成本
MQTT Broker是整个接入层的心脏。选型上,开源生态里EMQX、VerneMQ、HiveMQ比较常见,自研也不是不行,但成本极高,不推荐在核心链路上重复造轮子。
我实际用过EMQX集群承载百万级连接,它基于Erlang/OTP的进程模型在大量长连接场景下表现确实好。官方宣传单集群可以支撑更高规模,我们实际生产跑到数千万连接也稳定。关键是不能只买软件,要算清楚每个节点扛多少连接。
节点数的估算公式很简单:节点数 = 总连接数 / 单节点安全连接数。单节点能承载多少连接,取决于CPU、内存和会话类型。我们实测下来,一个16核64GB的节点,承载30万到50万轻量连接(cleanSession=true、QoS0/1)压力不大;如果大量设备开持久会话并且QoS2,内存占用会翻几倍,安全水位就得往下调。
这张表是我常用的粗略参考:
场景 | 单节点内存预算 | 单节点安全连接数 | 集群规模估算(1000万连接) 轻连接、QoS0/1、无持久会话 | 约60KB/连接 | 40万 | 25节点 混合QoS、部分持久会话 | 约120KB/连接 | 20万 | 50节点 重度持久会话、QoS2、大消息 | 约250KB/连接 | 8万 | 125节点
注意除了内存,还有网络带宽。单节点40万连接,每条连接即使只是心跳,每秒一个小包,带宽和CPU的中断开销也不小。所以Broker节点建议万兆网卡,系统层面也要调整文件描述符限制、TCP内核参数、并发连接数等。
2.3 鉴权与会话:千万级连接风暴之外的第二个坎
千万级设备接入,还有一个容易被低估的地方:连接鉴权。设备每次重连都要做一次鉴权,如果鉴权逻辑是同步查数据库,那么连接数一上来,数据库先被打挂。
我们当时把鉴权拆成了两层。第一层是网关侧的本地缓存,保存最近验证通过的设备Token,带TTL。第二层是集中鉴权服务,底层用分布式缓存,不直接查关系型数据库。设备证书或Token的吊销走独立的黑名单通道,同步到所有网关节点。这样大部分连接可以直接在网关层识别放行。
会话管理也要单独设计。对物联网平台来说,设备影子(Device Shadow)非常重要,它把设备的期望状态(desired)和上报状态(reported)分开,业务系统不直接操作设备,而是改影子,由影子服务同步下发。一千万设备就意味着一千万个影子对象,存储上要么用高性能分布式KV,要么用时序库加最新值缓存,不能把影子和业务数据库混在一起。
还有一个容易踩的坑是设备离线消息。很多平台默认开启持久会话保存离线消息,设备规模小时没感觉,到了百万级千万级,离线消息会像滚雪球一样堆积在Broker内存里,最终拖垮节点。我的建议是:默认cleanSession=true,关键设备的离线消息走轻量级队列或存储,而不是长期保存在Broker会话里。
3. 数据管道:从设备消息到可计算数据的核心链路
设备数据从Broker出来之后,真正决定系统吞吐和可扩展性的,是消息中间件这一层。这一层设计得好,后面所有下游都能各取所需;设计得不好,接入层再强也会被下游拖死。
3.1 消息中间件选型:Kafka与Pulsar之争
消息中间件领域最核心的选择是Kafka还是Pulsar。这两个我都跑过生产流量,谈不上谁碾压谁,关键是适不适合物联网场景。
Kafka的优势是吞吐高、生态成熟、运维资料多。如果你的核心场景是“高吞吐的流式日志管道”,Kafka是稳妥选择。它的短板是Topic一多、Partition一多,运维复杂度直线上升,而且扩Partition很麻烦,一旦初期分区规划不足,后面加分区很难做到无损。
Pulsar的优势是存算分离,Broker不保存数据,存储层用BookKeeper,扩容时可以独立扩展Broker或存储节点。它对多租户和大量Topic的支持比Kafka好,更贴近物联网平台“一个产品线一个Topic域”的隔离需求。代价是组件多、运维门槛高,集群占地面积也更大。
我们最终选了Kafka,核心原因是团队对它的运维经验最足,而且平台的核心流量适合用有限的Topic域承载。但如果你要做多租户物联网平台,每个客户都要独立隔离,我建议认真评估Pulsar。选型这块没有银弹,全看团队能驾驭哪个。
3.2 Topic和Partition规划:分区方案错了很难回头
Topic和Partition的规划,是数据管道设计里最容易在后期付出代价的事。
首先,不要给每台设备建一个Topic。物联网设备的消息是海量小消息,如果用一千万个Topic,Kafka的元数据压力会先让你崩溃。生产上我们按产品线(ProductKey)和业务域建Topic,例如设备原始数据一个Topic、设备事件一个Topic、设备生命周期一个Topic。这样的数量级控制在几十个Topic,才能真正利用好Kafka。
Partition数量要按目标吞吐预留。经验值是一条Partition的生产写入吞吐在小消息(几百字节)场景下每秒能扛几万条,峰值不要超过10MB/s。拿上面70万QPS的峰值为例,预留3到5倍余量,Partition总数设计在100到200个是合理的。
Partition数量直接决定消费者的并发上限,所以在建集群时就要想清楚:以后数据量翻倍了,是重建Topic还是继续加Partition。Kafka的Topic一旦创建,Partition只能增加不能减少,增加会触发数据重平衡,对生产影响不小。所以初期宁可多建一些,也不要抠抠搜搜。
分区键也有讲究。同一设备的消息如果分散到不同Partition,下游做顺序性处理会很头大。我们规定设备原始数据Topic的分区键一律用设备ID哈希,保证同一设备消息落在同一Partition;跨设备的全局顺序性在物联网场景里没意义,不用强求。
3.3 削峰与背压:当三十万台设备在同一分钟上报时
物联网流量有个显著特征:突发性强。最常见的场景是整点批量上报,比如智能电表在每小时第0分钟集中上报一次,瞬时流量可能是平时几十倍。当作流量的均值做设计,不做峰值做设计,系统在第一个整点就会被打穿。
消息中间件天然是削峰缓冲层。生产端把消息写进Kafka就返回,Kafka用磁盘顺序写和页缓存吸收流量尖峰;下游消费者按自己的速率拉取,不需要跟上游同步。这个机制让Kafka能承受几倍于均值的瞬时流量。
但削峰不是无限的。生产端要注意发送超时和批量参数。我们在线程模型上控制发送端的max.block.ms和linger.ms,避免高延迟场景下生产者线程全部阻塞。消费者端的背压指标是消费Lag,一旦Lag持续增长,说明下游消费能力不足,这时候不是加机器就能解决的,要先找下游的瓶颈在哪,比如存储写入慢、外部API响应慢。
还有一个策略很多人会忽略:拒绝与降级。当流量大到系统真的接不住时,要有优先级。低价值的高频数据可以降采样或暂时丢弃,保连接、保关键数据链路。我见过不少系统在极端流量下死扛,结果核心数据全丢了,还不如主动降级保住最重要的部分。
4. 存储与计算:每天几TB数据怎么放、怎么算
到了存储和计算这一步,核心矛盾变成了成本和实时性的权衡。千万级设备的原始数据量一天就是几TB,不做分层的话,存储成本会以惊人的速度膨胀。
4.1 时序数据库选型:写入模型决定天花板
物联网的数据绝大部分是时序数据,选型时我主要看三点:写入吞吐、压缩比、查询能力。常用选项包括TDengine、IoTDB、InfluxDB、TimescaleDB,各有用武之地。
TDengine和IoTDB是物联网原生设计的时序库。TDengine的超级表模型很契合“同类型设备批量建表、统一查询”的需求,写入和压缩性能都很突出,开源版支持集群。IoTDB在工业物联网里用得很多,双内核(时序+日志)设计适合复杂工业场景。InfluxDB生态好、上手快,但开源版主要面向中小规模,集群能力没有前两者强。TimescaleDB是PostgreSQL扩展,功能全面但写入吞吐上限相对有限,适合数据量不大但对SQL兼容性有要求的团队。
写入模型上要留意一个原则:批量写入,不要单条插入。时序数据库的写入引擎大多是LSM Tree风格,单条小消息写入会产生严重的写入放大,几十万QPS打到存储层会非常吃力。我们把Kafka里的消息攒一攒,按1000条或1MB一个批次写入时序库,写入吞吐提升了接近一个量级。
这里再提醒一下:不要因为时序数据库支持“每设备一张表”就真给每台设备建一张表。一千万台设备建一千万张表,管理和元数据开销都是灾难。要用超级表(或等效模型)统一管理同类型设备,把设备ID作为标签而不是表名。
4.2 冷热分离与TTL:存储成本的大头在这里
存储成本是我在千万级项目里最为关注的问题。如果所有数据都用同一套高吞吐存储,第一年就能吃掉整个预算的一大半。
我们的方案是把数据分成三层。热数据——最近7天——放在时序数据库里,提供秒级查询。温数据——最近90天——做压缩后转储到更便宜的存储,仍然可以按需查询。冷数据——超过90天——落对象存储加列式文件格式,配合批量查询引擎做分析。
TTL不是简单地设一个删除时间,而是要和业务保留需求对齐。有些数据业务上要留三年,有些数据本身只有几分钟的热度。我们按数据类型设不同TTL,原始高精度数据保留7天转冷,统计聚合数据保留90天,账单事件类数据保留三年。这样既不超卖存储,也不会在审计追查时拿不出数据。
冷数据存储格式建议使用列式存储,比如Parquet或ORC,并按时间分区存储。一列存一天的数据,查询可以只扫描需要的列和分区,分析性能比直接打开几TB的文本文件好得多。压缩率和成本优势在千万级规模下非常明显。
4.3 流式处理与告警:Flink在链路里的位置
数据落到存储之后,还有一大块工作是实时计算。告警、在线率统计、大屏指标、数据清洗,这些都靠流式计算完成。我们用的计算引擎是Flink,准确说,Flink承担了数据处理链路里“智能”的部分。
Flink从Kafka消费设备消息,经过规则引擎做告警判断,命中就把告警事件写入告警Topic,同时更新实时指标。整个过程是毫秒到秒级延迟。这里有个设计细节:不要把所有计算都压在一个Flink作业里,否则一个规则的升级要重启整个作业。我们按业务域拆成多个作业,一个作业负责规则告警,一个负责指标聚合,一个负责数据清洗回写。作业之间通过Kafka解耦。
状态管理要重视。Flink做窗口聚合时,窗口状态存在内存和后端状态存储里,假如每分钟统计一次全网设备在线状态,窗口状态和事件时间处理要仔细校准。我建议在非必要场景坚持processing time,不要一上来就搞event time加watermark,那套复杂度在千万级大流量下会被放大很多倍。
流式计算的反压问题也需要提前设计。下游存储写入慢时,Flink会把反压传导到Kafka消费者,最终体现在消费Lag增长。监控Lag变化趋势要比监控Flink各算子繁忙度更直观,这也是我们把消费Lag作为核心监控指标之一的原因。
5. 实测踩过的坑:连接风暴、乱序与扩展陷阱
设计稿再完美,都要经过真实流量的毒打。这一节复盘几个我们实际踩过、也花了不少时间才填平的坑,每一个都值得你在架构设计阶段提前预防。
5.1 凌晨四点,上百万设备同时重连
那是某次区域大面积停电恢复的凌晨。电网一恢复,几十万台设备几乎在同一时刻开始重连,紧接着相邻区域也陆续恢复,短时间内新连接请求数直接冲到每秒几十万。Broker的TCP握手、TLS握手、鉴权请求一下全堵在入口,CPU打满,大量正常连接反而被踢下线。
这次事故让我们上了三堂课。第一,设备端必须做错峰重连。连接断开后采用指数退避加随机抖动,第一次重连延迟几秒,后面逐步加大间隔,把“百万设备同时涌上来”拆成“分散在十几分钟里陆续恢复”。第二,Broker集群要预留新连接建立速率的容量。一个Broker节点每秒能处理的新连接数量是有限的,几万级别都算高,让所有设备同时重连本身就是不可能完成的任务。第三,接入网关要能自动摘除异常节点。当节点CPU或连接数超过安全水位时,负载均衡层要把它摘掉,避免雪崩。
5.2 消息乱序与去重:从“收到数据”到“数据可信”
设备上报数据的可靠性,不只是“消息有没有到”,还包括“到的数据对不对”。我们在上线初期就遇到过消息乱序导致的设备状态回跳:设备上报温度是35度,再上报是36度,但消费端收到的顺序反了过来,最后库里记成了35度。
乱序的来源很多。设备端网络异常重传、MQTT会话迁移、Broker分发到不同分区、消费者重启后的重平衡,都可能导致后发的消息先被处理。解决乱序没有万能药,只能分场景处理。
同一设备的数据流,我们通过固定分区键保证进同一Partition,同时消息体里带设备端递增的序列号,消费端根据序列号丢弃过期消息。不同设备之间无所谓顺序,不用管。对跨系统的数据一致性需求,比如设备上报和平台下发指令的顺序,靠的是业务层面的幂等和状态机校验,而不是消息系统保障。
去重是个容易被低估的成本。完全精确的一次语义(Exactly-Once)在物联网场景下代价很高,尤其是涉及外部存储和回调时。我们采用“至少一次消费 + 幂等写入”的方式:数据库表用“设备ID + 消息序列号”做唯一键,重复消息直接丢弃。去重表只保留最近几天的键,过期清理,避免存储无限膨胀。
5.3 “加了机器反而更慢”的扩展陷阱
千万级规模下,加机器并不总是生效。我们第一次横向扩容Broker时,连接数和消息量没有显著提升,反而出现了一部分设备连接不稳定的情况。
排查下来,问题出在负载均衡层。新加的Broker节点权重没有被正确识别,旧节点连接依然处于满负荷状态,新节点却空闲。把负载均衡策略从轮询改成基于活跃连接数的动态分配后,连接才真正均匀分布。
Kafka消费者加机器也有同样的问题。如果只是增加消费者实例,但Partition数量没变,新增实例并不会自动分担消费压力,因为一个Partition在同一时刻只会被一个消费者实例消费。消费者并发上限由Partition总数决定,所以提高下游消费能力,要么加Partition,要么保证每个消费者处理得更快,单纯加机器是没用的。
还有一个很多人容易忽略的坑:消费者组重平衡风暴。当消费者实例频繁加入退出(比如部署更新时滚动重启),或者心跳超时,会触发Rebalance,在Rebalance期间整个消费者组会停止消费,Lag瞬间飙高。我们把消费者实例数量和Partition数量对齐,并调大会话超时参数,才把重平衡风暴压下去。
6. 拿什么证明能扛千万级:压测与灰度的实战节奏
讲完设计思路和踩坑经历,最后聊一个几乎所有团队都会忽略、但恰恰最要命的问题:怎么证明你的系统能扛千万级。没有验证过的架构,只能叫纸面方案。
6.1 压测指标与场景设计:不能只测QPS
很多人做压测,只盯一个QPS,跑到目标值就认定系统达标了。在千万级物联网场景下,这是远远不够的。我建议至少盯四类指标:
- 吞吐与容量类:每秒消息数、每秒新建连接数、峰值连接数、消息积压数。
- 时延类:端到端时延(从设备上报到数据可查询)的P99和P999,以及Broker确认延迟。
- 可靠性类:消息丢失率、消息重复率、消费Lag变化趋势。
- 资源类:CPU使用率、内存水位、GC频率、磁盘IO、网络带宽、文件描述符用量。
压测场景不能只有“匀速压测”一种。我们在压测环境里设计了几个固定场景:目标QPS持续运行1小时看稳定性;3到5倍峰值流量打30分钟看削峰和积压恢复能力;百万设备模拟同时上线看连接风暴应对;随机kill一个Broker节点和Kafka节点看故障自愈。
压测工具方面,MQTT场景可以用emqtt-bench或者自研模拟网关。我们自研了一套设备模拟器,可以按真实业务的设备数量、上报频率、报文大小生成流量,比用通用压测工具更贴近实际。
6.2 从一万台到一千万台:渐进式灰度路径
压测通过不代表直接切生产。我们没敢让千万台设备一次性接入新平台,而是走了一条渐进式灰度路径,每一步都定一个验证目标。
第一阶段,接入1万台设备,跑一周,验证接入、消息管道、存储全链路功能。这时候可能暴露的是配置错误、字段问题,而不是性能问题。第二阶段,放到100万台,重点看Broker集群和Kafka的吞吐余量,以及消费端是否跟得上。第三阶段,放到500万台,这时候开始压真实瓶颈,比如整点峰值是否打穿削峰层、时序库批量写入是否稳定。第四阶段,放到千万台,这时候基本是查漏补缺和极限调优。
每一阶段的灰度,都配合流量镜像和影子模式。影子模式的意思是,新老平台同时接收消息,新平台只计算不返回结果,跑一段时间后对比两边数据是否一致。等影子数据完全对齐,再逐步切真实流量。这套做法看起来很保守,但在千万级规模下,保守是最快的路径。
6.3 监控大盘与容量水位:长期盯住这七类指标
系统上线不是终点,日常运维才是。我们内部长期盯着一套容量大盘,七个指标按优先级排列:连接数、每秒新连接数、生产QPS、消费Lag、存储写入速率、存储磁盘水位、端到端时延P99。
前三个指标反映接入层健康度。一旦连接数逼近集群安全水位,或者每秒新连接数异常上涨,就要提前预警。消费Lag是下游消费能力的温度计,Lag长时间增长,说明消费链路有瓶颈。存储写入速率和磁盘水位决定你能在故障发生前留出多少时间去扩容或清理数据。端到端时延P99则是最直观的用户体验指标,设备数据晚到几秒,对某些实时业务是致命的。
监控告警之外,每个月我们还会做一次容量复盘:这个最高峰值是多少、各环节余量还剩多少、按当前增速多久会触顶。这个习惯让我在多次流量翻倍前提前搭建扩容计划,而不是等到报警响了再救火。
千万级物联网设备接入的数据处理方案,说到底是把一个“看起来很大的数字”拆解成连接、吞吐、存储、计算四个具体问题,再逐个击破。架构没有完美的,只有能在你的场景下稳定扩张的。希望这篇文章能给你在设计路上提前排掉几个雷,剩下的,交给真实流量来检验。