数采插件
说明:本文以当前工作区内的源码实现为准。部分 README 与源码存在不一致的地方,例如
mssql-agent的 README 会提到 MySQL,但代码实际使用的是 SQL Server;toolsnet的 README 也更偏旧版描述,但源码已经明确了它的 HTTP 入站、Influx 写入与字段映射逻辑。
总体认知
这个仓库里的每个文件夹都可以看作一个独立的工业数采插件。它们的共同目标不是“采集”本身,而是把现场不同形态的数据源,统一变成可追踪、可落库、可回放的标准点位数据。
这些插件的差异主要体现在四个维度:
- 输入源不同:数据库、MQTT、CSV、HTTP、Access MDB、MSSQL CDC。
- 触发方式不同:定时轮询、文件监听、消息订阅、HTTP 请求、WebSocket 查询。
- 输出形态不同:InfluxDB、文件、HTTP 转发、RabbitMQ、TimescaleDB、浏览页面。
- 状态管理不同:有的依赖 offset 文件,有的依赖 SQLite,有的依赖 last timestamp,有的几乎无状态。
总体架构
实现形态分布
插件矩阵
| 插件 | 输入源 | 输出 | 适用场景 | 核心机制 |
|---|---|---|---|---|
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.txttoolsnetdb使用offset1.txt和offset2.txtmssql-cdc使用last_timestamp.jsonmssql-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做了时区/时间偏移修正promess用CycleDate + StartTime + RunTime还原曲线区间mssql-cdc根据t_stamp和last_timestamp做增量同步
5. 稳定性保护
几乎每个链路都考虑了现场常见的不稳定因素:
- 断路器:
bsnjdb、toolsnetdb - 全局异常捕获:
mqtt-influxdb、mssql-cdc - 日志轮转:多个项目都有
- 优雅关闭:
bsnjdb、mssql-cdc - 配置热重载或注册同步:
bsnjdb、mssql-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.js用node-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}这种通用结构转成 Influxpoint。
理解要点
这个插件的特点是“格式极固定”。它不做通用 CSV 解析,而是直接按业务约定的行列位置抽取,因此非常适合格式稳定的现场导出文件。
3. mqtt-influxdb
适用场景
- 设备通过 MQTT 推送 JSON 报文
- 报文中带有数组型时序数据,需要拆点后入库
- 需要按业务字段过滤消息、并可选择保存原始报文
实现原理
index.js使用mqtt.connect建立持久会话。- 连接成功后订阅配置中的 topics。
common/utils.js的parseToInfluxPoints完成核心转换:- 先按
msg.validates过滤消息。 - 再读取
dataArray.keyNameOfArray所指向的数组。 - 对每个数组项提取时间、值、编码。
- 根据值是否为数字,把字段放到
value或value_str。 - 把额外业务字段映射成 Influx tag。
- 先按
common/influxdb.js构造 InfluxDB 客户端并批量写点。- 如果配置了
influxdb.file = on,还会把原始消息另存成文件,便于追溯。
理解要点
这是典型的“消息总线 -> 时序库”网关。它的重点不在采集,而在于把异构 MQTT 报文统一成稳定的点位模型。
4. mssql-agent
适用场景
- 工厂边缘侧需要统一配置、统一注册、统一轮询的数据代理
- SQL Server 作为源库,InfluxDB 作为时序库
- 需要动态映射离散业务值、维护工位映射、做时间修正的场景
实现原理
- 这是一个 Egg.js 应用。
app.js在didReady阶段监听bind、refreshMap、refresh消息,用 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.js用t_stamp > lastTimestamp做增量查询。src/db/influxdb.js:- 从
mapping.json读取tagid -> item_code。 - 按不同字段类型选择
value或strValue。 - 用
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.mapping把ToolID转成前缀编码 - 对普通字段直接生成点位
- 对
TorqueValue/AngleValue这种曲线字段按TimeCoefficient拆成多个点 - 用
tagField生成额外 tag,如curve_id、sn
- 从
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,最终统一落到时序点位或结构化文件。
这也是它们能在工业现场长期运行的原因。