自建物联网接入平台:从MQTT到规则引擎的架构实战
给内部项目起名这件事我一直主张要有记忆点但又别把功能写在名字上。上个月我把跑了大半年的物联网接入平台正式定名为 Madeira。同事问为什么我说你想想马德拉群岛在航海时代的位置再想想马德拉酒为什么越陈越有味道。这个平台做的事是类似的把乱七八糟的设备数据接进来经过清洗、加工、按规则触发动作最后沉淀成业务真正能用的资产。项目缘起很朴素。团队要交付能耗监测、环境告警、设备预警几个业务线后端要接的设备五花八门有 Modbus 网关、蓝牙信标、MQTT 传感器、4G DTU还有只会上报 HTTP 的定制设备。以前每个业务各拉一套接入逻辑读写设备状态的重复代码写了好几遍告警阈值散落在不同代码仓库里想改一次阈值还要重新发布服务。那段时间最常听到的一句话是“这个数据是哪个服务写的口径怎么对不上”于是我们决定做一个偏底层的 IoT 基础平台把设备接入、数据解析、时序存储、规则告警这几件通用的事统一收口。业务侧不再关心设备怎么传数据只关心自己的业务逻辑和触发条件。这个平台就是 Madeira。如果你正被异构设备接入、数据孤岛、告警规则写死在代码里这些问题折腾这篇内容可以给你一套可参考的架构和实现思路。我会把设备接入层、数据模型、数据管道、规则引擎以及上线后踩过的坑完整拆开讲尽量说人话。1. 平台定位与技术选型为什么整套架构选择这三板斧1.1 平台边界只做通用能力不碰业务逻辑立项第一天我拉着后端同事做了一次“痛点访谈”把过去几个月重复造轮子的场景全部列了出来。问题高度集中在三块异构设备接入没有统一入口每个业务都自己维护一套 MQTT 客户端或 HTTP 回调服务。数据口径不统一同一种属性A 服务叫 temperatureB 服务叫 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 桥接写入 Kafkadata-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规则表达式里写的是字符串60Aviator 在严格模式的比较结果可能完全不符合预期。所以在上报数据进入规则引擎之前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 回头看哪些设计真正省了事平台上线三个月后再回头看有几个设计我认为是真正省了事的。一是统一物模型。虽然前期给每类设备建模有点繁琐但后期接入新设备时只需要定义一套模型后续所有解析、校验、规则、存储全部自动生效。新设备接入时间从以前的一周缩短到一天大部分时间都花在跟设备厂商确认字段定义上。二是死信队列。这个设计在前期差点被砍掉觉得“解析失败的数据直接丢就好了搞什么死信队列”。后来几次线上问题都靠它定位到了根因有的是设备固件字段名变了有的是单位从摄氏度变成了开尔文如果没有死信数据这些问题根本无从查起。三是规则引擎的独立部署。把规则引擎从>
上一篇/下一篇内容由系统自动关联
返回资讯列表 →