ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

DolphinDB多协议接入实战:MQTT、Modbus到统一测点流

DolphinDB多协议接入实战:MQTT、Modbus到统一测点流 1. 工业数据接入这件事为什么值得单独拎出来讲搞工业物联网和智能制造的朋友都有一个共识数据接入是整个数据链路的“第一公里”这一公里修不好后面的存储、计算、分析全是空中楼阁。我过去几年做过不少产线数据采集项目从注塑机、数控机床到环境传感器、电表水表设备品牌五花八门通信协议更是各说各话。有的设备只给你一个 RS485 口跑 Modbus RTU有的网关支持 MQTT 主动上报还有的老旧 PLC 只认 OPC UA。你要把这些数据统一收上来还要保证时序对齐、不丢点、能实时算光靠一个采集脚本是撑不住的。这个项目标题“DolphinDB 多协议接入实战从 MQTT、Modbus 到统一测点流”说的就是这件事用 DolphinDB 作为统一的数据底座把 MQTT、Modbus、OPC UA 这几类主流工业协议的数据接进来最终归一化成一张统一的测点流表。它解决的核心问题是协议异构、数据格式不统一、时间戳对不齐、写入吞吐上不去这几个老大难。适合谁看做工业数据采集的工程师、做设备联网的集成商、做产线数字化的技术负责人以及任何需要把现场设备数据接进时序数据库的开发者。哪怕你之前只玩过 Modbus Poll 点点鼠标或者只用过 MQTT 客户端订阅过消息这篇文章也能让你把整条链路串起来。我先把结论摆在这DolphinDB 在这套方案里的角色不只是“存数据”它同时承担了协议接入、流式预处理、测点模型统一和实时计算四个职责。这跟传统“采集程序写库、数据库只管存”的思路完全不一样也是我觉得这套方案值得写一篇实战总结的原因。下面我会从整体设计、协议细节、实操步骤到踩坑排查一层层拆开讲。2. 整体架构设计与协议选型思路2.1 为什么是“统一测点流”而不是“一设备一表”很多刚接触工业采集的人第一反应是每台设备建一张表字段就是设备自己的寄存器。这个思路在小规模场景下没问题但设备一多就崩了。你想想一条产线 50 台设备每台 20 个测点如果一设备一表就是 50 张表、50 套 schema查询的时候要 union 半天做跨设备对比分析更是噩梦。更麻烦的是同型号设备换个固件寄存器地址可能就变了表结构跟着改维护成本极高。统一测点流的核心思想是“窄表 标签”所有测点数据写进同一张流表用deviceId、pointName、tag这类维度字段区分来源用timestamp、value承载时序和数值。这样做的好处是 schema 稳定、写入路径统一、查询灵活想按设备过滤就过滤想按测点聚合就聚合。代价是单表数据量大但这恰恰是时序数据库擅长的事DolphinDB 的分区表和流表就是为这种场景设计的。我选这个方案还有一个现实原因现场协议太多如果每个协议单独建表接入层代码会重复到你想哭。统一测点流之后MQTT 接入、Modbus 接入、OPC UA 接入最终都往同一张流表里灌接入层只需要负责“解析成标准测点格式”后面的存储和计算完全复用。2.2 三种协议的定位差异与接入策略MQTT、Modbus、OPC UA 这三者不是竞争关系而是互补关系理解它们的定位差异才能设计出合理的接入策略。协议典型场景通信模式数据特点接入策略MQTT网关、智能传感器、无线设备发布/订阅主动上报消息驱动频率不定订阅主题回调解析ModbusPLC、电表、老设备主从轮询被动读取寄存器映射周期采集定时轮询批量读取OPC UA数控机床、高端 PLC客户端/服务端订阅或读取信息模型丰富带语义订阅节点或周期读取MQTT 是“设备推给你”Modbus 是“你去问设备”OPC UA 两者都支持。这个差异直接决定了接入层的线程模型MQTT 用回调或异步消费Modbus 用定时任务轮询OPC UA 用订阅回调或轮询。在 DolphinDB 里这三类接入最终都通过 API 写入流表但触发方式完全不同这是实操中最容易踩坑的地方后面会细讲。2.3 DolphinDB 流表在架构中的位置很多人以为流表只是“临时缓冲”其实在 DolphinDB 里流表是一等公民。它可以持久化、可以被多个订阅者消费、可以触发实时计算引擎。在这套方案里流表承担三个角色第一写入缓冲。现场设备可能瞬间上报大量数据流表作为内存表先接住再异步落盘到分区表避免直接写磁盘造成抖动。第二格式归一。不同协议解析出来的数据字段可能不一致流表的 schema 强制统一倒逼接入层做标准化。第三实时计算入口。流表可以挂接响应式状态引擎、横截面引擎、时间序列引擎做实时告警、滑动平均、异常检测。这是“统一测点流”最大的价值——数据一进来就能算不用等落盘再查。我实测下来单节点 DolphinDB 流表写入吞吐可以轻松跑到几十万条每秒对于绝大多数产线场景绰绰有余。关键是 schema 设计要合理字段类型要精简别把一堆字符串塞进流表。3. 核心细节解析与实操要点3.1 MQTT 接入订阅、解析、写入三步走MQTT 接入的核心是订阅主题 回调解析 写入流表。我以最常见的“网关上报 JSON 格式测点数据”为例。首先MQTT 主题设计要有层次比如factory/line1/device001/data这样可以用通配符factory/line1//data一次订阅整条产线。主题层级不要太深否则订阅匹配效率下降也不要用纯数字主题可读性太差。其次消息体解析要健壮。现场网关的 JSON 经常不规范比如数值有时是字符串有时是数字时间戳有时是毫秒有时是秒。我的做法是在回调里做类型强制转换和单位归一统一转成timestamp毫秒长整型和valuedouble。# Python 侧 MQTT 回调示例伪代码展示解析逻辑 import json import time def on_message(client, userdata, msg): payload json.loads(msg.payload.decode(utf-8)) # 时间戳归一兼容秒和毫秒 ts payload.get(ts, time.time() * 1000) if ts 1e12: # 秒级时间戳 ts ts * 1000 # 数值归一 value float(payload.get(value, 0)) # 写入 DolphinDB 流表 ddb_session.run( insert into mqttStream values(?, ?, ?, ?, ?) , [payload[deviceId], payload[pointName], ts, value, mqtt])注意MQTT 回调里不要做耗时操作解析完立刻写入复杂计算交给 DolphinDB 流计算引擎。回调阻塞会导致消息堆积甚至丢消息。3.2 Modbus 接入轮询策略与批量读取优化Modbus 接入的难点不在协议本身而在轮询效率和错误处理。Modbus RTU 走串口波特率通常 9600 或 19200一次读几十个寄存器就要几百毫秒Modbus TCP 走网口快很多但设备响应也可能慢。我的经验是按寄存器地址连续性分组批量读取。比如设备有 10 个测点地址分别是 40001、40002、40003、40010、40011、40012、40020……那就分成三组读40001-40003、40010-40012、40020 起。这样一次请求读多个寄存器比逐个读快好几倍。# Modbus 批量读取示例pymodbus from pymodbus.client import ModbusTcpClient client ModbusTcpClient(192.168.1.100, port502) # 批量读取 40001-40010共 10 个保持寄存器 result client.read_holding_registers(address0, count10, slave1) if not result.isError(): for i, reg in enumerate(result.registers): point_name fpoint_{i1} value reg * scale_factor # 根据量程做缩放 write_to_ddb(device_id, point_name, value, modbus)提示Modbus 寄存器地址有“协议地址”和“文档地址”之分文档写 40001 通常对应协议地址 0写代码时要注意偏移。这个坑我踩过不止一次读出来的数据全是错的。轮询周期要根据设备响应时间和测点数量算。假设一次批量读耗时 200ms有 20 台设备串行轮询一轮就是 4 秒。如果要求 1 秒采集一次就必须并行轮询或者用多串口。别指望单线程轮询能扛住大规模设备这是很多采集程序卡顿的根因。3.3 OPC UA 接入订阅模式与节点映射OPC UA 比前两者复杂但它的信息模型也最丰富。接入 OPC UA 有两种方式订阅Subscription和读取Read。订阅适合变化频繁的测点服务端主动推送读取适合变化慢或需要按需获取的场景。我一般用订阅模式设置publishingInterval为 1000mssamplingInterval为 500ms。注意这两个参数的关系采样间隔是服务端采集数据的频率发布间隔是服务端推送数据的频率。如果采样比发布快中间会做聚合如果发布比采样快可能推重复值。# OPC UA 订阅示例opcua-asyncio from asyncua import Client async def subscribe_opcua(): async with Client(urlopc.tcp://192.168.1.200:4840) as client: node client.get_node(ns2;sMachine1.Temperature) subscription await client.create_subscription(1000, handler) await subscription.subscribe_data_change(node)节点映射是 OPC UA 接入的关键。现场设备的节点 ID 往往很长很乱比如ns2;sChannel1.Device1.Tag1直接写进测点流可读性差。我的做法是建一张映射表把节点 ID 映射成统一的deviceId和pointName接入层查表转换。这张映射表可以放在 DolphinDB 的维度表里接入时关联查询也可以缓存在接入程序内存里。3.4 统一测点流的 Schema 设计Schema 设计是整套方案的地基我给出一个经过实战验证的版本字段名类型说明timestampTIMESTAMP测点时间毫秒精度deviceIdSYMBOL设备唯一标识pointNameSYMBOL测点名称valueDOUBLE测点数值protocolSYMBOL来源协议mqtt/modbus/opcuaqualityINT数据质量码0 正常非 0 异常用 SYMBOL 而不是 STRING 存 deviceId 和 pointName是因为 SYMBOL 在 DolphinDB 里做了字典编码存储和查询效率高很多。value 统一用 DOUBLE整数、浮点都兼容别为了省空间用 INT后面遇到浮点测点又要改表。注意quality 字段别省。现场数据经常有“通信正常但数值无效”的情况比如传感器断线返回 0 或 -9999这时候 quality 标记异常后续计算可以过滤。没有质量码脏数据会污染整个分析结果。4. 实操过程与核心环节实现4.1 环境准备与 DolphinDB 流表创建先创建流表和持久化分区表。流表用streamTable持久化表用createPartitionedTable按日期和设备哈希分区。// 创建流表 mqttStream streamTable(1000000:0, timestampdeviceIdpointNamevalueprotocolquality, [TIMESTAMP, SYMBOL, SYMBOL, DOUBLE, SYMBOL, INT]) enableTableShareAndPersistence(tablemqttStream, tableNameunifiedStream, cacheSize1000000) // 创建持久化分区表 db database(dfs://iot, VALUE, 2024.01.01..2025.12.31) pt db.createPartitionedTable(mqttStream, unifiedPoint, timestampdeviceId)流表缓存设为 100 万行够缓冲几分钟的高频数据。持久化表按日期分区方便按时间范围查询和删除旧数据。设备哈希分区是为了避免单设备数据倾斜。4.2 MQTT 接入完整链路实现MQTT 接入我用 Python 的 paho-mqtt 库配合 DolphinDB Python API。完整流程是连接 MQTT Broker → 订阅主题 → 回调解析 → 批量写入 DolphinDB。批量写入是关键优化点。不要一条消息写一次那样网络往返开销太大。我的做法是回调里先塞进本地队列另起一个线程每 100ms 或每 500 条批量写入。import paho.mqtt.client as mqtt from queue import Queue import threading import dolphindb as ddb queue Queue() session ddb.session() session.connect(localhost, 8848, admin, 123456) def batch_writer(): while True: batch [] while len(batch) 500: try: batch.append(queue.get(timeout0.1)) except: break if batch: session.run( insert into unifiedStream values(?, ?, ?, ?, ?, ?) , batch) threading.Thread(targetbatch_writer, daemonTrue).start() def on_message(client, userdata, msg): data parse_payload(msg.payload) queue.put([data[ts], data[deviceId], data[pointName], data[value], mqtt, 0]) client mqtt.Client() client.on_message on_message client.connect(broker.local, 1883) client.subscribe(factory///data) client.loop_forever()实测下来批量 500 条写入单线程可以跑到 5 万条每秒以上。如果还不够可以开多个写入线程但要注意 DolphinDB 单会话并发写入的限制。4.3 Modbus 轮询采集实现Modbus 我用 pymodbus 的同步客户端配合定时任务。核心是分组批量读 异常重试 质量码标记。from apscheduler.schedulers.background import BackgroundScheduler from pymodbus.client import ModbusTcpClient def poll_device(device): try: client ModbusTcpClient(device[ip], portdevice[port]) for group in device[groups]: result client.read_holding_registers( addressgroup[start], countgroup[count], slavedevice[slave]) if result.isError(): mark_quality(device[id], group, quality1) continue for i, reg in enumerate(result.registers): point group[points][i] value reg * point[scale] point[offset] write_to_ddb(device[id], point[name], value, modbus, 0) except Exception as e: mark_quality(device[id], None, quality2) finally: client.close() scheduler BackgroundScheduler() scheduler.add_job(poll_device, interval, seconds5, args[device_config]) scheduler.start()轮询周期我设 5 秒因为这条产线的工艺参数变化不快。如果是高速采集场景比如振动监测Modbus 就不合适了得换 OPC UA 或专用采集卡。4.4 OPC UA 订阅接入实现OPC UA 我用 asyncua 库异步订阅。注意异步框架和 DolphinDB 同步 API 的配合我一般用队列解耦。import asyncio from asyncua import Client, ua async def opcua_handler(node, val, data): ts int(data.monitored_item.Value.ServerTimestamp.timestamp() * 1000) value float(val) queue.put([ts, node_to_device[node], node_to_point[node], value, opcua, 0]) async def main(): async with Client(urlopc.tcp://192.168.1.200:4840) as client: subscription await client.create_subscription(1000, opcua_handler) for node_id in node_list: node client.get_node(node_id) await subscription.subscribe_data_change(node) while True: await asyncio.sleep(1) asyncio.run(main())OPC UA 订阅的坑在于断线重连。网络抖动或服务端重启后订阅会失效必须监听连接状态并自动重建订阅。我在这上面吃过亏设备明明在跑数据却断了半小时才发现。4.5 流表到持久化表的自动落盘流表数据要落盘到分区表用 DolphinDB 的订阅机制。可以写一个订阅把流表数据实时写入分区表。subscribeTable(tableNameunifiedStream, actionNamesaveToDFS, handlerappend!{loadTable(dfs://iot, unifiedPoint)}, msgAsTabletrue, batchSize10000, throttle1)batchSize 设 1 万throttle 设 1 秒意思是攒够 1 万条或等 1 秒就落盘一次。这个参数要根据数据量调太小落盘频繁影响性能太大内存占用高。5. 常见问题与排查技巧实录5.1 数据时间戳错乱怎么排查时间戳错乱是工业采集最常见的问题表现是数据在时序表里顺序不对或者同一时刻出现大量重复。原因通常有三类设备时钟不准、协议时间戳单位不统一、接入层用了本地时间。排查方法先看原始报文里的时间戳字段确认单位是秒还是毫秒再对比设备时钟和服务器时钟差多少最后检查接入层代码是不是有的地方用time.time()有的地方用设备时间。我的原则是统一用设备时间设备没时间就用网关时间都没有才用服务器时间并且全部转成毫秒。5.2 Modbus 读不到数据的几种典型情况现象可能原因排查方法超时无响应IP/端口错、从站地址错ping 通、telnet 端口、确认 slave id返回异常码寄存器地址越界、功能码不支持查设备手册确认地址范围数据全 0寄存器地址偏移错试 0 基和 1 基地址数据跳变量程缩放错、字节序错确认 scale 和字节序大端/小端间歇性失败串口干扰、轮询太快降波特率、加延时、换屏蔽线字节序这个坑特别隐蔽。Modbus 是 16 位寄存器32 位浮点要占两个寄存器有的设备高字在前有的低字在前读出来数值完全不对。我一般先用 Modbus Poll 手动读一遍确认字节序再写代码。5.3 MQTT 消息丢失与重复消费MQTT 的 QoS 等级决定消息可靠性QoS 0 最多一次可能丢QoS 1 至少一次可能重复QoS 2 恰好一次开销大。工业场景我一般用 QoS 1接入层做幂等处理用deviceId pointName timestamp做去重键。消息丢失还有一个原因是客户端 ID 冲突。两个客户端用同一个 clientId 连同一个 Broker会互相踢下线。现场部署时一定要保证 clientId 唯一我习惯用采集程序名_设备ID_随机数的格式。5.4 流表写入性能上不去怎么优化写入性能瓶颈通常在这几个地方单条写入、字段类型太重、流表缓存太小、DolphinDB 会话数不够。优化顺序是先改批量写入再精简 schemaSTRING 改 SYMBOL再调大流表 cacheSize最后考虑多会话并行。我实测过一个对比单条写入 5000 条每秒批量 500 条写入 5 万条每秒差了 10 倍。批量写入是性价比最高的优化没有之一。5.5 OPC UA 订阅断线重连处理OPC UA 订阅断线后不会自动恢复必须手动重建。我的做法是起一个监控线程定期检查订阅状态发现失效就重新创建订阅并重新订阅所有节点。同时记录断线时间断线期间的数据标记 quality 为异常避免误判。async def monitor_subscription(client, subscription, nodes): while True: await asyncio.sleep(10) try: await subscription.get_monitored_items() except: subscription await client.create_subscription(1000, handler) for node in nodes: await subscription.subscribe_data_change(node)6. 测点流后续扩展与个人经验统一测点流建好之后能做的事情很多。我目前在这套基础上挂了三个实时计算任务滑动平均滤波去除传感器毛刺、阈值告警超限立即推送、设备停机检测连续 N 个周期数值不变判定停机。这些都用 DolphinDB 的流计算引擎实现不用额外写程序。如果后续要接更多协议比如 BACnet、CANopen思路是一样的接入层解析成标准测点格式写入同一张流表。统一测点流的价值就在于接入层可扩展存储和计算层不用动。最后分享一个我踩过的坑别在接入层做复杂计算。我一开始图省事在 MQTT 回调里做滑动平均结果回调阻塞消息堆积最后丢数据。后来把所有计算移到 DolphinDB 流引擎接入层只做解析和写入稳定多了。接入层的职责就是“搬运工”搬得越快越稳越好计算的事交给专业引擎。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进