数采插件

说明:本文以当前工作区内的源码实现为准。部分 README 与源码存在不一致的地方,例如 mssql-agent 的 README 会提到 MySQL,但代码实际使用的是 SQL Server;toolsnet 的 README 也更偏旧版描述,但源码已经明确了它的 HTTP 入站、Influx 写入与字段映射逻辑。

总体认知

这个仓库里的每个文件夹都可以看作一个独立的工业数采插件。它们的共同目标不是“采集”本身,而是把现场不同形态的数据源,统一变成可追踪、可落库、可回放的标准点位数据。

这些插件的差异主要体现在四个维度:

  1. 输入源不同:数据库、MQTT、CSV、HTTP、Access MDB、MSSQL CDC。
  2. 触发方式不同:定时轮询、文件监听、消息订阅、HTTP 请求、WebSocket 查询。
  3. 输出形态不同:InfluxDB、文件、HTTP 转发、RabbitMQ、TimescaleDB、浏览页面。
  4. 状态管理不同:有的依赖 offset 文件,有的依赖 SQLite,有的依赖 last timestamp,有的几乎无状态。

总体架构

现场输入源
数据库 MQTT CSV MDB HTTP
触发与调度层
定时轮询 文件监听 消息订阅 HTTP 接入 WebSocket 查询
解析与映射层
字段映射 标签补充 时间修正 类型转换 offset / checkpoint
标准点位
measurement tag field timestamp
输出层
InfluxDB 文件留档 HTTP 转发 RabbitMQ 页面查询

实现形态分布

按触发方式归类
定时轮询型
5
事件监听型
2
HTTP 服务型
1

插件矩阵

插件 输入源 输出 适用场景 核心机制
bsnjdb PostgreSQL 紧固数据表 InfluxDB / 文件 / RabbitMQ Bosch 类拧紧结果、步骤、曲线的增量采集与归档 定时轮询、offset 文件、summary/step/graph 拆分、断路器、管理接口
csv-tesla 目录中的 CSV 文件 InfluxDB 现场设备导出的定制 CSV 曲线文件入库 文件监听、固定列号解析、曲线点展开、点位编码拼装
mqtt-influxdb MQTT topic 消息 InfluxDB / 文件 设备通过 MQTT 上送 JSON 数组数据的场景 Topic 订阅、消息校验、数组展开、额外标签映射、批量写入
mssql-agent SQL Server InfluxDB 需要集中配置、定时抽取、带注册与心跳的边缘采集代理 Egg 调度、SQLite 配置缓存、动态映射、时间修正、增量 offset
mssql-cdc MSSQL CDC 表 InfluxDB / WebSocket 页面 CDC 增量同步、历史点位查询、可视化监控 last_timestamp 持久化、tagid 映射、可选页面服务、Chart.js 展示
promess Access MDB 文件 / MDB 目录 工作目录内的结构化输出文件 Promess 紧固机曲线文件解析、批量备份、离线转换 ADODB 读取、Cycle/Curve 分组、XML 曲线解析、按 mapping 生成输出
toolsnet HTTP POST 数据报文 InfluxDB 作为数采接入网关,接收上游系统推送的拧紧报文 接口认证、请求校验、字段映射、曲线点时间展开、标签补充
toolsnetdb MySQL HTTP API / 文件 从数据库抽取结果后转发到中心服务,兼顾落盘留档 定时轮询、双 offset、单轴/多轴分流、断路器、Bearer token 转发

共性实现原理

1. 配置驱动

几乎所有插件都把业务差异放进配置里:

  • 数据源连接信息
  • 字段映射关系
  • 目标数据库或目标接口
  • 拉取频率
  • 文件输出开关
  • 认证信息

这样做的好处是,插件本体只保留“采集、转换、投递”的通用逻辑,现场适配主要靠配置完成。

2. 增量处理

