1. 接入层整体设计
做工业数据接入这件事,最头疼的往往不是某个协议有多难,而是设备太多、协议太杂。今天聊的这套方案,核心就是把 MQTT、Modbus 这两类最常见的工业协议,全部收口到 DolphinDB 里,统一成一测点一行的“测点流”,让后续的规则引擎、告警、报表、AI 训练都能直接对着同一张宽表干活。
先说为什么非要“统一”。我在实战里见过很多团队,设备数据进来的时候是什么协议,落库就按什么协议建表。MQTT 一份表、Modbus 一份表,看起来省事,但下游用的时候全是坑。时序关联要跨表 join,点位名称在两张表里对不上,单位还各写各的,每次出报表都要人工核对半天。与其让下游替你做脏活,不如在接入层就把数据规整成统一格式。
这套设计里,我推荐一个“三明治”结构:
- 底层是协议适配:MQTT 插件负责订阅、Modbus 插件负责轮询,它们只做一件事,就是把报文变成结构化记录。
- 中间是统一测点流:一张物理宽表,测点ID、设备ID、时间戳、值、质量戳、采集时间,六列打天下。
- 上层是应用层:告警、看板、机器学习,全部只认测点流这一张表。
这个思路的好处在于,以后再加新协议,比如 OPC UA、IEC 104,你只需要在底层新增一个适配器,中间层和上层完全不用动。这就是“统一测点流”的真正价值——它是整个数据链路的稳定接口。
提示:不要在底层做业务判断,过滤逻辑一律放到测点流之后。协议层只保证“收得到、解析对、落得下”,否则接入层会变得越来越重,最后变成没人敢碰的屎山。
2. MQTT 接入实战
2.1 MQTT 插件配置与订阅关系
DolphinDB 的 MQTT 插件,本质是一个订阅客户端。你需要告诉它三件事:Broker 在哪、订阅什么主题、消息来了如何处理。
先说 Broker 配置。生产环境通常不会只有一台,我建议在插件配置里就写多个地址,用失败转移。插件里直接配连接参数即可,关键是host和port必须能连通,这个听起来是废话,但我见过太多人换了服务器之后,忘了改防火墙规则,客户端一直报连接超时,折腾了半天才发现是安全组没放行 1883 端口。
订阅关系我用一段伪配置说明:
mqttClientId=iot_gateway_01 brokerAddress=192.168.1.101:1883,192.168.1.102:1883 subTopic=devices/{deviceId}/data qos=1这里有一个很重要的点:subTopic不建议固定写死一个完整主题,比如devices/gw01/data,而是建议用通配符,比如devices/+/data。这样新接入设备时不用改配置,插件能自动订阅到新设备的主题。实际项目中,我曾经因为主题写死,导致客户加了一台设备后数据一直收不到,排查了很久才发现是订阅列表没更新。
2.2 消息解析的两种模式
MQTT 消息的 payload 有几种常见格式,最常见的是 JSON 和二进制。DolphinDB 插件支持自定义解析函数,你可以写 DolphinDB 脚本,也可以直接调用内置 JSON 解析函数。
JSON 模式比较简单。比如一条设备上报数据是这样的:
{"deviceId":"dev_001","ts":"2024-05-20 10:00:00.000","values":{"temp":25.6,"humidity":60.2}}你只需要在订阅回调里写一个函数,把 JSON 拆开:
def mqttMsgHandler(msg){ parseJson = jsonParse(msg.payload) deviceId = parseJson.deviceId ts = timestamp(parseJson.ts) temp = parseJson.values.temp humidity = parseJson.values.humidity // 写入测点流 insert into loadTable("dfs://iot", "metric_stream") values(ts, deviceId, "temp", temp, 0, now()) insert into loadTable("dfs://iot", "metric_stream") values(ts, deviceId, "humidity", humidity, 0, now()) }这里我建议你注意性能。如果单条消息里带了几十个测点,一个测点一个 insert 显然不现实。正确做法是先把数据攒到一个表变量里,用tableInsert或者append!批量写入。我在现场测过,批量写入比单条写入快了两个数量级,毫秒级和秒级的差别。
二进制模式稍微麻烦一点。有些设备为了省流量,用自定义二进制格式上报,比如报文头 4 字节是设备 ID,接下来 4 字节是时间戳,再往后每 2 字节是一个测点值。这时候就需要写一个带位运算的解析函数。说句实在话,如果设备厂商能提供报文文档,二进制解析并不难;难的是厂商自己也说不清格式。我的经验是拿真实报文逐字节对照文档,用 DolphinDB 的脚本调试窗口一边看一边试,效率最高。
2.3 QoS 等级的选择
MQTT 的 QoS 有三个等级,很多初学者搞不清该用哪个,我把话说透:
- QoS 0:最多一次,消息可能丢,适合温度、湿度这种丢了也无所谓的遥测数据。
- QoS 1:至少一次,消息不丢但可能重复,适合大部分监控场景。
- QoS 2:恰好一次,性能开销最大,只适合计费、订单这类绝对不允许丢不允许重的数据。
工业采集场景里,测点数据用 QoS 1 就够了。重复消息虽然会带来少量重复数据,但我们在测点流设计里会加一个去重逻辑,后面细说。QoS 2 在工业遥测里性价比很低,我基本不推荐。
2.4 断线重连与消息堆积
MQTT 客户端断线是常态,不是异常。你需要在插件配置里设置keepAliveInterval,建议 30 到 60 秒,太短会给 Broker 增加无谓的负载,太长会导致断线发现不及时。还要设置autoReconnect=true和cleanSession=false,这样会话恢复时能拿到离线期间积压的消息。
这里有个坑,就是当设备端量大、网关多的时候,一旦网络抖动,所有客户端同时重连,Broker 会被打爆。所以重连逻辑一定要加退避策略,不要一断就连。现场实测过,30 台网关同时重连,那台单机版 EMQX 直接 CPU 100%,消息全部堆积。
注意:DolphinDB 侧写数据时如果遇到 Broker 推送大量积压消息,写入速度会跟不上。这时候监控写入延时非常重要,一旦发现订阅线程堆积,立刻处理。
3. Modbus 接入实战
3.1 Modbus RTU 与 Modbus TCP 的选择
Modbus 这个东西,做工业的都知道,老古董但生命力极强。分 RTU(串口)和 TCP(以太网)两种。RTU 走 RS485,一主多从,波特率从 9600 到 115200 不等;TCP 直接走以太网,默认端口 502。
接入 DolphinDB 的时候,我的建议是优先选 Modbus TCP。原因很简单,如果现场已经有协议网关把 RTU 转成了 TCP,你就不用关心串口参数了,直接拿 IP、端口、寄存器地址就能干活。如果只能走 RTU,那你需要先确认波特率、数据位、校验位、停止位,这些参数错了任何一个,报文都是错的。
DolphinDB 的 Modbus 插件做轮询采集,你要定义一个轮询任务,每个任务关联一个设备的一张寄存器映射表。某个设备要采哪些测点、每个测点对应哪个寄存器、数据类型是什么、字节序是什么,都在映射表里写清楚。
3.2 寄存器映射表的设计
Modbus 的寄存器是分区的:线圈(Coil)、离散输入(Discrete Input)、保持寄存器(Holding Register)、输入寄存器(Input Register)。我们实际采集的数据,绝大多数在保持寄存器和输入寄存器里。
映射表里最关键的三列:register_start、quantity、data_type。
register_start:从哪个寄存器开始读。注意,Modbus 协议里地址有 0 和 1 的差异,不同厂商的说明书写法不一样,有的从 0 开始,有的从 1 开始。插件读取时要统一用 0 基地址进行配置。quantity:连续读多少个寄存器。比如一个 32 位浮点数在 Modbus 里占 2 个 16 位寄存器,你就得连续读 2 个。data_type:数据类型,常见有INT16、UINT16、INT32、FLOAT32、FLOAT64。
这里最坑的是字节序。同样一个 32 位浮点数,有的设备是大端,有的是小端,还有一种是字序大端但字节序小端,四种组合你都得试。之前接一台老外的设备,怎么读都是 0,后来把字节序改成 Big-Endian,数据立刻正常。如果你不确定字节序,先在 DolphinDB 里手动读几个寄存器看一下原始 hex 值,再按实际格式定义解析规则。
3.3 轮询周期的设置
Modbus 是主从协议,从机不会主动上报,全靠主机轮询。轮询周期的设计直接影响数据时效性和总线负载。
我给的参考值:
| 场景 | 轮询周期 | 说明 |
|---|---|---|
| 设备数 10 台以下 | 500ms ~ 1s | 响应快,实时性好 |
| 设备数 10 ~ 50 台 | 1s ~ 3s | 平衡时效性与负载 |
| 设备数 50 台以上 | 5s ~ 10s | 防止总线拥堵,优先保证任务成功率 |
如果有个别测点需要毫秒级响应,建议单独拉一条专线,用独立网关来做,别跟大轮询混在一起。我有一次把一台高速设备的振动信号混在 PLC 轮询里,结果轮询一圈回来,振动数据已经被业务侧嫌弃延迟太大。
3.4 Modbus 轮询脚本实现
在 DolphinDB 里,Modbus 插件的核心是一个封装函数,循环遍历设备列表,逐设备执行读操作。读回来的是寄存器原始值数组,你要再根据映射表做一次换算,比如原始值乘以 0.1 才是真实温度,再写入测点流。
示例逻辑:
for device in deviceList { // 读保持寄存器 rawData = modbusRead(device.ip, device.port, device.slaveId, "HOLDING", device.registerStart, device.quantity) // 按映射表转换 for i in 0..(mapping.size()-1) { metric = convertValue(rawData, mapping[i].dataType, mapping[i].byteOrder) insert into loadTable("dfs://iot", "metric_stream") values(now(), device.deviceId, mapping[i].metricName, metric, 0, now()) } }这段思想很简单,但工程上要处理几个细节:
- 每次轮询结束之后,记录一下耗时,如果某次轮询超过了周期的 1.5 倍,就要告警。
- 如果某台设备连读三次失败,把该设备标记为离线,不要反复重试同一台设备,否则会拖垮整条总线的轮询节奏。
- 写入测点流之前,要把原始值换算成物理值。换算规则是写在映射表里的,不要写死在代码里,这样加设备的时候不用改一行代码。
3.5 RTU 串口的 485 总线注意事项
如果你还是绕不开 RTU,那我说几个 RS485 总线的坑:
- 终端电阻。总线上最远的两端必须各接一个 120 欧姆终端电阻,否则长线传输时信号反射会导致数据错乱。我见过现场调试的人,怎么调都不对,最后发现总线两端根本没接电阻。
- 手拉手接线。RS485 是总线拓扑,必须串联手拉手,不能星形连接。星形接法在短距离可能没问题,一旦距离超过几十米,反射就让通讯废了。
- 地线。RS485 的 A/B 两根信号线之外,最好再拉一根公共地线。因为工业现场干扰多,共地能显著降低共模干扰导致的误码。
提醒:Modbus 轮询脚本别写得过于“健谈”。两台设备之间的轮询间隔建议给一点缓冲,尤其在 WiFi 网桥连接的场景,连续高频发包很容易把网桥打挂。
4. 统一测点流的设计与落地
4.1 表结构设计
统一测点流这张表,是我这套方案的心脏。表结构如下:
| 字段名 | 类型 | 说明 |
|---|---|---|
| ts | TIMESTAMP | 数据采集点时间,准实时 |
| device_id | SYMBOL | 设备编号,全局唯一 |
| metric_id | SYMBOL | 测点编号,全局唯一 |
| metric_value | DOUBLE | 测点值 |
| quality | INT | 质量戳,0 正常,1 可疑,2 失效 |
| ingest_time | TIMESTAMP | 写入 DolphinDB 的时间 |
这里有一个设计细节想特别说明一下:为什么要device_id + metric_id,而不是一个联合 ID?因为从 MQTT、Modbus 来的数据,设备 ID 和测点名本来就分属不同字段,保留两个字段可以在上层应用做分组统计的时候省下大量字符串拼接。
ingest_time这个字段容易被人忽略,但我认为是这套表设计里最不该省的一列。它可以用来衡量采集链路的总延迟,也方便排查“采集器时间错了导致时序倒挂”这类问题。我见过太多团队因为不存摄入时间,遇到数据乱序时完全无从下手。
4.2 宽表与窄表的取舍
很多人说,既然叫“测点流”,为什么不设计成宽表,一列一个测点,一行一个设备?宽表查询确实直观,但工业采集的痛点在于设备型号不同、测点组合不同,宽表要频繁加列,运维成本极高。
我推荐的是“窄表存储 + 宽表查询”的组合。磁盘上存窄表,上层用pivot函数把需要的测点旋转成宽表来查看或计算。你不需要在下层就迁就上层某一个固定视图,因为视图是随时可以变的,但表结构一旦定下来,迁移成本极高。
4.3 写入性能优化
DolphinDB 的写入性能瓶颈,几乎永远不在磁盘,而在你的“写入方式”。
四个经验:
- 批量写入是必须的。单条写入在吞吐量上毫无竞争力。把 1000 条数据攒成一个块,一次写入,吞吐量能提升一个数量级。
- 分区粒度别太细。工业数据按天分区足够,如果按小时分区,每天的写入会触发大量小分区合并,反而拖慢查询。
- 排序键。
ts和device_id组合设置排序键,查询单个设备一段时间的数据会非常快。 - 副本策略。数据重要就设 2 副本,但在集群规模小于 3 台时,2 副本会拖慢写入性能。小集群先保写入,别盲目追求冗余。
4.4 数据去重与质量标记
MQTT QoS 1 会带来重复消息,Modbus 轮询偶尔读到瞬时错误值,这两类问题在统一测点流里用同一个办法兜底:基于ts + device_id + metric_id做去重。DolphinDB 的流计算引擎里,distinct窗口能解决绝大部分重复数据。
对于瞬时错误值,比如寄存器读到 0xFFFF,换算出来是 -9999 这种奇葩值,我的做法是给它标记quality=2,而不是直接删除。因为告警系统可能需要识别“数据失效”这一个状态,把它删了反而是掩耳盗铃。
4.5 多协议数据规整为统一测点流的处理流程
从原始协议数据到统一测点流,中间经过三层处理:
- 协议解析层:把 MQTT 消息、Modbus 报文解析成结构化记录。
- 字段映射层:把设备ID、测点名、时间戳、值,从协议自定义格式映射到统一字段。
- 质量控制层:做去重、做边界检查、做物理单位换算。
这个三层处理看起来简单,但每层都要做好“防御性设计”。比如第二层映射,如果某个新设备的 JSON 里deviceId字段叫dev不叫deviceId,映射层要做别名兼容,而不是改底层代码。第三层,如果温度传感器坏了读出来 85 度,但现场设备额定温度最高 60 度,质量层要能识别并标记异常。这些逻辑越靠前做,上层应用越干净。
5. 常见问题与排查技巧实录
5.1 MQTT 订阅了但收不到任何消息
先确认三件事:主题是否正确、Broker 是否允许该客户端订阅、消息是否真的发出来了。很多人第一步就错在 Topic 用$SYS开头的内容去测,$SYS主题本来就特殊,不是普通客户端能随便用的。
实测排查流程:
- 先用
mosquitto_sub或者 MQTT Explorer 订阅同一个主题,确认 Broker 上有消息在流动。 - 再检查 DolphinDB 插件里的
subTopic是否带了通配符,如果 Broker 端消息主题是devices/gw02/data,而你订阅devices/gw01/data,自然收不到。 - 查看插件日志,确认订阅成功之后有没有消费消息。很多情况下,订阅是成功的,但消息 handler 里解析报错,插件会将异常消息丢弃。
5.2 Modbus 读上来的值对不上
遇到读回来的值跟设备面板显示不一致,99% 是字节序、数据类型或者缩放因子的问题。
举例:某设备用 Modbus 保持寄存器存温度,寄存器 40001(协议地址 0)是浮点数,我一开始按FLOAT32_BE解析,读出来 4.17,但面板显示 41.7。后来查文档发现数据的缩放因子是 10,也就是说原始值是 417,面板除 10 显示 41.7。这种问题代码里加一个scale字段,定义映射表时注明缩放因子即可。
还有一种常见情况是寄存器地址偏移。很多厂商文档里的地址是 PLC 地址(40001 这种),而实际协议帧里的地址是 0,所以配置采集任务时如果不做偏移换算,读出来的值永远是另一个寄存器的。
5.3 轮询任务跑着跑着停了
Modbus 轮询任务突然停了,最常见的原因是脚本抛了异常,而你没有在任务外层包try-catch。比如某次读到非法响应时,整个轮询任务就崩了。我的习惯做法是:每个设备、每一轮读操作都包try-catch,捕获异常后记录日志,继续处理下一台设备。
另外还要注意设备响应超时。Modbus TCP 如果设备 IP 不可达,默认连接超时可能要几十秒,这期间轮询线程卡住,后续所有设备都被堵住了。所以一定要设置较短的超时时间,比如 1500ms,宁可读失败标记离线,也不能让整个轮询队列被一台死设备卡死。
5.4 MQTT 消息积压严重
这个问题通常发生在网关端。如果你发现 DolphinDB 写入速度正常,但 MQTT 客户端消费消息的延迟越来越大,那问题多半出在消息处理函数里做了重活,比如在 handler 里直接写库、或者做复杂 JSON 解析。
正确做法:handler 只做最基本的解析,然后把明细数据放进队列,由另一个线程批量落库。DolphinDB 的流表引擎天然适合这个场景。订阅消息后先写入流表,再用流计算引擎批量写入分布式表,让“接收”“解析”“落库”三者解耦。
5.5 时间戳不一致引起的时序错乱
MQTT 消息里的ts字段是设备端生成的,设备端时钟可能不准,导致写进来的数据时间顺序是乱的。我推荐两个字段都要存:设备时间ts和摄入时间ingest_time。做实时计算时用ts,做写入监控和排障时用ingest_time。
如果现场设备整体时钟偏差超过 5 秒,问清楚再决定是让设备端校时还是接入层做时间偏移补偿。有些 NTP 没配好的现场,设备时间比服务器快半小时,不处理的话所有趋势图都是斜的。
6. 实测体会与扩展思路
这套多协议接入方案在项目里跑了大半年,给我最大的感受是:真正的复杂度不在协议本身,而在你用什么结构去承接这些数据。MQTT 接入的排查成本主要在主题管理和消息格式兼容;Modbus 的坑主要集中在寄存器地址、字节序和轮询节奏;而统一测点流的设计,则是决定整个数据平台能否长期好用、不用返工的关键。
最后再分享一个扩展思路。如果后续要接 OPC UA、IEC 104,或者走边缘网关统一转换后再推 MQTT,这套测点流结构完全不需要变化,只是增加适配器而已。甚至你可以把测点流的最底层设计成“度量 + 标签 + 时间戳 + 值”的模式,未来对接 Prometheus 这类生态也能顺带打通。按这个节奏走下去,接入层会越来越薄、越来越稳定,你的精力就可以放到真正有价值的事情上:用数据做告警、做预测、做产线优化。
我个人在实际项目里吃过亏的地方也多说一句:一定要在上线第一天就监控“消息积压量”“轮询失败率”和“写入吞吐”这三个指标,别等出了事故再回头查。先把这层防御打好,多协议接入这件事就算成了大半。