给内部项目起名这件事,我一直主张要有记忆点,但又别把功能写在名字上。上个月,我把跑了大半年的物联网接入平台正式定名为 Madeira。同事问为什么,我说你想想马德拉群岛在航海时代的位置,再想想马德拉酒为什么越陈越有味道。这个平台做的事是类似的:把乱七八糟的设备数据接进来,经过清洗、加工、按规则触发动作,最后沉淀成业务真正能用的资产。
项目缘起很朴素。团队要交付能耗监测、环境告警、设备预警几个业务线,后端要接的设备五花八门,有 Modbus 网关、蓝牙信标、MQTT 传感器、4G DTU,还有只会上报 HTTP 的定制设备。以前每个业务各拉一套接入逻辑,读写设备状态的重复代码写了好几遍,告警阈值散落在不同代码仓库里,想改一次阈值还要重新发布服务。那段时间最常听到的一句话是:“这个数据是哪个服务写的?口径怎么对不上?”
于是我们决定做一个偏底层的 IoT 基础平台,把设备接入、数据解析、时序存储、规则告警这几件通用的事统一收口。业务侧不再关心设备怎么传数据,只关心自己的业务逻辑和触发条件。这个平台就是 Madeira。
如果你正被异构设备接入、数据孤岛、告警规则写死在代码里这些问题折腾,这篇内容可以给你一套可参考的架构和实现思路。我会把设备接入层、数据模型、数据管道、规则引擎,以及上线后踩过的坑完整拆开讲,尽量说人话。
1. 平台定位与技术选型:为什么整套架构选择这三板斧
1.1 平台边界:只做通用能力,不碰业务逻辑
立项第一天,我拉着后端同事做了一次“痛点访谈”,把过去几个月重复造轮子的场景全部列了出来。问题高度集中在三块:
- 异构设备接入没有统一入口,每个业务都自己维护一套 MQTT 客户端或 HTTP 回调服务。
- 数据口径不统一,同一种属性,A 服务叫 temperature,B 服务叫 temp,单位还一会儿摄氏度一会儿华氏度。
- 告警规则硬编码在业务代码里,不支持组合条件,也没有时间窗口的概念。
这些问题单个拎出来都不难,但叠在一起就是典型的“底层重复劳动 + 数据事故隐患”。所以 Madeira 的定位很明确:它不关心里面跑的是能耗数据还是环境数据,只负责把设备上报的数据变成干净、有序、可触发的标准事件。业务逻辑由上层消费方自己决定,平台只提供能力。
这个边界很重要。如果平台企图把业务规则也收进来,就会变成一个巨无霸系统,改一行业务逻辑要拉一堆人评审。合理的方式是平台提供通用机制,业务方通过配置规则、订阅数据流来接入自己的场景。
1.2 技术选型拆解:MQTT、Kafka、TimescaleDB 各自解决什么问题
设备接入层首选 MQTT,这个基本没有悬念。MQTT 的发布订阅模型、QoS 等级、遗嘱消息、心跳保活机制,几乎就是为物联网设备端设计的。传感器、DTU、网关这些设备大多低功耗、网络不稳定,MQTT 的断线重连和消息可靠投递能力比裸 HTTP 长连接要省心太多。
那只会报 HTTP 的设备怎么办?我在接入层加了一个轻量代理服务,把 HTTP 请求转换成内部 MQTT 消息,后面的链路完全一致。Modbus 和串口设备则由边缘网关做协议转换,网络侧同样走 MQTT。这样无论设备原始协议是什么,进入到平台核心链路时已经统一成了一种消息格式。
消息管道选择 Kafka,主要解决两个问题:削峰和解耦。设备上报的数据量天然忽高忽低,晚间有批量任务时可能导致几百台设备同时上报,如果让存储层直接面对这种尖峰流量,很容易被打穿。Kafka 类似一个快递中转站:发件方把包裹扔进集散中心,收件方按自己的节奏去取,两边不需要同时在线上,也不用担心高峰期门口排队。
时序存储选了 TimescaleDB,而不是直接用 ClickHouse 或 InfluxDB。理由很简单:团队规模不大,TimescaleDB 是 PostgreSQL 扩展,SQL 语义完全兼容,运维成本低,中小规模的数据量下读写性能足够。ClickHouse 更适合大批量离线分析场景,后期如果业务分析需求变大,可以把数据分析链路单独接出去,但作为平台的主存储,TimescaleDB 够用且稳定。
1.3 模块划分与数据流向
整个平台拆成六个模块:接入层(device-ingress)、消息管道(Kafka)、数据处理服务(data-processor)、规则引擎(rule-engine)、时序存储(tsdb)、管理后台(console)。
数据流向是一条直线加一个分支:设备通过 MQTT 上报 JSON 数据,EMQX 按 Topic 桥接写入 Kafka;data-processor 消费 Kafka 消息,做解析、校验、补全、去重,然后把数据分成两路:一路写入时序数据库,一路推给规则引擎;规则引擎根据预设条件判断是否触发动作,动作包括告警入库、下发设备指令、回调业务系统 HTTP 接口。
这套拆分的好处是每个环节都能独立扩缩容。接入层入口被设备流量打爆了,只扩接入层;消费速度跟不上了,给>{ "msgId": "uuid:xxxx-xxxx-xxxx", "productKey": "pk_xxx", "deviceName": "dev_001", "ts": 1710000000000, "properties": { "temperature": 23.6, "humidity": 58.2 } }
msgId是全局唯一消息 ID,它是整个数据处理链路幂等的基础,后面排查重复数据全靠它。ts是设备产生的原始时间戳,注意这里只做存证,真正写入时序库的时间标准是平台接收时间serverTimestamp,原因下文细说。properties里的字段名和数据类型,必须跟管理后台定义的物模型一致,不一致的报文会在接入层被拦截。
2.2 物模型与设备影子:让每台设备都有“标准脸”
物模型是物联网平台里的老概念,但对一个统一接入平台来说,它是绝对的基础。我们给每一类设备定义一个 JSON Schema,描述三类能力:
- 属性(property):设备的状态值,比如温度、湿度、开关状态。
- 事件(event):设备主动上报的异常或提示,比如电量低、按键被按下。
- 服务(service):平台可下发的指令,比如重启、校时、调档。
为什么要做物模型?以温湿度为例,A 团队接入的传感器上报二进制温度,B 团队接入的设备上报字符串温度,如果平台不做统一抽象,下游每个服务都要自己写一堆转换逻辑。物模型负责把设备上报的原始数据翻译成标准字段,同时还能做合法性校验:温度超过了定义的取值范围,直接走异常分支,不污染主链路。
设备影子则是另一个很有用的设计。它的本质是云端缓存的一份设备最新状态快照。为什么要缓存?因为查询设备状态时,设备可能离线,也可能网络延迟很高,如果每次都发指令去实时读,会非常慢且不稳定。有了设备影子,查询接口直接读缓存,毫秒级返回。设备上线或上报新数据时,影子会同步更新。这个机制对管理后台的设备列表页、规则引擎的条件判断都很关键。
2.3 时序存储:建表、压缩与保留策略
设备属性数据是典型的时序数据,不能全堆在普通业务库里。普通 MySQL 表在几千万行之后,按时间范围查询的响应时间会明显变长,索引维护成本也高。TimescaleDB 的 hypertable 默认按时间分块(chunk),数据写入和查询都能自动路由到对应分块,效率高很多。
建表时要注意几个点。第一,表结构尽量扁平,不要把 properties 整个 JSON 塞进一个字段里,否则后续查询和聚合很难走索引。第二,标签字段(deviceName、productKey)和字段字段(温度值、湿度值)分开设计。第三,合理设置 chunk 时间间隔,太短会导致分块过多,太长则不利于旧数据淘汰。我们线上 50 万点级的规模,按天分块效果比较理想。
保留策略也不得不做。原始数据保留 90 天,超过部分通过 TimescaleDB 的连续聚合按小时降采样,再保留一年。日常实时查询走原始表,趋势分析走聚合表,两边互不干扰。这里有个经验:接入层做一次字段裁剪能省一半存储空间。很多设备上报的 JSON 里带了大量的冗余字段,比如固件版本、WiFi 信号强度、预留位等等,如果全部写入时序库,磁盘增长会非常快。业务需要哪些字段,在物模型里就定义好,接入口直接过滤掉无关字段。
3. 端到端实操:从设备上报到规则告警最小链路
3.1 核心链路与最小流程跑通
说完了设计,看看真正落地时最小链路怎么跑通。我按实际调试顺序来讲,这套顺序我们自己踩过一遍,比一上来就怼全链路要稳得多。
第一步,把 EMQX 跑起来,创建一个桥接,把mfd/#下的数据转发到 Kafka 的device_reportTopic。EMQX 的规则引擎配置大致长这样(版本不同会略有差异):
bridges.kafka { type = kafka servers = "127.0.0.1:9092" topic = "device_report" }第二步,写>func handleReport(msg []byte) error { var report DeviceReport if err := json.Unmarshal(msg, &report); err != nil { writeDeadLetter(msg) return err } report.ServerTs = time.Now().UnixMilli() if deduplicated(report.MsgId) { return nil } if err := validateByModel(report); err != nil { writeDeadLetter(msg) return err } tsdb.Insert(report) ruleEngine.Fire(report) return nil }
这里插一句:writeDeadLetter是整套链路里最容易被忽略但又最重要的功能。解析失败、校验失败的数据不能直接丢,要写到死信队列里,方便事后排查是设备端 bug 还是模型配置问题。没有死信队列,你会经常面临“用户说数据丢了,但业务日志里什么都没有”的窘境。
第三步,把规则引擎单独做成一个服务而不是函数库,是为了让它可以独立扩容,也方便后续接入多种触发源。规则引擎消费的是一条内部的device_processedTopic,里面是已经清洗过的标准事件。业务方不用关心原始数据从哪来,只订阅清洗后的结果。
3.2 实站实现:规则引擎的 JSON DSL 怎么写
规则引擎最怕两件事:一是规则表达能力太弱,稍微复杂一点的逻辑就要写代码;二是规则表达太灵活,随便一个脚本引擎都能执行,结果安全和维护成本完全失控。
我用的是 JSON DSL 加轻量表达式引擎,规则结构分三部分:基本信息、匹配条件、执行动作。下面是一个“机房高温告警”的规则示例:
{ "id": "rule_tmp_01", "name": "机房高温告警", "match": { "productKey": "pk_xxx", "deviceType": "temperature-sensor", "expr": "properties.temperature > 60" }, "actions": [ { "type": "alert", "level": "warning", "target": "ops", "template": "设备 {deviceName} 温度 {temperature} 度,超过 60 度" }, { "type": "command", "topic": "mfd/{productKey}/{deviceName}/service/call", "payload": { "action": "open_fan" } } ] }条件里的expr用 Aviator 之类的轻量表达式引擎执行,只允许访问当前事件里白名单字段,不允许任意代码执行。这样既保证了表达灵活性,又不会引入高危的脚本执行能力。规则支持与、或、比较、算术运算,对绝大多数业务场景都够了。
真正复杂的是“时间窗口”类规则。比如“五分钟内温度连续超过 60 度才告警”,这不是简单的单次命中,需要记录滑动窗口内的连续状态。我们实现方式是引入一个状态计数器,按“设备维度 + 规则 ID”存 Redis,每次数据进来时更新,窗口内满足条件才触发,窗口内复位则清空计数。这个模式下,规则引擎天然就是一个状态流处理节点,所以它的存储依赖必须独立设计,不能把状态放到进程内存里,否则重启一次规则执行就乱了。
3.3 数据管道中的三个硬骨头:幂等、时间戳、背压
链路跑通之后,大头工作才开始。按我实际体验,真正决定线上稳定性的不是架构,而是三个细节。
第一是消息幂等。MQTT 的 QoS1 语义是“至少一次”,这意味着网络抖动时消息可能重复到达。Kafka 本身也提供了 at-least-once 保证。两层叠加,消费端必须在业务上做幂等。我的做法是用 Redis 对msgId做去重,设置合理的过期时间(比如 30 分钟),重复的msgId直接丢弃。不处理的后果很直接:重复数据重复入库,规则触发重复告警。
第二是时间戳统一。设备上报的ts来自设备本地时钟,很多设备没有 NTP 校时,时间可能差出几分钟甚至更多。如果按设备时间做时序查询,你会发现数据曲线时不时“倒挂”。我的做法是:写入时序库的排序时间一律用serverTimestamp,也就是平台接收时间,设备上报的ts作为原始字段单独存储,只在排查设备端问题时才看。这样才能保证数据的可靠性。
第三是背压。Kafka 消费慢,积压会越来越严重。消费速度由两个因素决定:分区数和消费逻辑的耗时。分区数一定要大于等于消费者实例数,否则多出的实例闲着没事干;同时,消费逻辑里不能有同步的远程调用,比如每条消息都去调一次告警 HTTP 接口,这种必须异步化,或者批量处理后统一发送。我们后来把告警发送改成异步队列,积压问题立刻缓解。
4. 踩坑实录:设备、规则、数据三个方向的高频问题
4.1 设备显示在线但消息断了,怎么查
上线初期收到最多的反馈是“后台显示设备在线,但数据不动了”。这个问题的根源在于“在线”的定义不一致。设备主动连接 EMQX 时,连接层当然知道它在线;但设备可能因为网络问题已经与服务器断开,只是没有及时发送遗嘱消息,或者连接的 session 还停留在过期边缘。
排查的时候第一件事是看 EMQX 的 dashboard,确认这个设备当前是否真的有活动连接。如果 EMQX 显示已断开,但管理后台还显示在线,问题多半出在状态同步逻辑上。我们的解决办法是:设备离线状态不完全依赖 MQTT 连接,而是结合心跳上报超时来判定。管理后台的“在线”定义就是“最近 N 分钟内有过属性上报”,这样即使连接层状态有误差,业务侧看到的数据仍然是可靠的。
另一个相关坑是网关转发。部分设备先连到边缘网关,再由网关转发到平台,设备与网关之间用的是私有协议。这种情况平台看到的“设备在线”其实只是“网关在线”,设备本身可能早已离线。排查时一定要查看网关的上行缓存日志,确认数据是从设备实时产生还是从缓存补发的。
4.2 规则命中了却没告警,问题出在哪
规则引擎上线之后,最让人抓狂的问题是“这条规则的测试数据明明触发了条件,为什么告警没发出来?”排查了几轮之后发现原因主要集中在三处。
第一是类型比较问题。设备上报的温度是数字 60,规则表达式里写的是字符串"60",Aviator 在严格模式的比较结果可能完全不符合预期。所以在上报数据进入规则引擎之前,data-processor 会做一次类型规范化,把物模型里定义的数值类型强制转成 float,确保到规则层的时候类型已经是可控的。
第二是规则引擎订阅的数据源不对。有些规则配置的是监听某个 productKey,但业务方实际测试时用了另一个 productKey 的设备上报,这类配置错误在日志里非常隐蔽。排查方式很简单:在规则引擎里加一个匹配日志,记录每条进入规则上下文的设备标识和命中情况,一眼就能看出数据流有没有走对。
第三是动作执行失败但没留下痕迹。告警发送到钉钉、企业微信这类回调接口,失败原因可能是回调地址变化、接口超时、模板变量解析异常。我们的做法是动作执行也走一张执行记录表,每一条规则命中的动作状态都能查到,失败原因写在日志里。没有这张表,告警缺失几乎只能靠猜。
4.3 Kafka 积压和时序库暴涨的处理思路
Kafka 积压是一个可见又可查的问题。监控告警发现device_report消费积压到几百万条时,第一反应不是加消费者,而是先判断消费链路里有没有慢操作。我们遇到过两次积压,原因都很典型:一次是规则引擎初始化时连接 Redis 超时,重试机制写成了同步阻塞,整个消费线程被卡住;另一次是数据量突增后,批量写 TimescaleDB 的 SQL 没加分批,单次写入量太大导致数据库阻塞。
解决积压的基本原则是:先恢复消费速度,再查根因。如果消费程序还活着,先把规则引擎的动态开关打开,临时跳过不重要的规则类型,只做数据入库,让积压尽快降下来。等积压清零后,再回头解决慢操作。时序库暴涨的问题则需要两个手段配合:接入口的字段裁剪和数据库层的保留策略,缺一个都扛不住长期数据增长。
4.4 常见问题速查表
我把上线以来被问得最多的几个问题整理成了一张表,方便你直接对号入座。
| 现象 | 可能原因 | 排查手段 | 解决方案 |
|---|---|---|---|
| 设备显示在线但无数据 | 设备与网关断连,网关缓存转发 | 看 EMQX 连接、查网关日志 | 在线状态按心跳上报判定 |
| 重复数据入库 | QoS1 重复投递,消费端未幂等 | 检查 msgId 去重记录 | Redis 按 msgId 去重 |
| 规则不触发 | 类型不匹配、订阅配置错误 | 查看规则命中日志 | 数据入规则前做类型标准化 |
| 规则触发但告警丢失 | 动作执行失败、回调异常 | 查动作执行记录 | 记录动作状态,失败重试 |
| Kafka 积压持续上涨 | 消费者阻塞、分区数不足 | 看消费 lag 和线程状态 | 异步化慢操作,合理扩分区 |
| 时序库磁盘暴涨 | 冗余字段过多、保留策略缺失 | 查 hypertable chunk 大小 | 接入口做字段裁剪,配置降采样 |
5. 复盘心得与后续演进方向
5.1 回头看,哪些设计真正省了事
平台上线三个月后再回头看,有几个设计我认为是真正省了事的。
一是统一物模型。虽然前期给每类设备建模有点繁琐,但后期接入新设备时,只需要定义一套模型,后续所有解析、校验、规则、存储全部自动生效。新设备接入时间从以前的一周缩短到一天,大部分时间都花在跟设备厂商确认字段定义上。
二是死信队列。这个设计在前期差点被砍掉,觉得“解析失败的数据直接丢就好了,搞什么死信队列”。后来几次线上问题都靠它定位到了根因:有的是设备固件字段名变了,有的是单位从摄氏度变成了开尔文,如果没有死信数据,这些问题根本无从查起。
三是规则引擎的独立部署。把规则引擎从>