做物联网平台开发的时候,“支持百万设备连接”这句话经常被写在路演PPT首页。但真要把IoT接入这件事做好,不是堆几台服务器、起个MQTT Broker就完事了。我自己这两年带团队做过一个物联网接入平台,从最开始只能撑几千台设备,一步步压到百万连接级别,中间踩过的坑、验证过的方案、被现实教育过的认知,都在这篇里了。如果你正在评估自研平台选型,或者刚接手类似项目,这篇文章应该能帮你少走好几个月弯路。
先说清楚,我讲的“百万连接”不是拿压测工具刷出来的假连接数,而是真实生产环境里,百万台设备能持续在线、按时上报数据、随时接受下发指令,并且在网络抖动、设备掉线、半夜集中上线这些混乱场景下不崩、不丢、不乱。这个目标拆开来看,每一个部分都有一堆细节要处理。
1. 百万连接这个目标,到底难在哪里
很多人一听到“百万连接”第一反应是上ARM服务器、改内核参数、堆带宽。其实这些只是最后一步,真正的难点藏在更底层的地方。
第一道坎是连接本身。一个TCP长连接驻留在网关节点上,不只是占一个文件描述符那么简单。它背后有Socket缓冲区、协议解析上下文、待发送队列、认证会话、订阅关系。一个空连接大概占20~40KB内存,100万个空连接就需要20~40GB内存,这还没算心跳报文和消息转发产生的临时对象。如果你的设备还订阅了大量Topic,这个内存开销还要翻倍。
第二道坎是消息吞吐。连接在线不等于平台有用,设备要上报数据,平台要下发指令。100万台设备,假设每台每30秒上报一条数据,平均每秒大概3万条消息;如果业务场景要求每台每秒一条,那就是100万TPS。这个量级下,任何数据库直写、同步调用都是灾难,必须从一开始就把链路设计成异步管道,中间每一级都要能缓存、能削峰、能重试。
第三道坎是状态管理。连接不只是一条socket,还牵扯到设备在线状态、离线消息、遗嘱消息、Session会话、订阅关系。100万个连接同时断线重连时,Broker要处理海量会话恢复,业务系统要收到精确的在线/离线事件,不能靠全量扫描来判断“谁掉线了”。我之前见过一个平台用数据库轮询方式判断设备在线状态,每两分钟扫一次,百万设备之后数据库直接被打爆——这种设计思维从一开始就不该出现。
第四道坎是网络现实的复杂性。真实设备不是在干净的机房里,可能在工厂角落、地下室、高速移动的车辆上。网络随时会断、IP随时会变、设备随时会被运营商NAT踢掉。弱网环境下TCP连接经常处于半开状态,服务器以为连接还在,设备其实早已失联。这种问题Web开发基本不会遇到,但IoT平台必须把心跳超时、链路探测、断线重连当成一等公民来设计。
所以百万连接的本质,不是“能建多少条TCP链路”,而是“在真实混乱的网络环境下,平台能稳定维护多少条有业务的连接”。想清楚这点,架构设计才不会跑偏。
2. 接入层设计:协议选择与网关架构
2.1 先看清楚设备端长什么样
接入层方案不是拍脑袋定的,主要看设备端的硬件条件和网络环境。
| 协议 | 连接方式 | 报文开销 | 实时性 | 典型场景 |
|---|---|---|---|---|
| MQTT | 长连接,发布订阅 | 低 | 高 | 智能家居、工业设备、车联网 |
| CoAP | UDP,类HTTP | 极低 | 中 | 低功耗传感器、NB-IoT设备 |
| HTTP(S) | 短连接 | 高 | 低 | 上报型设备、无状态简单采集 |
| 私有TCP | 长连接,自定义报文 | 最低 | 高 | 嵌入式设备、已有存量协议 |
我的经验是:除非设备端已经被硬件成本压到极致,只能用CoAP或者早就固化好了私有协议,否则优先用MQTT。原因很实际,MQTT的生态最成熟,Broker开源实现多,客户端SDK覆盖几乎所有的MCU、RTOS、Linux、Android平台,QoS机制能解决消息可靠性,心跳、遗嘱、持久会话这些特性简直是IoT场景量身定做的。CoAP更适合想省功耗、带宽受限的节点,但要做消息可靠传输时,很多东西得自己补,开发成本不低。HTTP轮询是最省事的,但一来实时性差,二来百万设备低频轮询也是每秒上万的短连接请求,对网络的冲击不一定比长连接小。
2.2 网关集群怎么搭才能水平扩展
连接接入层必须做到无状态,这是水平扩展的前提。每个网关节点只负责维护当前这台机器上的TCP连接,不跨节点存储任何会话数据,节点可以随时摘掉、随时加新。负载均衡层用四层转发,LVS、HAProxy或者云上的NLB都行,把设备的TCP流量均匀分发到后面一组接入节点上。
Broker方面,除非你的团队有足够的网络协议栈功底,确定能长期维护一个自研的MQTT Broker,否则我建议第一版直接用开源Broker,EMQX或者VerneMQ都是经过大规模场景验证的。EMQX的集群方式是把节点组成一个分布式消息骨干,设备连到任意节点都能访问完整的设备订阅矩阵,对业务侧透明。版本选择我个人不会去追新,而是选一个长期维护的稳定分支,集群规模上去之后,升级Broker是一个非常痛苦的动作。
自研网关作为Broker前置代理的情况也可以考虑,很多团队为了统一接入私有协议,会写一层协议网关,把私有TCP/CoAP转成MQTT再进Broker。这个方案没问题,但要注意网关本身不能有状态,连接信息和消息转换要纯函数化,否则节点一重启,下面几十万个设备连接全部断掉,业务侧会被重连风暴淹没。
2.3 设备认证与连接鉴权
百万设备在公共网络上,不可能放开门禁随便连。我见过设计得极其简陋的方案:设备连接的时候发一个明文clientId和密码,Broker侧查一下数据库有没有这个设备。设备少的时候没问题,设备多了之后,每次连接都查库,连接风暴时数据库先挂掉。
更稳的做法是认证前置。接入网关收到设备连接请求时,先做一次高效的校验,比如HMAC签名或JWT,一次性把认证结果放进内存或Redis,而不是让每个连接都穿透到关系数据库。设备连接时推荐使用三元素:ProductKey、DeviceName、DeviceSecret,设备端用DeviceSecret对当前时间戳或者随机数做签名,网关秒级验签。
连接鉴权通过后,还要做连接限流。设备侧的Bug、网络运维的误操作,或者纯恶意攻击,都可能造成每秒数万次连接握手。服务器端必须有令牌桶限流,每节点每秒钟最多处理多少握手,超过的直接断开,并记录日志。最终效果是:即使全网设备同一时间重启,平台也不至于被自己的设备打死。
# 伪代码示例:设备连接认证逻辑(网关侧) def authenticate_device(product_key, device_name, sign, timestamp, expected_window=300): device_secret = get_device_secret(product_key, device_name) if device_secret is None: return False, "device not registered" raw = f"{product_key}{device_name}{timestamp}" calc_sign = hmac_sha256(device_secret.encode(), raw.encode()).hexdigest() if abs(time.time() - timestamp) > expected_window: return False, "timestamp expired" if hmac.compare_digest(calc_sign, sign): return True, "ok" return False, "sign mismatch"3. 消息链路:从设备到业务的实时数据管道
3.1 Topic规划是第一步,也是最容易马虎的一步
接入层稳定只是地基,真正让平台有业务价值的,是设备消息从接入层到业务系统这条数据管道。而管道的第一站,是Topic怎么设计。
参考实践里,我一般分三层:
- 上行消息:
{productKey}/{deviceName}/up,设备上报属性、事件、日志。 - 下行指令:
{productKey}/{deviceName}/down,业务系统下发控制、配置、OTA指令。 - 内部系统消息:
sys/{productKey}/{deviceName}/connected|disconnected|status,Broker产生的生命周期事件。
Topic层级不要放太多动态字段,动态字段全放payload。原因是订阅匹配是有成本的,层级越深、通配符越多,Broker的路由压力越大,百万连接配百万级订阅树的时候,这个性能差异会被放大。
消息体建议统一格式,至少包含消息ID、设备标识、上报时间、业务数据。消息ID特别重要,它是去重的唯一依据。设备端弱网重传、Broker转发重试、消费端重复消费,这些情况在分布式系统里是必然的,没有一个全局唯一的消息ID,后面任何幂等逻辑都无从做起。
{ "msgId": "a89f0c2e-4c6a-4ae2-9b1e-8aef1e5b023d", "productKey": "pk001", "deviceName": "dev_0001", "timestamp": 1716000000000, "data": { "temperature": 26.5, "humidity": 48.3 } }3.2 为什么消息必须先进消息队列,而不是直接写库
这是整个数据链路设计里最关键的一步,也是最容易走弯路的一步。
设备上报的流量是突发性的。凌晨两点设备集中升级、厂家做活动集中开启、故障恢复后大批设备上报,瞬时流量可能比平时高几十倍。如果消息链路是“接入层 -> 业务接口 -> 数据库”,那数据库QPS直接跟着设备流量走了,平时能扛,尖峰一来就撕裂。
正确做法是接入层只负责把消息投递到消息队列,比如Kafka或RocketMQ,让队列做削峰填谷,业务系统按自己的消费能力平稳拉取。这里Kafka胜在吞吐高、分区机制适合并行消费;RocketMQ胜在事务消息和延迟消息能力更强。没有绝对的谁好,看团队熟悉哪个,但队列必须选一个,不能省。
消费端拿到消息以后再做分发,分成几路:
- 规则引擎链路:根据消息内容触发告警、联动。
- 时序数据链路:写InfluxDB、TDengine或云上的时序数据库,存历史趋势。
- 实时状态链路:更新Redis里的最新设备属性,供App查询。
- 离线计算链路:进数仓做设备画像、质量分析。
这些链路各自独立,互不阻塞。任何一路挂了也不能影响平台主干,这才叫可运维的消息管道。
# 伪代码示例:消费端幂等处理 def process_message(msg): msg_id = msg["msgId"] if redis.sismember("consumed_msgs", msg_id): return # 已消费过,丢弃 try: handle_business_logic(msg) redis.sadd("consumed_msgs", msg_id) redis.expire("consumed_msgs", 86400) except Exception as e: log.error("handle msg error: %s, msgId=%s", e, msg_id) raise # 依赖队列的重试3.3 序列化和流量治理
消息体用什么序列化,直接影响带宽和性能。全链路用JSON最直观、排障最容易,但100万设备、每台每秒一条消息,JSON解析的CPU开销是个不小的数字。我的建议是内部链路(网关到Kafka、Kafka到存储)用Protobuf或MessagePack这类二进制格式,只有到业务展示层再还原成JSON。如果做不到全链路改造,至少要在网关做一次格式转换,把设备端上报的紧凑二进制转成标准消息,能省不少存储和传输成本。
流量治理上,单设备消息频次要限流。比如一台温度传感器正常1秒一条,如果因为固件Bug改成每毫秒一条,它的连接不会变,但整个链路的吞吐会被浪费。网关侧配置每个设备的QoS和消息速率上限,超速的丢弃一定比例或直接告警,不能让它拖垮整平台。另外,设备上报的字段也要做大小限制,我见过嵌入式工程师把调试日志打到业务Topic里,一条消息几十KB,百万设备跑起来直接引爆带宽。
4. 设备注册与状态管理:百万台设备如何不丢不漂
4.1 身份设计与注册流程
每台设备的身份要从出生那一刻就确定,并且不可变。我的做法是三元组:ProductKey + DeviceName + DeviceSecret。ProductKey对应产品型号,DeviceName是设备唯一标识,DeviceSecret用于签名认证。设备量只有几十台时,随便存个自增ID没问题;到百万级,deviceName要采用有业务含义的编码规则,比如厂商代码+产线代码+流水号。这样排查问题时,看到一个设备名就能大概知道它是哪条产线、哪批次的,否则日志里全是无意义的随机串,排障效率极低。
设备注册有两种主流方式。一种是产线预注册,设备出厂前批量导入系统,激活时校验。另一种是动态注册,设备第一次上电时用产品级密钥申请一个设备级密钥,平台生成DeviceName和DeviceSecret返回。动态注册比预注册灵活得多,适合终端用户自己买设备、自己激活的场景,但需要设计好防滥用机制:同一台设备反复注册,返回同一个DeviceName;不同硬件但用同一产品密钥批量申请,要限制频率。
设备基础信息库不建议放在Redis里,因为百万设备冷数据占内存太多。MySQL单表存百万级可以,但建议直接按产品分表,或者用TiDB这类分布式数据库。查询路径要设计好,设备列表按产品分页查询,设备详情按DeviceName精确查询,尽量避免模糊匹配。Redis只放热数据,比如最近一小时有上报的设备属性快照、连接状态、影子数据。
4.2 在线状态不能靠轮询
在线状态这事,最忌讳的就是“定时扫库”。百万设备每5秒来一次心跳,如果你每5秒select一次所有设备,哪台设备多久没心跳,数据库早就被拖垮了。正确做法是事件驱动。
Broker在设备连接、断开、Session过期时会发出生命周期事件,网关把事件投递到状态Topic。一个状态消费者负责写Redis,用hash结构维护online:{productKey},每台设备的在线状态用一个key表示,value是连接时间戳或逻辑时钟。业务系统查询设备状态时直接读Redis,毫秒级返回。
“在线”不能只看TCP连接,还要结合业务心跳。很多设备连接还在,但应用层已经卡死、物联卡流量耗尽,数据一个都不上报。这时候设备状态应该算“半在线”。我的处理是:连接层面由Broker管,业务层面由平台管,超过N分钟没有收到任何业务消息,就标记为“数据失联”并告警。在线状态的分级比二元的在线/离线更实用。
4.3 设备影子和断线重连
下发指令的时候,设备刚好不在线怎么办?不能直接丢。业界通用的做法是设备影子(Device Shadow)。平台保存两个状态:期望状态(desired)和实际状态(reported)。业务系统修改desired,如果设备在线立刻推送;如果不在线,状态先存影子,等设备一上线,自动从影子拉取最新期望并执行。
这样一来,下发指令和设备在线状态之间就解耦了,业务系统不用担心“设备现在不在线,指令丢了怎么办”这种问题。
断线重连是压垮接入层的经典场景。请一定要在设备SDK里做随机退避,最简单的是指数退避加抖动:第一次重连延迟1秒,第二次2秒,第三次4秒,最大到5分钟,然后加一个随机0~1000毫秒的偏移量。不能所有设备同时重连,否则就是一台聪明的“垃圾”设备,背后放枪打烂平台。平台侧也要配合,连接认证接口做限流,把请求排到令牌桶里,允许少量等待,超出直接拒绝。
5. 压测实践:百万连接是怎么一步步压出来的
5.1 连接压测必须真实建连
做压测最忌“估算”,最怕拿几个客户端开多线程跑,数据一填就算百万。真实的长连接压测,一定要有足够多的压测机,每台机器用不同源端口建立真实TCP连接,然后维持心跳、发送真实消息。
工具方面我常用的有两类:
- EMQX自带的
emqtt_bench,支持并发连接、发布、订阅压测,简单直接,适合验证Broker容量。 - 基于Netty自研的压测客户端,可以模拟真实设备的连接行为(连接、心跳、不定时上报、掉线重连),适合验证整个平台链路,不只是Broker。
如果你愿意折腾,JMeter的MQTT插件也能用,但要模拟百万连接不现实,JMeter本身会成为瓶颈。不要听别人说“JMeter能压百万”,测试工具扛不住的时候,你会误以为是平台扛不住,白折腾半天。
# emqtt_bench 建立十万连接的示例(单台压测机) emqtt_bench conn -h 10.0.0.10 -p 1883 -c 100000 -i 10 emqtt_bench pub -h 10.0.0.10 -p 1883 -c 1000 -q 1 -t ptest/up -m "hello"5.2 系统内核参数不改,压测必翻车
压测前请先检查操作系统,百万连接对系统资源的占用远超常规web服务。下面这组参数在我压测时反复验证过,是基础中的基础。
ulimit -n 1048576 sysctl -w net.core.somaxconn=65535 sysctl -w net.ipv4.ip_local_port_range="1024 65535" sysctl -w net.ipv4.tcp_max_syn_backlog=65535 sysctl -w net.ipv4.tcp_fin_timeout=30 sysctl -w net.ipv4.tcp_tw_reuse=1这些参数的含义不用全都背下来,但至少要理解:文件描述符限制决定进程能开多少TCP连接;端口范围决定客户端机器能发起多少连接;backlog决定高并发握手时内核能排多长的队。改动完记得确认配置文件持久化,否则重启机器一夜回到解放前。
服务端也要调JVM和线程池参数。如果你用的是Java系网关,注意Netty的boss/worker线程数不必跟着核数堆,8核机器设成4个boss、8~16个worker通常就够。内存池用PooledByteBufAllocator,堆外内存要预留,否则高并发下GC会非常难看。日志级别压测时调到WARN,INFO级别的连接日志在百万连接场景下本身就是纯IO灾难。
5.3 压测过程中的典型问题和指标解释
我实测下来,8核心16G的机器跑EMQX,10万空连接很轻松,内存大概占用2~3GB。但连接从10万往20万走的时候,内存增速会明显变快,因为每个连接关联的Channel、消息缓冲、订阅结构都开始膨胀。这时如果消息吞吐再上来,GC时间会直线上升。先不要急着加机器,先查内存泄漏、查未释放的缓冲,很多团队在50万连接级别遇到的线上问题,不是Broker不行,而是没把单机调到位。
压测时一定要测两个维度:连接数、消息吞吐。连接数再高,如果消息吞吐上不去,也是死路。我们当时的经验值如下表,结合自己的机器可以做参考:
| 场景 | 单机资源 | 结果 |
|---|---|---|
| 10万空连接 | 8C16G | CPU 20%,内存 2.5G |
| 10万连接 + 每秒5千条消息 | 8C16G | CPU 65%,延迟P99 < 50ms |
| 集群4节点(40万连接) | 4台8C16G | 消息吞吐峰值每秒2万条,正常工作 |
| 集群10节点(百万连接) | 10台8C16G | 消息吞吐每秒8万条,P99延迟 < 100ms |
注意,百万连接在10台8C16G上是可以跑的,但前提是消息模型是典型IoT上报模型(设备上报到平台为主,下发比例很低)。如果你的场景是大量设备实时接收视频或高频控制指令,那消息吞吐会上升好几个数量级,节点数还要继续加。所以压测不能只背一个“百万连接”的目标就完事,把消息模型明确写下来,后面扩容才有依据。
6. 生产环境的几个深坑和补救
6.1 设备时间不准导致消息乱序
设备端的时间是最不可信的。很多设备没有RTC电池,时间全靠瞎跑,NTP又通不到公网,导致上报的timestamp远远偏离真实时间。业务系统一旦直接用设备时间做排序,就会出现消息在时间线上来回跳,甚至未来时间跑到前面去。
我的处理方法是双时间戳:平台网关在收到消息时打一个服务器接收时间(serverTime),设备上报数据保留设备时间(deviceTime)。时序存储用serverTime做主时间轴,deviceTime作为参考字段存起来。如果业务上需要精确到设备本地时间的事件顺序,那就让设备在消息里带一个单调递增的序列号,用序列号排序,不要用时间排序。这个经验几乎每个物联网项目都会遇到,千万别等到线上数据乱成一锅粥再补。
6.2 同一设备重复连接导致互相踢下线
弱网环境下,设备以为旧连接断了,发起新连接,但旧连接实际上还挂在一台网关节点上。于是同一个设备有两三条有效连接,平台可能把指令下发到旧连接上,设备却已经不听旧连接的指令了,指令就莫名其妙丢了。
这个问题要在网关层就解决。每个设备只能有一个“当前有效连接”,新连接认证通过后,网关通知旧连接的那个节点,要求它关闭这条旧的Channel。实现上可以在Redis里记录设备当前连接的唯一标识(连接ID),新连接注册成功后,旧的连接ID对应的Channel查出来主动关闭并清理订阅。要注意关闭旧连接和清理订阅的动作要全链路走完,尤其是离线事件不能重复生产,避免业务侧收到两次离线通知,逻辑乱套。
6.3 连接风暴与Broker实例重启的安全操作
生产环境不可能不重启Broker,但重启前必须摘流量。如果直接重启节点,连接在它下面的几十万设备会因为TCP断连而立刻发起重连,重连流量瞬间打到其他节点上,其他节点可能一片连锁崩溃。
我习惯的安全操作是:
- 先在负载均衡层面把该Broker节点摘掉,不让新连接进来。
- 观察存量连接平均多少秒后会自然掉线并重连到其他节点。
- 等连接数降到接近0再重启节点。
- 重启完先不挂回负载均衡,先健康检查,再逐步放流量。
如果业务要求不能等,那就必须确保所有节点都有足够余量承接其他节点的连接,并且连接认证限流全局生效。这套流程听起来简单,但每次线上演练我都发现不少团队根本没这意识,节点重启时直接硬来,结果就是平台雪崩。
6.4 监控和告警不能只看连接数
最后说监控。很多IoT平台大盘上只放了“当前连接数”,这远远不够。连接数只是表象,它没反映设备是否真的在业务上报。监控指标至少要包括:
- 连接数、连接成功率、连接失败原因分布。
- 上行消息TPS、下行指令TPS,按产品拆分。
- 消息转发延迟P99、端到端延迟P99。
- Kafka积压数量、消费组Lag。
- Redis命中率、时序库写入速率。
- 设备离线事件风暴、设备影子同步失败次数。
告警规则也不是设置一个固定阈值就完事,要结合历史趋势。比如“连接数低于过去24小时平均值的80%”往往比“连接数低于100万”更早发现问题。有一次我们业务量很高,连接数也正常,但Offline事件暴涨,查下来是某个型号的设备固件Bug导致反复断连。如果只监控连接数是看不出这个问题的,必须把事件速率纳入告警。
最后再说一点我自己这几年做IoT平台的体会:百万连接这个目标,真正考验的不是某个中间件能扛多少,而是你怎么把一个复杂系统切分成可独立扩展的模块,并且在每一层都提前想好流量失控、节点故障、设备异常这三个场景。架构上没有银弹,但只要你把连接模型、消息模型、设备模型梳理清楚,选型不迷信,压测不糊弄,百万设备连接是可以稳定运行的。设计一套能增长的IoT平台,远比拼出一套能演示的百万连接更有价值。