为了避免重复采集,大部分插件都做了 checkpoint:

  • bsnjdb 使用 offset.txt
  • toolsnetdb 使用 offset1.txtoffset2.txt
  • mssql-cdc 使用 last_timestamp.json
  • mssql-agent 把偏移写进 SQLite 的 t_cfg_offset

这类设计非常适合工业现场的“持续采集、允许断点恢复”需求。

3. 统一点位模型

无论输入是 SQL 行、MQTT JSON、CSV 行还是 MDB 结果,最终都要落到类似下面的模型:

  • measurement
  • tags
  • fields
  • timestamp

在 InfluxDB 体系里,点位编码通常由设备前缀、工位、螺栓号、字段编号组合而成。这样可以把现场复杂的业务对象压缩成稳定可查询的时序键。

4. 时间处理

工业数采里,时间不是“附属字段”,而是组织数据的主轴。

  • bsnjdb 直接用微秒级时间范围计算 summary、step、graph 的结束时间
  • csv-tesla 把 CSV 中的时间字段转换成微秒时间戳
  • mssql-agent 做了时区/时间偏移修正
  • promessCycleDate + StartTime + RunTime 还原曲线区间
  • mssql-cdc 根据 t_stamplast_timestamp 做增量同步

5. 稳定性保护

几乎每个链路都考虑了现场常见的不稳定因素:

  • 断路器:bsnjdbtoolsnetdb
  • 全局异常捕获:mqtt-influxdbmssql-cdc
  • 日志轮转:多个项目都有
  • 优雅关闭:bsnjdbmssql-cdc
  • 配置热重载或注册同步:bsnjdbmssql-agent

插件说明

1. bsnjdb

适用场景

  • Bosch 类拧紧工位的数据归档
  • 需要把“概要数据 + 步骤数据 + 曲线数据”统一转成 InfluxDB 点位的场景
  • 采集后还要把原始数据保存成文件,或在时间范围上触发后续任务的场景

