TDengine 零代码接入 SparkplugB:通过 taosExplorer 将 IIoT 设备数据实时写入时序数据库
2026/9/13 13:31:43 网站建设 项目流程

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,实现实时数据流入库。整个流程如下:

  1. MQTT 代理(Broker)中持续接收来自 SparkplugB 设备的消息;
  2. TDengine 侧部署的 taosX 作为 MQTT 客户端订阅对应主题;
  3. taosX 按 SparkplugB 规范解析消息(protobuf → JSON),并通过 taosExplorer 中配置的解析、提取、过滤、映射规则完成数据转换;
  4. 转换后的数据写入 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 证书的校验方式,共有三种模式:
    1. 不开启:表示不进行 TLS 证书认证。在连接 MQTT 时,会先进行 TCP 连接,如果连接失败,会进行无证书认证模式的 TLS 连接。
    2. 单向认证:开启 TLS 连接,并验证服务端证书,此时需要上传 CA 证书。
    3. 双向认证:开启 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 编码,因此需要先将原始消息体转换为结构化字段。有三种获取示例数据的方法:

  1. 点击从服务器检索按钮,从 MQTT 获取示例数据;
  2. 点击文件上传按钮,上传 CSV 文件,获取示例数据;
  3. 消息体中填写 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_emptycontainsstarts_withends_withlen)等进行组合,具体表达式写法见 过滤。

点击删除可以删除当前过滤规则;点击放大镜图标可查看预览过滤结果。

表映射

目标超级表的下拉列表中选择一个目标超级表,也可以先点击右侧的创建超级表按钮创建新的超级表。

当超级表需要根据消息动态生成时,可以选择创建模板。超级表名称、列名、列类型等均可以使用模板变量;当接收到数据后,程序会自动计算模板变量并生成对应的超级表模板:

  • 当数据库中超级表不存在时,会使用此模板创建超级表;
  • 对于已创建的超级表,如果缺少通过模板变量计算得到的列,也会自动创建对应列。

映射中填写目标超级表中的子表名称,例如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 - keep1now + 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询