
TDengine 零代码接入 MQTT通过 taosExplorer 构建 MQTT 到 TDengine 的实时数据同步任务【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengineMQTTMessage Queuing Telemetry Transport是物联网与工业物联网场景中最流行的轻量级消息传输协议。本文以 07-mqtt.md 为主体系统讲解如何借助 TDengine 内置的可视化管理工具 taosExplorer 与数据接入组件 taosX通过零代码的向导式界面创建从 MQTT Broker 到 TDengine 集群的数据迁移任务。读完本文你将掌握 MQTT 连接的认证与 TLS 配置、采集参数、主题解析、JSON Payload 解析、字段拆分、数据过滤、表映射、高级选项与异常处理策略的完整配置方法并能结合仓库源码理解其底层机制。概述为什么用 MQTT 接入 TDengineMQTT 是一种基于发布/订阅Publish/Subscribe模式的轻量级消息协议专为低开销、低带宽占用、易实现的即时消息场景设计被广泛应用于物联网、小型设备、移动应用等领域。TDengine 可以通过 MQTT 连接器订阅 MQTT Broker 上的数据并写入 TDengine实现实时的数据流接入。TDengine 的零代码数据接入平台以 taosExplorer可视化数据管理工具和 taosX数据接入引擎为核心。用户无需编写代码只需在浏览器中完成简单的配置即可提交任务实现多种数据源向 TDengine 的无缝导入。在导入过程中TDengine 自动完成数据的提取解析、过滤与转换保证写入数据的质量用户无需额外部署 ETL 工具从而简化整体架构设计。整个零代码平台的系统架构可参见 index.md。根据 index.md 中支持数据源清单MQTT 数据源已通过验证的 Broker 版本范围包括emqx3.0.0 至 5.7.1hivemq4.0.0 至 4.31.0mosquitto1.4.4 至 2.0.18值得说明的是TDengine 仓库中还内置了一个 MQTT Broker 的实现位于 source/libs/tmqtt它作为 TDengine 的 BnodeBroker Node能力存在除网关接入外MQTT 采集配置中指定初始订阅位置等能力即与将 TDengine Bnode 用作 MQTT Broker 的场景相关本文后续会展开说明。前置准备在创建 MQTT 数据源任务之前需要完成两项准备确认 taosX 与数据源之间的网络可达性。如果 taosX 无法直接访问 MQTT Broker例如 Broker 位于隔离的 OT 网络或受限内网需要先安装并配置 taosX-Agent将其部署在数据源所在的网络中作为代理。具体安装步骤Windows 与 Linux 两个平台见 01-install-agent.md组件完整参考见 taosx-agent 组件文档。需要注意的是taosX-Agent 属于 TDengine 企业版TSDB Enterprise组件需单独下载不在社区版安装包中。登录 taosExplorer 并进入数据源页面。在左侧导航栏点击Data In数据接入进入任务列表页点击 Add Data Source新增数据源进入任务创建页面输入任务名称在数据源类型中选择MQTT随后可选择新建代理或复用已创建的代理。配置连接与认证信息在任务创建页的MQTT 连接配置区域填写 MQTT Broker 的基础连接信息配置项说明示例MQTT AddressMQTT 地址MQTT Broker 的 IP 地址或主机名192.168.1.42MQTT PortMQTT 端口MQTT Broker 的监听端口1883User用户连接 MQTT Broker 使用的用户名若 Broker 开启了认证taos_userPassword密码连接 MQTT Broker 使用的密码******其中MQTT 地址 端口也可合并为192.168.1.100:1883这样的地址端口形式输入对应 index.md 中创建任务的第 4 步。配置完成后可点击Check Connection检查连接按钮验证数据源是否可用。配置 TLS 加密如果 MQTT Broker 启用了 SSL/TLS 加密需要在TLS VerificationTLS 校验区域选择校验模式共三种Disabled禁用不校验 TLS 证书。连接器将首先尝试 TCP 直连若失败则尝试不校验证书的 TLS 连接。One-way authentication单向认证启用 TLS 并校验服务器证书需要上传 CA 证书。Mutual authentication双向认证 / 双向 TLS启用双向 TLS需要同时上传 CA 证书、客户端证书和客户端私钥文件。单向认证适用于仅需确认 Broker 身份的场景双向认证则进一步向 Broker 证明客户端身份适用于安全要求更高的生产环境。配置采集信息Collection Configuration采集配置区域集中了订阅行为的核心参数逐项说明如下。MQTT 协议版本从MQTT ProtocolMQTT 协议下拉框选择协议版本共三个选项3.1、3.1.1、5.0默认值为3.1。注意MQTT 5.0 在协议层面新增了用户属性User Properties、共享订阅等能力只有选择5.0时才可配置对应的高级选项见下文。Client ID在Client ID客户端标识中输入客户端标识系统会在其前面自动拼接前缀taosx。例如输入foo生成的客户端 ID 为taosxfoo。如果开启末尾的开关则当前任务的 task id 会被拼接在taosx之后、所输入标识之前生成的客户端 ID 形如taosx100foo。连接同一 MQTT 地址的所有客户端 ID 必须唯一——在为同一 Broker 创建多个同步任务时若客户端 ID 冲突会导致任务无法正常运行因此建议开启 task id 拼接开关以确保唯一性。Keep AliveKeep Alive保活间隔是客户端与 Broker 之间协商的用于检查客户端是否活跃的时间间隔。如果 Broker 在保活间隔内未收到客户端的任何消息会认为客户端已断开并关闭连接同理若客户端在保活间隔内未向 Broker 发送消息Broker 也会断开连接。配置时应大于消息上报的典型周期避免因消息稀疏被误判离线。Clean SessionClean Session清理会话用于选择是否清理会话状态默认值为true。开启时Broker 不保留该客户端的会话状态含离线期间的遗嘱消息与未确认 QoS 消息关闭时Broker 会保留会话客户端重连后可恢复未完成的订阅与消息投递。结合 Keep Alive 与 Clean Session可以控制断线重连后的消息行为。连接用户属性与订阅用户属性MQTT 5.0当选择 MQTT5.0协议时可以配置自定义的Connection User Properties连接用户属性和Subscription User Properties订阅用户属性。这是 MQTT 5.0 新增的特性允许在 CONNECT 与 SUBSCRIBE 报文中携带自定义键值对常用于传递租户标识、环境信息等扩展数据。当使用 TDengine Bnode 作为 MQTT Broker 时还可以在订阅中指定初始订阅位置Initial Subscription Position实现从指定位置开始消费配合断点续传能力保证数据不丢不重。Topics Qos Config主题与 QoS在Topics Qos Config中填写要订阅的主题名称与 QoS 级别格式为{topic_name}::{qos}例如my_topic::0。QoS 取值为0、1、2分别表示至多一次、至少一次、恰好一次的投递语义语义强度与开销依次递增。MQTT 协议 5.0 还支持共享订阅Shared Subscription允许多个客户端订阅同一主题实现负载均衡。共享订阅的格式为$share/{group_name}/{topic_name}::{qos}其中$share是固定前缀表示启用共享订阅group_name是客户端组名称其作用类似 Kafka 中的消费组Consumer Group。同一组内的多个消费者分摊该主题的消息适合通过多个 taosX 任务横向扩展消费吞吐。Topic Analysis主题解析Topic Analysis主题解析用于将 MQTT 主题的每一级解析为对应的变量名格式与 MQTT 主题本身一致使用/分隔层级_表示解析时忽略当前层级。例如MQTT 主题a//c解析规则v1/v2/_其含义是将第一级a赋值给变量v1将第二级的值通配符可匹配任意值赋值给变量v2第三级c被忽略不赋给任何变量。解析得到的主题变量可以继续参与后续 Payload 解析中的各类转换与计算例如用作表名模板变量或过滤条件。Compression消息体压缩Compression压缩配置消息体的压缩算法。taosX 收到消息后会使用对应的压缩算法解压消息体以还原原始数据。可选值包括none不压缩默认gzipsnappylz4zstd需要与消息发布方使用的压缩算法保持一致否则无法正确解压。Char Encoding字符编码Char Encoding字符编码配置消息体的编码格式。taosX 收到消息后使用对应的编码格式解码消息体以还原原始数据。可选值包括UTF_8默认GBKGB18030BIG5对于中文等非 ASCII 内容务必根据发布端实际编码选择否则会出现乱码。连接检查完成上述配置后点击Check Connection检查连接按钮检查数据源是否可用若检查失败页面会返回具体的错误提示可根据提示修改配置后重试。配置 MQTT Payload 解析MQTT Payload ParsingMQTT 消息体解析区域是数据 ETL 的核心。taosX 使用 JSON 提取器解析数据并允许用户指定数据库中的数据模型包括指定表名与超级表名、设置普通列与标签列等。该区域的完整交互逻辑与 index.md 中数据提取、过滤与转换一节描述的 ETL 能力一一对应。获取样例数据有三种方式获取样例数据点击Retrieve from Server从服务器获取按钮从 MQTT 获取样例数据点击File Upload文件上传按钮上传 CSV 文件获取样例数据在Message Body消息体文本框中直接填写 MQTT 消息体的样例数据。每条样例数据以回车结尾。JSON 数据支持 JSONObject 或 JSONArray 两种形态JSON 解析器可解析以下数据{id: 1, message: hello-word} {id: 2, message: hello-word}或[{id: 1, message: hello-word},{id: 2, message: hello-word}]解析后的结果以表格形式展示如图 mqtt-04 所示点击放大镜图标可以查看解析结果的预览mqtt-05。关于 JSON 解析的细节index.md 补充了更多能力JSON 解析支持嵌套对象与数组如data.voltage、location[0].province这类层级字段会被自动解析可自由选择需要解析的字段并为解析字段设置别名解析出的字段类型由 JSON 属性值自动推断——布尔值推断为 bool 类型、整数推断为 int、小数推断为 float、字符串推断为 string。此外对于 nginx 日志这类非结构化文本可使用带命名捕获组的正则表达式提取多字段对于需要定制逻辑的场景还可使用 rhai 脚本UDT实现自定义解析输入为 JSON 解析后的对象映射输出必须为数组。字段拆分Extract or Split from Column解析出的字段未必直接满足目标表的数据要求。在Extract or Split from Column从列中提取或拆分中可以填写要从消息体中提取或拆分的字段。例如将message字段按-拆分为message_0和message_1选择Split拆分提取器分隔符Separator填写-拆分数Number填写2。拆分字段的命名规则为{原字段名}_{序号}。此外也可以使用正则提取器要求使用命名捕获组为提取出的字段命名。点击Delete删除当前提取规则点击Add添加更多提取规则点击放大镜图标预览提取/拆分结果。index.md 中给出的典型示例是智能电表上报{voltage: 221V, ...}电压值带单位字符串可通过正则提取器将数值与单位分开再配合类型转换函数parse_int/parse_float可将字符串转为整数/浮点数转换为数值列后写入。数据过滤Filter在Filter过滤中填写过滤条件过滤条件表达式的结果必须是布尔类型。例如填写id ! 1则只有 id 不等于 1 的数据会被写入 TDengine。点击Delete删除当前过滤规则点击放大镜图标预览过滤结果。过滤表达式的书写与字段类型强相关index.md 给出了完整的语法规则bool 类型直接使用变量或!取反如inuse、!inuse数值类型int/float支持、!、、、、比较运算符字符串类型支持is_empty()、contains(sub)、starts_with(prefix)、ends_with(suffix)等函数以及len属性返回字符数需与比较运算符配合如s.len 5复合表达式可用逻辑运算符、||、!组合多个条件例如location.starts_with(beijing) voltage 200表示筛选北京地区电压大于 200 的智能电表数据。表映射Table Mapping在Target Supertable目标超级表下拉框中选择目标超级表或点击右侧的Create Supertable创建超级表按钮新建。如果超级表需要根据每条消息动态生成则选择Create Template创建模板。超级表名称、列名称、列类型都可以包含模板变量数据到达时 taosX 会先对变量求值缺失的超级表会被自动创建已有超级表则自动补充缺失的列。在Mapping映射中填写目标超级表中的子表名称例如t_{id}。随后按照需求填写各字段的映射规则映射支持设置默认值。映射规则的完整类型见下表来源index.md规则说明mapping直接映射需选择映射源字段value常量可输入字符串常量或数值常量直接存储所输入的常量值generator生成器目前仅支持时间戳生成器 now存储时写入当前时间join字符串连接可指定连接字符拼接选中的多个源字段format字符串格式化如${year}-${month}-${day}${}为占位符占位符可以是源字段或字符串函数处理结果sum对多个数值字段做加法计算expr数值运算表达式支持更复杂的函数处理与数学运算format中可用的字符串处理函数包括pad(len, pad_chars)填充到指定长度、trim去除首尾空白、sub_string(start_pos, len)截取子串起始位置为负时从末尾计数、replace(substring, replacement)替换子串。expr除支持、-、*、/四则运算外还支持sin、cos、tan、sqrt、exp、ln、log、floor、ceiling、round、int、fraction等数学函数例如采集温度为摄氏度的场景可用表达式temperature * 1.8 32转换为华氏度后写入。子表名本身也是字符串可以直接用format表达式定义如t_{id}。点击Preview预览查看映射结果。特殊的透视Pivot语义如果超级表列名本身是模板变量则映射执行的是透视操作——模板变量的值成为列名被映射的字段提供列的值。即由数据动态决定列的形态适用于将某些枚举值横向展开为多列的场景。高级选项Advanced Options高级选项区域默认折叠点击展开。该区域的参数详见 resources/_02-advanced-options-mqtt.mdx如下参数说明Message Queue Size消息队列大小接收缓冲区大小。若队列已满且未开启缓存实时数据新到达的数据会被丢弃设为0表示禁用缓冲Maximum In-Process Batches最大处理中批次可并发处理的批次数上限。达到上限后连接器停止从接收队列取消息消息在队列中累积。最小值为1Batch Size批大小每次送入处理管道的消息数量与批延迟配合即使延迟未到批满也会立即发送。最小值为1Batch Delay批延迟每批的毫秒级超时从该批第一条消息到达开始计时。超时后即使未达到批大小也会发送该批。最小值为1Write Concurrency写并发并发写入 TDengine 的任务数Cache Realtime Data缓存实时数据开启后消费到的数据先写入本地文件由后台任务转发给下游在下游处理跟不上时起到流量整形作用积压消费完后缓存文件被清理。默认关闭。该能力即 taosX-Agent 的 Store and Forward存储转发机制详见 store-and-forward.mdCache Storage Directory缓存存储目录覆盖缓存文件的存储目录仅在开启缓存实时数据时生效否则使用 taosX 启动时配置的数据目录Save Raw Data保存原始数据开启后可进一步配置Maximum Retention Days最大保留天数与Raw Data Storage Directory原始数据存储目录其中Cache Realtime Data对 MQTT 数据源尤为重要。根据 store-and-forward.md 的说明Agent 与 taosX 之间网络中断时Agent 仍会持续从数据源采集数据并写入本地磁盘的持久队列persist_queue网络恢复后自动从断点位置补发缓存数据实现零丢失缓存数据只受磁盘空间限制已确认发送成功的缓存文件会被自动清理。同时该文档也提醒开启缓存后数据需先落盘再上报端到端写入延迟会有所增加对写延迟高度敏感且网络稳定的场景可选择不开启。此外index.md 中的健康监控Health Status相关参数也在此区域配置Health Check Duration健康状态统计周期、Busy State Threshold忙状态阈值写队列中排队项与队列容量的比值默认 100%、Max Write Queue Length最大写队列长度、Write Error Threshold健康检查周期内允许的写错误数超限上报错误。异常处理策略Exception Handling Strategy异常处理策略区域默认折叠点击展开。通用处理策略详见 resources/_03-exception-handling-strategy.mdx有四种Archive归档将无效数据写入归档文件默认位于${data_dir}/tasks/id/datetime下不写入目标数据库Discard丢弃忽略无效数据Error报错报告错误Cache缓存当目标连接失败或资源不足时将数据写入缓存文件待目标恢复后再写入。针对不同异常条件可配置的策略组合如下异常条件可选策略目标连接超时归档、丢弃、报错、缓存目标数据库不存在归档、丢弃、报错表不存在归档、丢弃、报错、自动建表并重试主时间戳超出范围now - keep1至now 100y归档、丢弃、报错主时间戳为 null归档、丢弃、报错、使用当前时间复合主键为 null归档、丢弃、报错表名超过 192 字符归档、丢弃、报错、截断、截断并归档表名含非法字符如.归档、丢弃、报错、用配置的字符串替换非法字符表名模板变量为 null丢弃、留空、用配置的字符串替换列不存在归档、丢弃、报错、自动补列并重试列名超过 64 字符归档、丢弃、报错列值超出定义长度归档、丢弃、报错、截断、截断并归档或开启自动扩列Automatic Column Expansion修改表结构后重试其他数据错误归档、丢弃、报错额外的全局设置包括Connection Timeout连接超时目标连接超时秒数取值范围1至600Temporary Storage Location临时存储位置相对于${data_dir}/tasks/id/的路径Archive Retention Days归档保留天数非负整数0表示不限制Archive Available Space归档可用空间取值0至655350表示不限制Archive Location归档位置相对于${data_dir}/tasks/id/的路径Archive Write Failure Strategy归档写入失败策略删除旧文件、丢弃数据、或报错并停止任务。合理组合这些策略可以在数据质量异常时既保证主流程不中断又保留可追溯的原始数据。提交任务与任务管理完成上述所有配置后点击Submit提交按钮完成 MQTT 到 TDengine 数据同步任务的创建系统自动返回Data Source List数据源列表页面查看任务执行状态。任务管理能力详见 index.md 的任务管理一节包括对任务执行启动、停止、查看、删除、复制等操作并查看各任务的运行状态已写入记录数、流量等健康状态Health Statusv3.3.5.0 起任务列表为每个运行中任务展示健康状态包括Ready就绪、Idle空闲、Active活跃、Pending等待、Busy繁忙写队列超过阈值需调参或扩容、Bounce写错误超过阈值可能意味着大量无效数据或数据丢失、SourceError源不可读自动重连、SinkError目标不可写恢复后回到 Ready、Fatal严重不可恢复错误等状态健康状态为空表示任务尚无数据进入断点恢复Resume Tasks from Checkpoints大部分数据源可从最后持久化的断点恢复。MQTT 与 OPC UA、OPC DA 一样在开启Cache Realtime Data时会将消息持久化到磁盘网络中断或任务重启后从断点位置继续消费从而避免数据重复或丢失。这也是使用 TDengine Bnode 作为 MQTT Broker 时可指定初始订阅位置这一能力的底层支撑。小结通过 taosExplorer 与 taosXTDengine 为 MQTT 数据接入提供了完整的零代码路径从 Broker 连接与 TLS 认证、订阅主题与 QoS 控制、主题解析到 JSON Payload 解析、字段拆分、数据过滤与动态表映射再到批处理、存储转发、异常处理与健康监控全部通过浏览器界面可视化配置完成。结合 TDengine 内置 Bnode 形式的 MQTT Broker 能力source/libs/tmqtt与断点续传机制MQTT 场景下的实时数据可稳定、低延迟、不丢不重地汇入 TDengine 时序数据平台为工业物联网与物联网应用提供可靠的数据底座。【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考