TDengine 零代码接入 SparkplugB:通过 taosExplorer 将 IIoT 设备数据实时写入时序数据库
【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine
导读
SparkplugB 是面向工业物联网(IIoT)场景的开放消息规范,基于 MQTT 协议定义设备、边缘网关与应用之间的主题命名与载荷格式,广泛用于 SCADA 与工业数据采集。TDengine 通过 taosX 与 taosExplorer 提供零代码数据写入能力,用户无需编写任何代码,即可在浏览器界面中创建任务,从 MQTT 代理订阅 SparkplugB 消息、完成 Payload 解析与字段转换,并实时写入当前 TDengine 集群。本文将完整讲解基于 零代码数据写入 体系创建 SparkplugB 数据同步任务的每一步:连接认证、订阅配置、Payload 转换(解析/拆分/过滤/表映射)、高级选项与异常处理策略,并结合仓库文档资源说明断点恢复、存储转发等关联能力,帮助你在实际项目中快速落地 SparkplugB 数据入库。
功能概述
SparkplugB 是一种开放消息规范,专为工业物联网(IIoT)应用设计,基于 MQTT 协议。它约定了设备、边缘网关、应用(如 SCADA、历史数据库)之间的主题命名空间(namespace)与消息载荷格式,其中载荷(Payload)使用 protobuf 编码,从而保证不同厂商设备之间的互操作性与数据语义一致性。
TDengine 通过 SparkplugB 连接器从 MQTT 代理订阅 SparkplugB 数据并将其写入 TDengine,实现实时数据流入库。整个流程如下:
- MQTT 代理(Broker)中持续接收来自 SparkplugB 设备的消息;
- TDengine 侧部署的 taosX 作为 MQTT 客户端订阅对应主题;
- taosX 按 SparkplugB 规范解析消息(protobuf → JSON),并通过 taosExplorer 中配置的解析、提取、过滤、映射规则完成数据转换;
- 转换后的数据写入 TDengine 的目标数据库与超级表。
该能力属于 零代码数据写入 体系的一部分,支持的数据源版本为 Sparkplug B 3.0。整个过程无需编写代码,配置均在 taosExplorer 界面完成。
提示:本节描述的 SparkplugB 数据写入(taosX 连接外部 MQTT 代理消费数据)与 MQTT 订阅(客户端连接 TDengine Bnode)不是同一功能,两者在协议能力与连接方向上有本质区别。
创建任务
进入 taosExplorer,在左侧导航栏点击数据写入进入数据源列表页面,然后按以下步骤创建 SparkplugB 数据同步任务。
新增数据源
在数据写入页面中,点击+新增数据源按钮,进入新增数据源页面。
配置基本信息
在名称中输入任务名称,例如:test_spb。
在类型下拉列表中选择SparkplugB。
代理是非必填项。如有需要,可以在下拉框中选择已经创建好的代理,也可以先点击右侧的+创建新的代理按钮,创建一个专用的 Agent 来承载该采集任务。关于代理(Agent)的安装,可参考 安装 Agent。
在目标数据库下拉列表中选择一个目标数据库,也可以先点击右侧的+创建数据库按钮,在界面中直接完成建库。
配置连接和认证信息
SparkplugB 基于 MQTT,因此这里配置的是与 MQTT 代理之间的连接参数:
- Brokers:填写 MQTT 代理的地址,例如
localhost:1883。可以填写多个,用,分隔,用于连接多个 broker。 - MQTT 协议:选择使用的 MQTT 协议版本,默认 5.0 版本。作为对照,MQTT 数据源 支持 3.1、3.1.1、5.0 三个版本,默认值为 3.1;SparkplugB 场景默认使用 5.0。
- 客户端 ID:填写连接到每个 broker 所使用的客户端标识符。注意:连接到同一个 MQTT 地址的所有客户端 ID 必须保证唯一,否则会造成客户端 ID 冲突,导致任务无法正常运行。从 MQTT 数据源 的说明可知,填写后通常会在其基础上生成带
taosx前缀的客户端 ID。 - Keep Alive:输入保持活动间隔。保持活动间隔是指客户端和代理之间协商的时间间隔,用于检测客户端是否活动。如果代理在保持活动间隔内没有收到来自客户端的任何消息,它将假定客户端已断开连接,并关闭连接。
- 用户:填写 MQTT 代理的用户名。
- 密码:填写 MQTT 代理的密码。
- TLS 校验:选择 TLS 证书的校验方式,共有三种模式:
- 不开启:表示不进行 TLS 证书认证。在连接 MQTT 时,会先进行 TCP 连接,如果连接失败,会进行无证书认证模式的 TLS 连接。
- 单向认证:开启 TLS 连接,并验证服务端证书,此时需要上传 CA 证书。
- 双向认证:开启 TLS 连接,并与服务端进行双向认证,此时需要上传 CA 证书、客户端证书以及客户端密钥。
配置完成后,点击检查连通性按钮,检查数据源是否可用。如果连通性检查失败,请按照页面上返回的具体错误提示进行修改。
订阅配置
SparkplugB 采用层级化的主题命名空间,本节配置需要订阅哪些 group、节点与设备以及消息类型。
- Group ID:填写 SparkplugB 规定的 group id 字段。通常一个 group id 代表一个集团/公司/工厂/流水线等概念,是主题命名的第一层。
- 节点/设备列表:填写需要订阅的节点和设备的列表,以逗号分隔。其中节点直接填写 ID 即可;设备需要按照
节点 ID/设备 ID格式填写。 - 消息类型:填写需要订阅的 SparkplugB 消息类型,以逗号分隔,可选项包括
NBIRTH/NDEATH/NDATA/NCMD/DBIRTH/DDEATH/DDATA/DCMD/STATE。其中:- 前缀 N 表示节点(Node)相关消息:NBIRTH(节点出生)、NDEATH(节点死亡)、NDATA(节点数据)、NCMD(节点命令);
- 前缀 D 表示设备(Device)相关消息:DBIRTH(设备出生)、DDEATH(设备死亡)、DDATA(设备数据)、DCMD(设备命令);
- STATE 表示 SparkplugB 的宿主应用状态消息。
- 订阅时,
NBIRTH/NDEATH/NDATA/NCMD类型的消息只会匹配“节点/设备列表”中的节点;而DBIRTH/DDEATH/DDATA/DCMD只会匹配“节点/设备列表”中的设备。
- 下发 REBIRTH 命令:开启后,taosX 会自动下发 NCMD 中的
Node Control/Rebirth命令,获取节点和设备的所有 metric 信息,从而可以得到 metric name 与 metric alias 的对应关系。如果节点/设备在上报数据时不使用 alias 别名机制,可以不开启此选项。
配置 Payload 转换
在Payload 解析区域填写 Payload 解析相关的配置参数。这是零代码接入的核心 ETL 环节,包含“解析 → 字段拆分 → 数据过滤 → 表映射”四步,可参考 数据提取、过滤和转换 的通用说明。
解析
SparkplugB 消息使用 protobuf 编码,因此需要先将原始消息体转换为结构化字段。有三种获取示例数据的方法:
- 点击从服务器检索按钮,从 MQTT 获取示例数据;
- 点击文件上传按钮,上传 CSV 文件,获取示例数据;
- 在消息体中填写 MQTT 消息体中的示例数据。
由于 SparkplugB 消息使用 protobuf 进行编码,因此从服务器检索的数据是经过编码为 json 格式的数据。json 数据支持 JSONObject 或者 JSONArray 两种形态,可以用于解析 SparkplugB 中的 metadata 和 properties 等 json 格式的字段。
点击放大镜图标可查看预览解析结果,确认解析出的字段是否符合预期。
字段拆分
在从列中提取或拆分中填写从消息体中提取或拆分的字段。典型场景是:SparkplugB 的 metric 中携带数据类型字符串字段datatype_str,而目标 TDengine 表需要的是对应的 TDengine 原生类型,此时可以用转换规则将datatype_str字段的值转换为 TDengine 类型。
例如,在rule输入框中填写如下 json 值,在name中填写td_datatype:
{ "Int8": "TINYINT", "UInt8": "TINYINT UNSIGNED", "Int16": "SMALLINT", "UInt16": "SMALLINT UNSIGNED", "Int32": "INT", "UInt32": "INT UNSIGNED", "Int64": "BIGINT", "UInt64": "BIGINT UNSIGNED", "Float": "FLOAT", "DOUBLE": "DOUBLE", "Boolean": "BOOL", "String": "VARCHAR(128)", "DateTime": "TIMESTAMP" }该规则会将datatype_str列的值(如"Int8")转换为对应的 TDengine 类型(如"TINYINT"),并生成新列td_datatype。
- 点击删除,可以删除当前提取规则。
- 点击新增,可以添加更多提取规则。
- 点击放大镜图标可查看预览提取/拆分结果。
数据过滤
在过滤中填写过滤条件,只有满足条件的行才会写入 TDengine。例如填写datatype_str != "Int8",则只有datatype_str不为Int8的数据才会被写入 TDengine。
过滤条件表达式的结果必须是 boolean 类型,可依据解析字段的类型使用比较操作符(>、>=、<=、<、==、!=)、逻辑操作符(&&、||、!)以及字符串函数(is_empty、contains、starts_with、ends_with、len)等进行组合,具体表达式写法见 过滤。
点击删除可以删除当前过滤规则;点击放大镜图标可查看预览过滤结果。
表映射
在目标超级表的下拉列表中选择一个目标超级表,也可以先点击右侧的创建超级表按钮创建新的超级表。
当超级表需要根据消息动态生成时,可以选择创建模板。超级表名称、列名、列类型等均可以使用模板变量;当接收到数据后,程序会自动计算模板变量并生成对应的超级表模板:
- 当数据库中超级表不存在时,会使用此模板创建超级表;
- 对于已创建的超级表,如果缺少通过模板变量计算得到的列,也会自动创建对应列。
在映射中填写目标超级表中的子表名称,例如t_{id}。根据需求填写映射规则,其中 mapping 支持设置缺省值。映射规则支持 mapping(直接映射)、value(常量)、generator(生成器,如时间戳now)、join(字符串连接)、format(字符串格式化,${}占位符)、sum(数值求和)、expr(数值运算表达式)等多种类型,详细规则见 映射规则。
点击预览可以查看映射的结果,确认源字段与目标表字段的对应关系是否正确。
高级选项
高级选项区域默认折叠,点击右侧>可展开。MQTT / SparkplugB 等协议类数据源常见项如下(界面字段名可能略有差异):
- 消息等待队列大小:接收消息的缓存队列大小;队列满且未开启缓存实时数据时,新到达的数据会直接丢弃。可设为
0表示不缓存。 - 处理中批次上限(部分界面写作“处理批次上限”):可同时进行数据处理的批次数量;到达上限后不再从消息缓存队列取消息,会导致队列积压。最小值为
1。 - 批次大小:每次发送给数据处理流程的消息数量,与批次延时配合使用:达到批次大小时即使未到延时也会立即发送。最小值为
1。 - 批次延时:每批消息的超时时间(单位:毫秒),从该批第一条消息起算;超时后即使未达批次大小也会立即发送。最小值为
1。 - 写入并发数量:同时写入 TDengine 的并发任务数量。
- 缓存实时数据:开启时,消费数据会先写入本地文件,再由后台任务读出并发送给下游,用于流量削峰;消费完成后会自动清理文件。默认关闭。原理与配置见 存储转发。
- 缓存数据存储目录:缓存目录,仅在开启缓存实时数据时生效。默认为 taosX 启动时配置的数据目录。
- 保存原始数据:开启时可配置最大保留天数与原始数据存储目录。
此外,从v3.3.5.0开始,数据源的高级选项中还增加了健康状态监测相关配置项,包括健康监测时段(Health Check Duration)、Busy 状态阈值(Busy State Threshold,默认 100%)、写入队列长度(Max Write Queue Length)、写入错误阈值(Write Error Threshold),任务健康状态的完整说明见 健康状态。
异常处理策略
异常处理策略区域默认折叠,点击右侧>可展开,用于配置数据出现异常时的处理策略。通用策略说明如下:
- 归档:将异常数据写入归档文件(默认路径为
${data_dir}/tasks/_id/.datetime),不写入目标库; - 丢弃:将异常数据忽略,不写入目标库;
- 报错:任务报错。
各异常项及可选策略如下表:
| 异常项 | 可选处理策略 | 说明 |
|---|---|---|
| 目标库连接超时 | 归档、丢弃、报错、缓存 | 目标库连接失败时可选;缓存:当目标库状态异常(连接错误或资源不足等)时写入缓存文件(默认路径${data_dir}/tasks/_id/.datetime),目标库恢复正常后重新入库 |
| 目标库不存在 | 归档、丢弃、报错 | 写入报错目标库不存在 |
| 表不存在 | 归档、丢弃、报错、自动建表 | 自动建表:自动建表,建表成功后重试 |
| 主键时间戳溢出 | 归档、丢弃、报错 | 检查数据中第一列时间戳是否在正确的时间范围内(now - keep1,now + 100y) |
| 主键时间戳空 | 归档、丢弃、报错、使用当前时间 | 使用当前时间:使用当前时间填充到空的时间戳字段中 |
| 复合主键空 | 归档、丢弃、报错 | 写入报错复合主键空 |
| 表名长度溢出 | 归档、丢弃、报错、截断、截断且归档 | 子表表名长度限制最大 192 字符;截断:截取原始表名的前 192 个字符作为新的表名;截断且归档:截断并同时将记录写入归档文件 |
| 表名非法字符 | 归档、丢弃、报错、非法字符替换为指定字符串 | 检查子表表名中是否包含特殊字符(如符号.等);替换模式例如a.b替换为a_b |
| 表名模板变量空值 | 丢弃、留空、变量替换为指定字符串 | 留空:变量位置不做处理,例如a_{x}转换为a_;变量替换:例如a_{x}转换为a_b |
| 列名不存在 | 归档、丢弃、报错、自动增加缺失列 | 自动增加缺失列:根据数据信息自动修改表结构增加列,修改成功后重试 |
| 列名长度溢出 | 归档、丢弃、报错 | 列名长度限制最大 64 字符 |
| 列自动扩容 | 开关选项 | 打开时,列数据长度超长将自动修改表结构并重试 |
| 列长度溢出 | 归档、丢弃、报错、截断、截断且归档 | 截断:截取数据中符合长度限制的前 n 个字符;截断且归档:截断并写入归档文件 |
| 数据异常 | 归档、丢弃、报错 | 其他未在上方列出的数据异常 |
此外,异常处理策略区域还包含以下全局配置:
- 连接超时:目标库连接超时时间,单位“秒”,取值范围 1~600;
- 临时存储文件位置:缓存文件的位置,实际生效位置为
${data_dir}/tasks/:id/{location}; - 归档数据保留天数:非负整数,0 表示无限制;
- 归档数据可用空间:0~65535,其中 0 表示无限制;
- 归档数据文件位置:归档文件的位置,实际生效位置为
${data_dir}/tasks/:id/{location}; - 归档数据失败处理策略:当写入归档文件报错时的处理策略,可选删除旧文件(删除旧文件,若仍无法写入则报错并停止任务)、丢弃(丢弃即将归档的数据)、报错并停止任务。
创建完成与任务管理
点击提交按钮,完成创建 SparkplugB 到 TDengine 的数据同步任务,回到数据源列表页面可查看任务执行情况。
在任务列表页面,可以对任务进行启动、停止、查看、删除、复制等操作,也可以查看各个任务的运行情况,包括写入的记录条数、流量等,并通过任务活动日志排查失败原因,通过运行指标页面(每 2 秒自动刷新,分为累计指标与本次运行指标)观察运行状态。
断点恢复说明
根据 任务断点恢复 的说明,SparkplugB 目前不支持消息持久化和数据恢复。因此在任务重启或网络中断等场景下,中断期间的消息可能丢失。如需降低网络中断导致的数据丢失风险,可以在高级选项中开启缓存实时数据(存储转发),将消费数据先持久化到本地磁盘、网络恢复后自动补发——该机制由 taosX-Agent 侧的persist_data_enable/persist_data_dir配置驱动,详细原理见 存储转发;同时要注意,开启缓存会引入额外磁盘读写,增加端到端延迟,对延迟敏感且网络稳定的场景可以不开启。
健康状态
从v3.3.5.0开始,任务管理列表新增“健康状态”列,用于指示任务运行过程中的健康状态,状态包括 Ready、Idle、Active、Pending、Busy、Bounce、SourceError、SinkError、Fatal 等,当健康状态为空时表示尚未有数据开始入库。详细状态语义见 健康状态。
小结
本文完整介绍了在 TDengine 中通过 taosExplorer 创建 SparkplugB 零代码数据写入任务的完整流程:从新增数据源、配置 MQTT 连接与 TLS 认证,到订阅 Group/节点/设备与消息类型,再到 Payload 转换四步曲(解析、字段拆分、数据过滤、表映射),以及高级选项、异常处理策略与提交后的任务管理。整体能力基于 零代码数据写入 体系,与同体系下的 MQTT 数据源共享连接、ETL 与任务管理机制。掌握这些配置后,即可在数分钟内完成 SparkplugB 设备数据的实时入库,为 SCADA、产线监控等 IIoT 场景提供统一的时序数据底座。
【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考