十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

TDengine 零代码接入 SparkplugB:通过 taosExplorer 将 IIoT 设备数据实时写入时序数据库

TDengine 零代码接入 SparkplugB:通过 taosExplorer 将 IIoT 设备数据实时写入时序数据库 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.1SparkplugB 场景默认使用 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 - 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),仅供参考
返回列表