干工业数采的都知道,设备接口最能磨人的不是采集本身,而是多协议共存后的数据归一。同一个车间里,智能电表走 Modbus RTU 挂在 485 总线上,光伏逆变器走 Modbus TCP,新加的环境传感器又都走 MQTT、通过网关汇聚到消息服务器。这些数据如果各存各的,后面做可视化、做报警、做分析的时候,光对表就能对到怀疑人生;更别说 DolphinDB 这种时序数据库,数据进库之前没有一个统一的测点模型,分区、查询、流计算全都施展不开。
我今天要复盘的就是这件事:把 MQTT 和 Modbus 这两类最常见的采集协议,完整接入 DolphinDB,最终形成一张统一的测点流表。这个项目解决的是工业物联网里最普遍的问题——多协议数据归一化。如果你正在做工业数采平台、能耗监测系统,或者想把 DolphinDB 用到生产环境,这篇内容应该能帮你少走不少弯路。
1. 多协议接入的整体思路:统一测点流才是核心
1.1 多协议并存的真实场景与痛点
在工业现场,设备接口五花八门是常态。一个中等规模的园区或者工厂,设备清单往往长这样:
- 配电柜里的智能电表,本质上是 Modbus RTU 设备,挂在 RS485 总线上,波特率 9600、8-N-1,几十块表串在一起。
- 光伏逆变器、空调主机、空压机控制器,这类设备一般带网口,走 Modbus TCP,端口 502。
- 这两年新装的温湿度传感器、烟感、水浸探测器,基本都走 MQTT,通过一个物联网网关统一接入到自建的 EMQX 或者云上的 MQTT Broker。
问题随之而来:不同协议的设备,数据格式完全不一样。Modbus 那边读出来的是原始寄存器数值,可能是一串 16 位整数,需要你根据设备手册换算成真实的电压、电流、功率;MQTT 那边消息体里全是 JSON,不同厂家还喜欢用不同的字段名,这家叫humidity,那家叫rh,时间戳有的发北京时间字符串,有的发 UTC 毫秒。
如果每一个协议各建一套存储逻辑,短期内好像没啥,时间一长就崩了。命名混乱、时区不统一、设备 ID 对不上,做报表的时候要写一大堆 if-else 去适配历史数据。我在项目里吃过这个亏,后来痛定思痛,所有采集数据必须先落到统一的测点模型上,再谈存储和分析。
1.2 统一测点流的数据模型与三层结构
所谓统一测点流,其实就是把千奇百怪的设备数据,最终收敛成一个标准四元组:
(时间戳,设备 ID,测点名,测点值)
额外再加一个质量码字段,用来标记数据是否正常。设备 ID 全局唯一,比如PLANT1-MTR-001;测点名统一用小写加下划线,比如voltage_rms、active_power、temperature;值是浮点数,别用 int、字符串混着来,后面做窗口计算会方便很多。
架构上,我习惯分成三层:
- 采集接入层:负责和各协议打交道。MQTT 这边是订阅消息,Modbus 这边是轮询设备,这一层只负责把原始数据变成标准四元组。
- 汇聚处理层:在 DolphinDB 里,用一张流表接收所有接入层推送的数据,统一去重、补时间戳、映射设备 ID。
- 存储计算层:流表持久化,加上流计算订阅,把原始测点数据变成分钟均值、报警事件等次级结果。
这样分层以后,新增一种设备协议,只需要在接入层写一个适配器,上面的存储和计算完全不用动。这是统一测点流最大的价值——接入成本的边际递减。
1.3 为什么选 DolphinDB 做汇聚和存储
可能有人会问,Kafka 加关系型数据库不是也这么干?确实可以,但在这个场景里,DolphinDB 有几个很实在的优势:
- 原生支持流表、表共享和持久化、流计算订阅,相当于把消息队列、实时计算引擎和时序数据库三个组件合成一个,部署和运维成本明显低。
- 时序场景写入和聚合性能强,尤其是按时间分区的设计,压缩比也不错。
- 内置丰富的时间序列聚合函数,比如滑动窗口、重采样、ffill/bfill,统计数据直接一条 SQL 搞定。
- 跟 Python API 衔接很好,采集层用 Python 写适配器,推数据到 DolphinDB 非常顺滑。
当然 DolphinDB 也有学习曲线,但一旦表格建好、订阅挂好,整个数据链路非常稳定。下面我把两条接入路径分别拆开讲,先从 MQTT 开始。
2. MQTT 接入实战:从 Broker 到 DolphinDB 流表
2.1 MQTT 的关键概念与主题规划
MQTT 是发布-订阅模型,Broker 是核心中转站。生产者(设备/网关)发布消息到某个主题,消费者(DolphinDB 接入程序)订阅主题,就能实时收到消息。这个模型天然适合设备海量、上行数据为主的工业场景。
主题规划非常关键,它直接决定后续过滤和分流的成本。我在项目里跟设备厂家反复对齐后,约定了如下格式:
{site}/{device_type}/{device_id}/{metric}实际消息类似这样:
plant1/inverter/INV-001/output_power plant1/environment/TH-102/temperature每个主题下发布的消息 payload 就是对应的测点值,可以是数值、JSON 或者带状态的报文。
这样设计的好处是:订阅端可以根据需求精确订阅,比如只看某个逆变器的数据,就用plant1/inverter/INV-001/#;也可以用plant1/#把整站数据都收下来,再在解析逻辑里从主题中提取 device_id 和 metric。主题层级控制在 4 层以内,避免过深增加路由负担,这一点在设备量大时尤其重要。
2.2 订阅接入的具体配置与消息解析
DolphinDB 官方有一套 MQTT 插件,但不同版本差异较大,而且生产环境里我更喜欢用一个独立的 Python 采集网关来控制消息处理逻辑。原因很简单——消息解析、格式转换、异常处理这些东西,用 Python 验证和调整要快得多,不用每次改动都回 DolphinDB 重挂插件。
采集网关的核心逻辑就两件事:收消息、推数据。收消息用 paho-mqtt 客户端,推数据用 DolphinDB 的 Python API。下面是一个经过了简化但结构完整的示例:
import json import time import dolphindb as ddb import paho.mqtt.client as mqtt DB = ddb.session() DB.connect("127.0.0.1", 8848, "admin", "123456") TOPIC = "plant1/#" def parse_payload(topic, payload): # 按约定主题拆字段:site/type/id/metric parts = topic.split("/") device_id = f"{parts[0]}-{parts[1]}-{parts[2]}" metric = parts[3] blob = json.loads(payload) value = float(blob.get("value", blob.get("v", 0))) ts = blob.get("ts", time.time()) return device_id, metric, value, ts def on_message(client, userdata, msg): try: device_id, metric, value, ts = parse_payload(msg.topic, msg.payload) ts_ms = int(ts * 1000) if isinstance(ts, float) else int(ts) # 组装成服从统一测点流模型的记录 record = (ts_ms, device_id, metric, value, 0) DB.tableInsert("metricStream", record) # 流表追加 except Exception as ex: print("parse error:", ex) client = mqtt.Client() client.on_message = on_message client.connect("192.168.1.100", 1883, 60) client.subscribe(TOPIC, qos=1) client.loop_forever()这里面有两个细节值得单独说。
第一,parse_payload里的时间戳处理。设备厂家发过来的时间戳五花八门,有的是 ISO 字符串,有的是秒级时间戳,有的干脆不带。我统一在接入层就转成 epoch 毫秒整数,避免脏时间戳流到 DolphinDB。
第二,tableInsert的批量性问题。这个示例是单条插入,测试没问题,生产上如果消息量很大,建议攒一批再批量插入。比如用 list 累积 1000 条,或者间隔 1 秒 flush 一次,性能能差一个数量级。
2.3 QoS 选型与消息可靠性权衡
MQTT 的 QoS 有 0、1、2 三档,很多新手直接选 2,觉得最可靠。但我在生产环境中踩过坑:QoS 2 的协议交互开销很大,而且在 Broker 实现不完善时,可能出现比 QoS 1 更多的重复投递,反而增加了数据去重的难度。最终我全线使用 QoS 1,配合业务层去重。
去重怎么做?如果 payload 里带全局唯一的消息 ID,就在 DolphinDB 侧记录最近一段时间的 msg_id 做缓存;如果没带,就只能用时间戳加设备 ID 加测点名组合判断,但对高频率重复数据效果有限。所以我强烈建议需求对接时,要求设备网关侧每一条消息都带唯一的 msg_id,这在工业数采里是一个非常实用且容易被忽略的约定。
断线重连方面,paho-mqtt 自带重连机制,设置好reconnect_delay即可。更重要的问题是重连期间消息丢失。具体到我这套架构,我一般让网关在内存里做一个环形缓存,断线期间的消息先放缓存,重连后再补推。DolphinDB 不会因为网关重启就丢数据,除非网关进程直接挂了。
3. Modbus 接入实战:从寄存器到统一测点
3.1 RTU 还是 TCP:先梳理清楚现场网络
Modbus 在全球工业现场的地位不用多说,几乎所有 PLC、电表、传感器都支持。但 Modbus 有两个大分支:老派的 Modbus RTU 走串口(RS232/RS485),新派的 Modbus TCP 走网口。两者的报文内容基本一致,区别在于串口报文有 CRC16 校验,TCP 报文没有。
我在项目里的选择是:尽可能统一走 Modbus TCP。原因有三点:
- 现场已经布好了局域网络,网线直连或者交换机组网,不用再拉 485 总线。
- TCP 报文可以直接用 pymodbus 库读取,不需要额外处理串口转发的时序问题。
- 后续扩展设备方便,只要设备有网口,插上就能接入,不用考虑 485 总线挂载数量上限和干扰问题。
当然,纯 485 设备也不得不处理。常见的做法是加一个串口服务器,比如有人物的 USR-TCP232 系列,把 RS485 的电平信号转换成 TCP 服务,DolphinDB 侧的采集程序只需要像访问 Modbus TCP 设备一样访问串口服务器的 IP 和端口。此时串口服务器负责串口链路管理,上层代码完全不用区分 RTU 还是 TCP。
3.2 Modbus 报文结构与寄存器模型
不管是 RTU 还是 TCP,Modbus 的核心是读写设备内部的寄存器空间。常见功能码如下:
| 功能码 | 含义 | 典型用途 |
|---|---|---|
| 0x01 | 读线圈 | 开关状态、启停信号 |
| 0x02 | 读离散输入 | 无源触点、限位开关 |
| 0x03 | 读保持寄存器 | 可读写的参数和测量值 |
| 0x04 | 读输入寄存器 | 只读测量值,如电压电流 |
Modbus TCP 报文结构很固定:事务 ID(2 字节)、协议 ID(2 字节)、长度(2 字节)、单元 ID(1 字节)、功能码(1 字节)、数据区。报文格式看着简单,但实际项目中最容易错的不是报文格式,而是寄存器地址映射和字节序。
寄存器地址映射是最常见的坑。设备手册里写"保持寄存器 40001,对应电压值",编程时地址实际是 40001 - 40001 = 0。很多新手直接用 40001 去读,肯定会报错,因为 Modbus 的协议地址是从 0 开始的寄存器编号。不同厂家的手册习惯还不一样,有的用 4xxxx 表示保持寄存器,有的直接写十六进制地址,对表的时候一定要仔细。
字节序问题更隐蔽。Modbus 寄存器是 16 位,如果一个测点值需要 32 位精度,就要连续读两个寄存器,然后拼成一个 32 位整数或浮点数。拼接顺序有四种组合:大端字序加大端字节序、小端字序加小端字节序,以及两种混排。设备不同,组合就不同。我遇到过一个温控器,文档说 IEEE 754 浮点,但实际是低字在前、高字在后,跟默认的大端解析出来的数值差了十万八千里。最稳妥的办法是在接入层用四种组合都试一遍,哪个数值符合物理常识就用哪个,然后硬编码到配置里。
3.3 轮询策略与 CRC 校验细节
Modbus RTU 是一主多从协议,总线上同一时间只能有一个主站发起请求,从站只能在收到针对自己的请求时回复。TCP 模式其实也保留了单元 ID 的"伪从站"概念,可以一台 TCP 设备后面挂多个逻辑设备。
轮询策略直接影响数据实时性和总线载荷。我最初的做法是每台设备顺序轮询,读完全部测点再读下一台,结果在 32 台电表的总线上发现一轮下来要十几秒,部分设备抢答超时导致频繁重试。后来改成两层策略:高优先级测点(电压、电流、功率)用短周期轮询,低优先级测点(电量累计、温度)用长周期轮询,分开两个任务跑,高优数据延迟降到了 2 秒以内。
下面是基于 pymodbus 的简化轮询代码,包含字节序处理:
import time import struct from pymodbus.client import ModbusTcpClient def parse_float(data_bytes): # 依次尝试4种字节序组合 patterns = [">f", "<f", ">f", "<f"] # 实际需要对应4种字序/字节序组合 for pat in patterns: try: val = struct.unpack(pat, data_bytes)[0] if -1e6 < val < 1e6: # 物理合理性检查 return val except Exception: pass return None client = ModbusTcpClient("192.168.1.20", port=502) client.connect() devices = [ {"id": "MTR-001", "unit": 1, "points": [("voltage", 0), ("current", 2)]}, {"id": "MTR-002", "unit": 2, "points": [("voltage", 0), ("current", 2)]}, ] while True: start = time.time() records = [] for dev in devices: for point, addr in dev["points"]: resp = client.read_holding_registers(addr, 2, slave=dev["unit"]) if resp.isError(): continue raw = resp.registers data_bytes = b"".join([x.to_bytes(2, "big") for x in raw]) value = parse_float(data_bytes) if value is not None: records.append((int(time.time() * 1000), dev["id"], point, value, 0)) # 批量写入DolphinDB,这里省略表连接细节 DB.tableInsert("metricStream", records) elapsed = time.time() - start time.sleep(max(0, POLL_INTERVAL - elapsed))CRC 校验只在 RTU 串口通信中出现。pymodbus 内部已经实现了 CRC16 计算,但如果你要在 DolphinDB 内部直接解析 RTU 报文,就需要自己实现。Modbus CRC16 的核心流程是:初始值 0xFFFF,每来一个字节跟当前 CRC 异或,然后右移 8 次,如果最低位是 1 就再异或 0xA001。代码不长,但容易在移位次数上出错,建议先用在线计算工具验证几个标准报文再上逻辑。
4. 统一测点流落地:表建模、写入与消费
4.1 测点流表的设计与建表
不管上游是 MQTT 还是 Modbus,最后都要落到同一张表。DolphinDB 里我的核心表结构是这样设计的:
# 以DolphinDB脚本语言创建流表,具体函数按实际版本微调 t = streamTable(1000000:0, ["ts", "device_id", "point_name", "value", "quality"], [TIMESTAMP, SYMBOL, SYMBOL, DOUBLE, INT]) enableTableShareAndPersistence(table=t, tableName="metricStream", cacheSize=1000000, retentionMinutes=1440)几个关键点:
device_id和point_name用 SYMBOL 类型而不是 STRING。DolphinDB 对 SYMBOL 的字典编码索引效率高得多,在全表扫描和分组聚合时差距非常明显。quality字段是质量码。正常值 0,超量程 1,通信异常 2,解析失败 3。这个字段平时查询可能用不上,但关键时刻排查数据质量问题非常有用。enableTableShareAndPersistence可以把流表同时共享给多个订阅者和客户端,注意 cacheSize 和保留时间的设置,避免内存涨幅失控。
分区方面,如果按device_id哈希分区,同一设备的数据落在同一个分区,查询单设备数据时快;如果你的查询更多是按时间范围跨设备全量统计,按天分区或按小时分区更合适,取决于业务主查询模式。
4.2 批量写入与幂等去重
写入性能是大规模接入的生死线。我的经验是:能批量绝不分条。DolphinDB 的tableInsert支持传入 list of tuples,批量追加效率远高于逐条 insert。采集网关里攒批的逻辑很简单,Python 端用一个列表收集记录,达到 1000 条或者 1 秒超时就批量写入一次。
MQTT 的 QoS 1 可能带来重复消息,处理方案是流表增加一列msg_id,写入前在流计算订阅里做去重,或者维护一个 Redis 去重缓存。如果不想引入 Redis,可以在 DolphinDB 流计算订阅里用duplicate函数,或者对最近时间窗口内的 msg_id 做集合判断。注意去重操作要在流表上游完成,不要在存储层做完再改,否则重复数据已经落库了。
4.3 流计算订阅与物化结果
统一测点流直接存原始数据是最常见的需求,但实际业务往往要的是加工后的结果。DolphinDB 的流计算订阅非常方便,比如我可以挂一个订阅,实时计算每台设备每分钟的平均功率:
def calc_minute_avg(mutable table, msg): avg = select avg(value) as avg, max(value) as max, min(value) as min from msg where point_name in ["active_power"] group by device_id, bar(ts, 60000) tableInsert(minuteStats, avg) subscribeTable(tableName="metricStream", actionName="minuteAvg", offset=-1, handler=calc_minute_avg)这样原始测点流和分钟统计表并行维护,业务报表直接查分钟表,性能压力小很多。类似的订阅还可以做越限报警、趋势异常检测,整个实时计算框架非常统一。
5. 实战踩坑记录与优化心得
5.1 MQTT 侧典型问题排查
| 现象 | 排查思路 | 解决方案 |
|---|---|---|
| 订阅后收不到消息 | 确认 Broker 地址、端口、Topic 是否一致;检查是否有通配符权限限制 | 用 MQTTX 客户端先订阅验证;确认 Broker ACL 配置 |
| 消息偶尔丢 | QoS 设置过低;Broker 负载高时丢弃 | QoS 提到 1;检查 Broker 最大连接数和消息堆积情况 |
| 同一条消息多次写入 | QoS 1/2 重复投递;网关重连后补发重放 | 增加 msg_id 字段,DolphinDB 侧去重 |
| 时间戳乱 | 设备时区不一致;字符串解析错误 | 接入层统一转 epoch 毫秒,时区固定为 UTC+8 |
实际排查时,我最大的经验是:先隔离,再处理。比如消息丢了,先用命令行的 mosquitto_sub 直接订阅该主题,确认是 Broker 侧的推送问题,还是接入网关的解析问题,别一上来就怀疑 DolphinDB。
5.2 Modbus 侧典型问题排查
| 现象 | 排查思路 | 解决方案 |
|---|---|---|
| 读寄存器报错 3(非法数据地址) | 地址越界;寄存器编号换算错误;设备固件不同 | 核对设备原始地址表;用 Modbus Poll 工具实测 |
| 浮点数解析完全不对 | 字节序/字序组合不对 | 四种组合逐一测试,选取物理合理值 |
| 轮询太慢 | 设备多、超时重试累积 | 缩短超时到 500ms;分优先级轮询 |
| 网关掉线后无法恢复 | 串口服务器看门狗未开启;TCP 连接泄漏 | 开启串口服务器看门狗;采集端定期心跳检测 |
Modbus 调试时,Modbus Poll 和 Modbus Slave 这两款工具几乎是标配,一个模拟主站,一个模拟从站。建议先让 Modbus Slave 模拟一台设备,把报文结构、地址映射、字节序都调通,再对接真实设备,能省掉大量现场排查时间。
5.3 几条实操心得
最后分享几条我在这个项目里的亲身体会。
第一,协议适配层一定要像墙一样隔离开。接入层可以五花八门,Python 写也好、Go 写也罢,但往 DolphinDB 推的数据结构必须严格统一。我见过同事把 Modbus 的原始寄存器值直接塞进测点表,美其名曰"保留原始数据",结果后面所有分析逻辑都要加一层转换,非常痛苦。
第二,日志要打全。MQTT 的 broker 日志、采集网关的 log、Modbus 报文级日志都要留。一主多从的设备轮询,越混乱的时候越需要报文日志来定位。我用的是 logging 的 RotatingFileHandler,按天滚动,保留 30 天,排查历史问题时救命。
第三,用模拟器先把链路验证完再上现场。MQTT 这边用 MQTTX 发消息,Modbus 这边用 Modbus Slave 模拟设备,先把 DolphinDB 的流表和订阅调通,再连真实设备。真实设备往往达不到协议文档说的那样规范,尤其是一些小众国产设备,现场排查成本很高,能提前验证的绝不等到现场。
第四,DolphinDB 侧的性能优化,优先考虑批量写入、合理分区、数据压缩,这三项做完基本就够了。不要迷信单机性能,数据规模上来以后,分区策略和索引设计比什么都重要。
我的感受是,多协议接入这个事,说到底是把"杂乱数据"变成"有序数据"的过程。MQTT 和 Modbus 只是起点,后面还会有 OPC UA、行业私有协议或者其他新的接入方式,但只要统一测点流这个抽象层立住了,新协议接入就是写一个适配器的事,不需要动存储和分析的骨架。这套架构我目前跑得比较稳,如果后续项目里再多几种协议,我大概率还会沿用这个思路,只是把接入层的适配器再往上加一层罢了。