实现原理

  • configure/configure.yaml 中的 fetch-interval 定时触发。
  • 从 PostgreSQL 查询 overview / step / wave 三张关联表,按 offset.txt 做增量。
  • common/processor.js 会:
    • resultid 分组,确保 summary 只生成一次。
    • 对 summary、step、graph 三种数据分别映射。
    • mapping.yaml 把业务字段编码成 item_code
    • 把数值和字符串分流到 value / value_str
    • 把曲线数组按时间拆成多个点。
  • common/sender.js 支持两种输出:
    • 保存到 wd/multi/*.txt
    • 通过 HTTP 写入 InfluxDB
  • common/notify.js 在成功处理后可向 RabbitMQ 投递时间范围消息,格式是 triggerCode,start,end
  • 内置 opossum 断路器和 HTTP 管理接口,提供 /reload/status/health

理解要点

这个插件本质上是一个“高完整度紧固结果转换器”。它不是简单搬运数据库行,而是把一条拧紧记录拆成三个层次的数据模型,再统一映射成时序点位。


2. csv-tesla

适用场景

  • 设备或上位机导出的固定格式 CSV 文件
  • 需要按文件名和固定列号读取数据的离线/准实时场景

实现原理

  • app.jsnode-watch 监听目录 D:/temp
  • 新文件到达后:
    • 从文件名中解析设备编号。
    • 调用 process.js 解析 CSV。
    • 调用 util.js 转换为 Influx 点位对象。
    • 调用 writeData.js 批量写入 InfluxDB。
  • process.js 通过固定列号切分数据块:
    • buildResult 处理结果区
    • buildProcess 处理过程区
    • buildSummary 处理摘要区
    • buildCurve 处理曲线区
    • buildSequence 处理事件序列区
  • util.js 负责把 {time, item_code, value} 这种通用结构转成 Influx point

理解要点

这个插件的特点是“格式极固定”。它不做通用 CSV 解析,而是直接按业务约定的行列位置抽取,因此非常适合格式稳定的现场导出文件。


3. mqtt-influxdb

适用场景

  • 设备通过 MQTT 推送 JSON 报文
  • 报文中带有数组型时序数据,需要拆点后入库
  • 需要按业务字段过滤消息、并可选择保存原始报文

实现原理

  • index.js 使用 mqtt.connect 建立持久会话。
  • 连接成功后订阅配置中的 topics。
  • common/utils.jsparseToInfluxPoints 完成核心转换:
    • 先按 msg.validates 过滤消息。
    • 再读取 dataArray.keyNameOfArray 所指向的数组。
    • 对每个数组项提取时间、值、编码。
    • 根据值是否为数字,把字段放到 valuevalue_str
    • 把额外业务字段映射成 Influx tag。
  • common/influxdb.js 构造 InfluxDB 客户端并批量写点。
  • 如果配置了 influxdb.file = on,还会把原始消息另存成文件,便于追溯。

理解要点

这是典型的“消息总线 -> 时序库”网关。它的重点不在采集,而在于把异构 MQTT 报文统一成稳定的点位模型。


4. mssql-agent

适用场景

  • 工厂边缘侧需要统一配置、统一注册、统一轮询的数据代理
  • SQL Server 作为源库,InfluxDB 作为时序库
  • 需要动态映射离散业务值、维护工位映射、做时间修正的场景

实现原理

  • 这是一个 Egg.js 应用。
  • app.jsdidReady 阶段监听 bindrefreshMaprefresh 消息,用 messenger 在进程内同步配置和动态映射。
  • app/schedule/heartbeat.js 负责注册逻辑;当前源码里远端注册请求被简化为本地固定 nodeId/token 绑定,属于可继续接入中心服务的预留位置。
  • app/schedule/collectData.js 定时执行:
    • 读取 SQLite 中的节点配置。
    • 逐表查询 SQL Server。
    • 调用 mappingData.doMapping 做字段映射。
    • 写入 InfluxDB。
    • 如有新的动态映射值,则同步回 SQLite 并广播更新。
  • configureNode.js 把配置、设备映射、字段映射、偏移量统一存入 SQLite。
  • dynamicMapping.js 负责把离散业务值转成连续序列号,并缓存到内存中。
  • timeRule.js 对时间做 8 小时偏移修正。

理解要点

它更像“可集中管理的边缘代理”,不是纯粹的采集脚本。重点在于:

  • 配置可以下发和缓存
  • 采集状态可以增量恢复
  • 业务值可以动态编码
  • 设备与字段映射可以独立维护

5. mssql-cdc

适用场景

  • 直接消费 MSSQL CDC 表
  • 需要按 t_stamp 增量同步到 InfluxDB
  • 想要一个轻量的浏览页面查看最近历史点位

实现原理

  • src/index.js 读取 .env,决定是否启用页面服务。
  • 程序启动后:
    • 连接 MSSQL。
    • 根据 last_timestamp.json 读取上次处理位置。
    • 定时轮询 CDC 表。
    • 成功后写入 InfluxDB,并更新 last_timestamp.json
  • src/db/mssql.jst_stamp > lastTimestamp 做增量查询。
  • src/db/influxdb.js
    • mapping.json 读取 tagid -> item_code
    • 按不同字段类型选择 valuestrValue
    • t_stamp 作为点位时间。
  • src/page-server/ 提供可选浏览服务:
    • WebSocket 接口返回采集项列表与最近记录
    • 前端用 Chart.js 做散点图展示

理解要点

这个插件的关键是“CDC + checkpoint + 可视化”。它既是同步器,也是一个小型历史数据观察窗口。


6. promess

适用场景

  • Promess 紧固机生成的 MDB / Access 数据文件
  • 需要把离线数据库内容解析成标准化曲线输出
  • 需要批量备份或实时处理落地文件的场景

实现原理

  • index.js 通过 node-adodb 打开 MDB。
  • 支持两种模式:
    • --dd:监控目录,发现新 .mdb 文件就批处理
    • --db:处理单个数据库文件
  • common/parser.js 做三件事:
    • CycleId 分组
    • mapping.json 映射基本字段
    • 解析 XmlCurves 中的曲线点
  • common/influx.js 负责把解析结果写成结构化行,实际上落的是工作目录中的结果文件,而不是直接写 InfluxDB。
  • mapping.json 里定义了:
    • 普通字段如何映射到 AAxxx / ABxxx
    • 曲线字段按单位匹配到 BAxxx

理解要点

它是“离线数据库解析器 + 曲线还原器”。和其他插件不同,它更偏文件加工和结果落地,核心价值在于把 MDB 里的曲线信息结构化地抽出来。


7. toolsnet

适用场景

  • 上位系统通过 HTTP POST 推送拧紧报文
  • 需要把报文转换成 InfluxDB 点位
  • 需要对输入做认证和字段校验

实现原理

  • 这是一个 Egg.js API 服务。
  • 路由只暴露了一个资源:/OpenApi/data/bolt_curve_report
  • topics.js 中的 create() 完成完整链路:
    • Basic 认证校验
    • 请求体结构校验
    • 调用 ctx.service.influx.transformer.transform
    • 再写入 InfluxDB
  • app/service/influx/transformer.js
    • config.fields 中读取字段规则
    • config.mappingToolID 转成前缀编码
    • 对普通字段直接生成点位
    • TorqueValue / AngleValue 这种曲线字段按 TimeCoefficient 拆成多个点
    • tagField 生成额外 tag,如 curve_idsn
  • app/service/influx/writeData.js 定义了写入格式,支持数值和字符串两个字段族。

理解要点

它本质上是一个“工业报文接入 API”。和 mqtt-influxdb 相比,它不是订阅消息,而是提供 HTTP 接口给上游主动推送。


8. toolsnetdb

适用场景

  • MySQL 中已有拧紧结果,需要定时抽取并转发到中心服务
  • 需要同时保留原始 JSON 文件作为留档
  • 需要把单轴和多轴数据分别处理

实现原理

  • index.js 使用 node-schedule 定时触发。
  • 维护两个 offset 文件:
    • offset1.txt 对应单轴
    • offset2.txt 对应多轴
  • processData(type)
    • 按类型选择 SQL
    • 查询后如果开启 target.file,先落地到 single/multi/
    • 再通过 sendToRemote() 把 JSON 数组 POST 到目标 HTTP API
    • 最后更新 offset
  • common/processor.js 只负责三件事:
    • 单轴结果文件保存
    • 多轴结果文件保存
    • HTTP 转发
  • 这个插件也用了 opossum 做断路器,避免中心服务异常时把本地采集拖垮。

理解要点

它更像“数据库中继器”。数据并不在本地深度转换,而是以较轻量的方式抽取出来,交给中心系统继续处理。

选型建议

  • 如果源头是 PostgreSQL 紧固结果,优先看 bsnjdb
  • 如果源头是 CSV 文件,优先看 csv-tesla
  • 如果源头是 MQTT JSON,优先看 mqtt-influxdb
  • 如果需要 边缘代理 + 中央配置 + SQL Server,优先看 mssql-agent
  • 如果需要 CDC 增量同步,优先看 mssql-cdc
  • 如果数据来自 Access MDB / Promess 曲线文件,优先看 promess
  • 如果要做 HTTP 入站网关,优先看 toolsnet
  • 如果要做 数据库到中心服务的中继,优先看 toolsnetdb

一个更实用的结论

这些插件虽然面向不同协议和数据源,但它们几乎都遵循同一个工程范式:

采集入口尽量简单,业务复杂度集中在映射层,历史连续性靠 offset 或 checkpoint,最终统一落到时序点位或结构化文件。

这也是它们能在工业现场长期运行的原因。


数采插件
https://luischen.github.io/2026/06/18/manulism-work/20_Domain_Knowledge/21 工业数据接入/03 数采插件/
作者
Luis Chen
发布于
2026年6月18日
许可协议