TDengine 零代码接入 SparkplugB:基于 taosExplorer 的 IIoT 数据同步任务配置指南
【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine
本文基于 TDengine 开源仓库 docs/en/08-data-ingest-and-delivery/01-no-code-ingestion/18-sparkplugb.md 编写。SparkplugB 是专为工业物联网(IIoT)设计、构建于 MQTT 之上的开放消息规范。借助 TDengine 的零代码数据接入平台,你可以在 taosExplorer 图形界面中直接创建数据源任务,让 taosX 连接器从 MQTT Broker 订阅 SparkplugB 消息并实时写入 TDengine 集群——无需编写任何代码。读完本文,你将掌握从新增数据源、配置连接认证、订阅过滤、Payload 转换(解析/拆分/过滤/表映射)到高级选项与异常处理策略的完整任务创建流程。
背景:为什么需要 SparkplugB 数据接入
SparkplugB 是一种开放的消息规范,专为工业物联网(IIoT)应用设计,底层基于 MQTT 协议。它定义了 IIoT 场景下设备与 MQTT Broker 之间的标准消息格式(如 NBIRTH、NDATA、DDATA 等),被广泛用于工厂、产线、设备监控等工业数据采集场景。
TDengine 通过内置的 SparkplugB 连接器,可以从 MQTT Broker 订阅 SparkplugB 消息,并将数据实时写入 TDengine,实现工业数据的实时入库。整个流程在 taosExplorer 的数据写入(Data In)页面通过零代码配置完成,无需额外部署 ETL 工具。
从 零代码接入平台总览 可知,taosExplorer 是 TDengine 的可视化数据管理工具,支持在浏览器中通过简单配置向 TDengine 提交任务,实现多种数据源到 TDengine 的零代码导入,并在导入过程中自动完成数据的抽取、过滤与转换。SparkplugB 正是该平台支持的众多数据源之一。
说明:SparkplugB 连接器在消息断点续传方面有明确限制——从 任务断点续传 一节可知,与 MQTT、Kafka 等数据源不同,SparkplugB 当前不支持消息持久化与恢复,任务重启后无法从上次断点续传。
创建数据写入任务
新增数据源
登录 taosExplorer 后,进入数据写入页面,点击+ 新增数据源(Add Data Source)按钮,进入任务创建页面。
配置基本信息
在任务创建页面中配置以下基本信息:
- 名称:输入任务名称,例如
test_spb; - 类型:在下拉列表中选择SparkplugB;
- 代理(可选):选择一个已创建的 Agent 代理,或点击右侧+ 创建新的代理按钮新建;
- 目标数据库:在下拉列表中选择一个目标数据库,或点击右侧+ 创建数据库按钮新建。
配置连接与认证信息
在连接配置区域,需要填写 MQTT Broker 的接入参数:
| 配置项 | 说明 |
|---|---|
| Brokers | MQTT Broker 地址,例如localhost:1883。可以填写多个,用逗号,分隔,用于连接多个 Broker |
| MQTT 协议 | 使用的 MQTT 协议版本,默认5.0(可选 3.1、3.1.1、5.0) |
| 客户端 ID | 连接到每个 Broker 时使用的客户端标识符 |
| Keep Alive | 保持活动间隔。如果 Broker 在该间隔内没有收到来自客户端的任何消息,会假定客户端已断开并关闭连接。该间隔是客户端与 Broker 之间协商的、用于检测客户端活跃状态的时长 |
| 用户 / 密码 | MQTT Broker 认证所需的用户名与密码(若 Broker 开启了认证则必填) |
实操提示:同一 MQTT Broker 下如果创建多个同步任务,各任务的客户端 ID 必须互不相同,否则会造成冲突导致任务无法正常运行。
TLS 校验模式
TLS 校验(TLS Verification)支持三种模式:
- 不开启(Disabled):不进行 TLS 证书认证。连接 MQTT 时会先尝试 TCP 连接,若失败则改为无证书认证模式的 TLS 连接。
- 单向认证(One-way):开启 TLS 连接并验证服务端证书,此时需要上传CA 证书。
- 双向认证(Mutual):开启 TLS 连接并与服务端进行双向认证,此时需要上传CA 证书、客户端证书以及客户端私钥。
配置完成后,点击检查连通性(Check Connectivity)按钮验证数据源是否可用;若检查失败,请根据页面返回的具体错误提示修改配置。
订阅配置
订阅配置决定连接器从哪些主题、哪些设备、哪些消息类型中消费数据。
- Group ID:填写 SparkplugB 规范定义的 group id。通常一个 group id 代表一个集团/公司/工厂/流水线等概念。
- 节点/设备列表(Node/Device List):填写需要订阅的节点和设备列表,以逗号分隔。其中节点直接填写 ID 即可;设备需要按照
节点ID/设备ID的格式填写。 - 消息类型(Message Types):填写需要订阅的 SparkplugB 消息类型,以逗号分隔。支持的类型包括:
NBIRTH、NDEATH、NDATA、NCMD、DBIRTH、DDEATH、DDATA、DCMD、STATE。- 其中
NBIRTH、NDEATH、NDATA、NCMD类型的消息只会匹配"节点/设备列表"中的节点; - 而
DBIRTH、DDEATH、DDATA、DCMD只会匹配"节点/设备列表"中的设备。
- 其中
- 下发 REBIRTH 命令(Send REBIRTH Command):开启后,taosX 会自动下发 NCMD 中的
Node Control/Rebirth命令,从而获取节点和设备的所有 metric 信息,包括 metric name 与 metric alias 的对应关系。如果节点/设备在上报数据时不使用 alias 别名机制,可以不开启此选项。
配置 Payload 转换
Payload 转换是 SparkplugB 数据接入的核心环节,包含解析、字段拆分、数据过滤、表映射四步,是 taosX 内置 ETL 能力的具体体现(参见 数据抽取、过滤与转换)。
解析 Payload
Payload 解析区域提供三种获取示例数据的方式:
- 点击从服务器检索:从已配置的 MQTT Broker 获取示例数据;
- 点击文件上传:上传文件获取示例数据;
- 在消息体中手动填写 MQTT 消息体的示例数据。
由于 SparkplugB 消息使用Protocol Buffers(protobuf)编码,从服务器检索到的数据会先被解码为 JSON 格式。JSON 数据支持 JSONObject 或 JSONArray 两种形态,可以用于解析 SparkplugB 中的 metadata、properties 等 JSON 格式字段。
点击放大镜图标可预览解析结果:
从列中提取或拆分字段
解析后的数据可能仍不满足目标表的要求,此时可以在从列中提取或拆分(Extract or Split)区域填写提取/拆分规则。
典型的场景是将datatype_str字段的值转换为 TDengine 数据类型。选择映射(mapping)提取器,在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。
你可以点击新增添加更多提取规则,点击删除移除当前规则,点击放大镜图标预览提取/拆分结果。
数据过滤
在过滤(Filter)区域填写过滤表达式,只有满足条件的数据行才会被写入 TDengine。
例如填写datatype_str != "Int8",则只有datatype_str值不为Int8的数据才会被写入。
过滤表达式的结果必须为布尔类型,支持基于字段类型的判断函数与比较运算符(>、>=、<=、<、==、!=),多个条件可通过逻辑运算符(&&、||、!)组合。例如location.starts_with("beijing") && voltage > 200表示只同步北京地区电压大于 200 的智能电表数据。相关过滤语法细节可参考 零代码接入平台的过滤章节。
点击删除可移除当前过滤规则,点击放大镜图标可预览过滤结果。
表映射
表映射将解析、提取、拆分后的源字段映射到 TDengine 目标表。
- 目标超级表:在下拉列表中选择一个目标超级表,或点击右侧创建超级表按钮新建。
- 创建模板:当超级表需要根据消息动态生成时,选择创建模板。此时超级表名称、列名、列类型等均可以使用模板变量。接收到数据后,程序会自动计算模板变量并生成对应的超级表模板:
- 当数据库中该超级表不存在时,使用模板创建超级表;
- 对于已创建的超级表,如果缺少通过模板变量计算得到的列,也会自动创建对应列。
- 映射:填写目标超级表中的子表名称,例如
t_{id};根据需求填写映射规则,其中 mapping 支持设置缺省值(默认值)。
点击预览可查看映射结果,确认子表名称、列与标签的映射是否符合预期。
配置高级选项
高级选项(Advanced Options)区域默认折叠,点击>展开。MQTT 与 SparkplugB 数据源常用的选项如下(字段名可能因连接器而异,参见 高级选项详解):
| 选项 | 说明 |
|---|---|
| Message Queue Size(消息队列大小) | 接收缓冲区大小。队列满且未开启缓存实时数据时,新到达的数据会被丢弃;设为0表示禁用缓冲 |
| Maximum In-Process Batches(最大进行中批次) | 可并发处理的批次数上限。达到上限后连接器停止从接收队列取消息,消息会在队列中累积;最小值为1 |
| Batch Size(批量大小) | 每次送入处理管道的消息条数。与批量延迟配合使用:即使延迟未到,批量已满也会立即发送;最小值为1 |
| Batch Delay(批量延迟) | 每批次的超时时间(毫秒),从该批次第一条消息到达开始计时。超时后即使未达到 Batch Size 也会发送该批次;最小值为1 |
| Write Concurrency(写入并发) | 并发写入 TDengine 的任务数 |
| Cache Realtime Data(缓存实时数据) | 开启后,消费到的数据先写入本地文件,由后台任务转发下游,当下游处理跟不上时起到流量整形作用;积压消费完毕后缓存文件会被清除。默认关闭。详见 Store and Forward |
| Cache Storage Directory(缓存存储目录) | 覆盖缓存文件的存储目录,仅在开启缓存实时数据时生效,否则默认使用 taosX 启动时配置的数据目录 |
| Save Raw Data(保存原始数据) | 开启后可进一步配置最大保留天数与原始数据存储目录 |
此外,高级选项中还包含健康监控设置(Health Check Duration、Busy State Threshold、Max Write Queue Length、Write Error Threshold),用于任务列表页的健康状态展示,具体说明参见 Health Status。
配置异常处理策略
异常处理策略(Exception Handling Strategy)区域默认折叠,点击>展开。taosX 为各类写入异常提供了统一的分流策略(参见 异常处理策略详解):
- 归档(Archive):将无效数据写入归档文件(默认位于
${data_dir}/tasks/<id>/<datetime>下),不写入目标数据库; - 丢弃(Discard):忽略无效数据;
- 报错(Error):报告错误;
- 缓存(Cache):目标连接失败或资源不足时,将数据写入缓存文件,待目标恢复后再行入库。
可针对以下场景分别配置策略:
- 目标连接超时:归档 / 丢弃 / 报错 / 缓存;
- 目标数据库不存在:归档 / 丢弃 / 报错;
- 表不存在:归档 / 丢弃 / 报错 / 自动建表并重试;
- 主时间戳超出范围(
now - keep1至now + 100y):归档 / 丢弃 / 报错; - 主时间戳为空:归档 / 丢弃 / 报错 / 使用当前时间;
- 复合主键为空:归档 / 丢弃 / 报错;
- 表名超过 192 字符:归档 / 丢弃 / 报错 / 截断 / 截断并归档;
- 表名含非法字符(如
.):归档 / 丢弃 / 报错 / 用配置的字符串替换非法字符; - 表名模板变量为空:丢弃 / 变量留空 / 用配置的字符串替换;
- 列不存在:归档 / 丢弃 / 报错 / 自动补列并重试;
- 列名超过 64 字符:归档 / 丢弃 / 报错;
- 列值超出定义长度:归档 / 丢弃 / 报错 / 截断 / 截断并归档,也可通过自动扩列修改表结构后重试;
- 其他数据错误:归档 / 丢弃 / 报错。
附加设置项:
- 连接超时(Connection Timeout):目标连接超时时间(秒),取值范围
1~600; - 临时存储位置:相对
${data_dir}/tasks/<id>/的路径; - 归档保留天数(Archive Retention Days):非负整数,
0表示不限; - 归档可用空间(Archive Available Space):取值范围
0~65535,0表示不限; - 归档位置(Archive Location):相对
${data_dir}/tasks/<id>/的路径; - 归档写入失败策略:删除旧文件 / 丢弃数据 / 报错并停止任务。
提交任务
完成上述所有配置后,点击提交(Submit)按钮,即完成 SparkplugB 到 TDengine 的数据同步任务创建,自动回到数据源列表(Data Source List)页面。
提交成功后,可在任务列表页查看任务执行情况,包括写入记录数、流量等运行指标;任务状态会切换为 Running。你也可以在任务列表页对任务进行启动、停止、查看、删除、复制等管理操作,并查看每个任务的健康状态(Ready、Idle、Active、Pending、Busy、Bounce、SourceError、SinkError、Fatal 等,详见 任务管理)。
小结
SparkplugB 数据接入任务的核心链路可概括为:taosX 连接器订阅 MQTT Broker → 解码 protobuf 为 JSON → 解析/拆分/过滤 → 映射到超级表与子表 → 实时写入 TDengine。整个过程完全通过 taosExplorer 的零代码界面完成,涵盖连接认证(含 TLS 单向/双向认证)、订阅配置(Group ID、节点/设备、消息类型、REBIRTH)、Payload 转换(四种 ETL 步骤)以及高级选项与异常处理兜底策略。配置时需特别注意:SparkplugB 当前不支持消息持久化与断点续传,对于需要高可靠连续采集的工业场景,建议结合网络稳定性保障与异常归档策略共同使用。
相关参考文档:
- SparkplugB 接入指南(英文原档)
- SparkplugB 接入指南(中文原档)
- 零代码数据接入平台总览
- taosX Agent 存储转发(Store and Forward)
【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考