TDengine 通过 taosExplorer 零代码接入 Pulsar-Tuya:从连接配置到数据迁移全流程指南
【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine
Pulsar-Tuya 是涂鸦智能基于开源 Apache Pulsar 定制的高可用消息集群,承载着大量物联网设备上报数据。本文基于 TDengine 企业版的 taosExplorer 图形界面,系统讲解如何以"零代码"方式创建从 Pulsar-Tuya 到 TDengine 的数据同步任务,实现历史数据迁移与实时数据接入,涵盖数据源添加、平台认证、采集参数、Payload 解析、字段拆分、数据过滤、表映射、高级选项与异常处理等完整环节。
说明:本指南对应的功能属于TDengine TSDB-Enterprise 企业版专属能力,TDengine TSDB-OSS 社区版不包含 taosX 数据接入相关组件。完整操作流程参见 Pulsar-Tuya 接入文档,通用 Apache Pulsar 的接入方式参见 Pulsar 接入文档。
功能概述:为什么选择 Pulsar-Tuya 作为数据源
Apache Pulsar 是一款云原生、开源的分布式消息与流处理平台,采用存储与计算分离架构,天然支持多租户、跨地域复制与海量 topic。Pulsar-Tuya 则是涂鸦智能基于开源 Apache Pulsar 定制的集群版本,为物联网场景提供了更强的托管能力与平台化认证体系。
TDengine 的 taosX 数据接入引擎可以从 Pulsar-Tuya 高效读取消息数据并写入当前 TDengine 集群,从而支撑两类典型场景:
- 历史数据迁移:将 Pulsar-Tuya 中已积压的历史消息批量回放并落库;
- 实时数据接入:持续订阅新到达的消息,形成流式数据管道。
整个接入过程完全在 taosExplorer 的 Web 图形界面中完成,无需编写任何连接或解析代码。
前置条件
在开始创建任务前,需要确保以下条件已满足:
- 已部署 TDengine 企业版集群,并可通过浏览器访问 taosExplorer(默认位于运行 TDengine 的主机或 IP 的6060 端口);
- 已获得涂鸦平台的Access Id与Access Key认证信息;
- 已准备好 Pulsar-Tuya 的Broker Server 地址(如
mqe.tuyaus.com:6650); - 若 taosX 无法直接访问你的 Pulsar-Tuya 网络,需要先安装 taosX-Agent,使 taosX 通过 Agent 间接连接数据源。
操作流程概览
创建 Pulsar-Tuya 数据同步任务的整体流程如下:
- 添加数据源(Data In → Add Task);
- 配置 Broker 连接信息;
- 填写涂鸦平台认证信息;
- 配置采集信息(超时、初始位置、字符编码);
- 配置 Payload 解析(解析、字段拆分、数据过滤、表映射);
- 配置高级选项;
- 配置异常处理策略;
- 提交任务并查看运行状态。
添加数据源
在浏览器中打开 taosExplorer 后,按下述步骤创建一个 Pulsar-Tuya 数据接入任务:
- 在左侧主菜单中选择Data In(数据接入),点击Add Task(添加任务);
- 在Name字段中为该数据接入任务输入唯一名称;
- 从Type下拉列表中选择Pulsar-Tuya;
- (可选)如果该任务需要借助 Agent 访问数据源,从Agent下拉列表中选择已创建的 Agent,也可以点击Create New Agent现场创建;
- 从Target DB下拉列表中选择该任务写入的目标数据库,也可以点击Create Database现场创建。
配置连接信息
在Broker Server中填写 Pulsar-Tuya 的 Broker 服务地址,例如:
mqe.tuyaus.com:6650Pulsar 客户端会从该地址获取整个集群的元数据与连接信息,因此只需填写一个有效的 Broker 地址即可,无需逐个罗列集群内所有 Broker。
认证机制:涂鸦平台专用认证
与通用 Pulsar 接入(支持 Basic Auth / JWT / mTLS / Custom Authentication 四种认证)不同,Pulsar-Tuya 使用涂鸦平台特有的认证体系。在认证区域需要填写:
- Access Id:涂鸦平台提供的访问标识;
- Access Key:涂鸦平台提供的访问密钥;
- 执行环境:根据你的涂鸦账号所属环境选择对应的执行环境(如国内或海外环境)。
需要特别注意:Pulsar 所需的 Topic、Consumer Name(消费者名称)、Subscription Name(订阅名称)会由系统根据你填写的 Access Id 与 Access Key 自动生成,因此本环节不需要手动指定这些 Pulsar 概念,你只需关注下述采集参数。
配置采集信息
Collection Configuration(采集配置)区域集中了与采集任务相关的参数,共包含以下几项。
Timeout(超时时间)
填写超时时间。当连续从 Pulsar 消费不到数据且持续时间超过该阈值时,数据采集任务将自动退出。默认值为 0 ms;当设置为 0 时,任务会无限期等待,直到有数据到达或发生错误为止。该参数适用于消息稀疏的场景,可避免长时间空转。
Initial Position(初始消费位置)
通过下拉列表选择开始消费数据的位置,共两个选项,默认值为 Earliest:
| 选项 | 含义 |
|---|---|
Earliest | 从最早的位置开始消费,适用于历史数据迁移或全量回放 |
Latest | 从最新的位置开始消费,适用于只关心新数据的实时接入 |
Character Encoding(字符编码)
配置消息体的编码格式。当 taosX 收到消息后,会使用该编码对消息体进行解码以还原原始数据。可选值与默认值如下:
| 编码 | 说明 |
|---|---|
UTF_8 | 默认值,通用 UTF-8 编码 |
GBK | 中文环境常用编码 |
GB18030 | 国标扩展编码 |
BIG5 | 繁体中文编码 |
完成上述配置后,点击Connectivity Check(连通性检查)按钮,即可验证数据源是否可用。
配置 Payload 解析
Payload Parsing(负载解析)区域用于定义如何将 Pulsar 消息体转换为可写入 TDengine 的结构化数据,依次包含解析、字段拆分、数据过滤、表映射四个环节。
Parsing(解析):获取与解析样本数据
获取样本数据有三种方式:
- 点击Retrieve from Server(从服务器获取)按钮,直接从 Pulsar-Tuya 拉取样本数据;
- 点击File Upload(文件上传)按钮,上传 CSV 文件获得样本数据;
- 在Message Body(消息体)中手动输入 Pulsar 消息体的样本数据。
JSON 数据支持JSONObject或JSONArray两种形态,均可使用 JSON 解析器解析,例如:
{"id": 1, "message": "hello-world"} {"id": 2, "message": "hello-world"}或:
[{"id": 1, "message": "hello-world"},{"id": 2, "message": "hello-world"}]解析完成后,点击放大镜图标即可预览解析结果,确认字段与类型是否符合预期。
Field Splitting(字段拆分)
在Extract or Split from Columns(从列中提取或拆分)中填写需要从消息体中提取或拆分的字段。例如,要将message字段按分隔符拆分为message_0和message_1两个字段,操作如下:
- 选择split(拆分)提取器;
- 在分隔符(Separator)中填入
-; - 在数量(Number)中填入
2。
点击Add可添加更多提取规则,点击Delete可删除当前提取规则。完成后同样可通过放大镜图标预览提取/拆分结果。
Data Filtering(数据过滤)
在Filter(过滤条件)中填写过滤条件,例如输入:
id != 1则只有id不等于 1 的数据才会被写入 TDengine。点击Add可叠加多条过滤规则,点击Delete可删除当前规则,并可通过放大镜图标预览过滤后的结果。该功能适合在源头剔除脏数据或指定设备数据,减少无效写入。
Table Mapping(表映射)
在Target Supertable(目标超级表)下拉框中选择目标超级表,也可以点击右侧的Create Supertable(创建超级表)按钮现场创建。
在Mapping(映射)区域:
- 填写目标超级表下的子表名称,例如
t_{id}({id}为取自消息字段的动态模板变量); - 按需填写字段映射规则,映射支持设置默认值,即当消息中缺少某字段时使用默认值兜底。
点击Preview(预览)可以查看映射结果,确认子表名与字段对应关系正确后再进入下一步。
配置高级选项
Advanced Options(高级选项)区域默认折叠,点击右侧的>符号即可展开。其中通常包含并发度、批量大小、任务日志等进阶调优参数,可根据实际吞吐量需求进行设置。展开后修改任意参数并提交,即可让任务按新的配置运行。
配置异常处理策略
Exception Handling Strategy(异常处理策略)区域同样默认折叠,点击>展开后可见完整的策略配置。完整策略说明详见 异常处理策略资源文件。
通用处理策略
针对无效数据或异常,可选择以下四种通用策略:
| 策略 | 行为 |
|---|---|
| Archive | 将无效数据写入归档文件(默认位于${data_dir}/tasks/<id>/<datetime>),不写入目标数据库 |
| Discard | 直接忽略无效数据 |
| Error | 上报错误 |
| Cache | 当目标连接失败或资源不足时,将数据写入缓存文件,待目标恢复后再行写入 |
分场景策略映射
可针对以下具体异常场景分别配置策略:
- 目标连接超时:可配置为归档、丢弃、报错或缓存;
- 目标数据库不存在:可配置为归档、丢弃或报错;
- 表不存在:可配置为归档、丢弃、报错,或自动建表后重试;
- 主时间戳超出范围(
now - keep1至now + 100y):可配置为归档、丢弃或报错; - 主时间戳为 null:可配置为归档、丢弃、报错或使用当前时间;
- 复合主键为 null:可配置为归档、丢弃或报错;
- 表名超过 192 个字符:可配置为归档、丢弃、报错、截断或截断后归档;
- 表名包含非法字符(如
.):可配置为归档、丢弃、报错,或替换为配置的字符串; - 表名模板变量为 null:可配置为丢弃、留空或替换为配置的字符串;
- 列不存在:可配置为归档、丢弃、报错,或自动补列后重试;
- 列名超过 64 个字符:可配置为归档、丢弃或报错;
- 列值超过定义长度:可配置为归档、丢弃、报错、截断或截断后归档;自动扩列(Automatic Column Expansion)可先改表再重试;
- 其他数据错误:可配置为归档、丢弃或报错。
附加设置
- Connection Timeout(连接超时):目标连接超时时间(秒),取值范围
1~600; - Temporary Storage Location(临时存储位置):相对于
${data_dir}/tasks/<id>/的路径; - Archive Retention Days(归档保留天数):非负整数,
0表示永久保留; - Archive Available Space(归档可用空间):取值范围
0~65535,0表示不限; - Archive Location(归档位置):相对于
${data_dir}/tasks/<id>/的路径; - Archive Write Failure Strategy(归档写入失败策略):可配置为删除旧文件、丢弃数据或报错并停止任务。
合理配置异常处理策略可以显著提升数据管道在目标库结构变更、网络抖动等异常场景下的自愈能力,避免单条脏数据导致整个任务中断。
提交任务并查看状态
完成以上所有配置后,点击Submit(提交)按钮,即完成 Pulsar-Tuya 到 TDengine 数据同步任务的创建。随后返回Data Source List(数据源列表)页面,即可查看任务的执行状态。
相关文档
- 通用 Apache Pulsar 接入(含 Basic Auth / JWT / mTLS / 自定义认证):Pulsar 接入指南
- 数据接入总览与更多数据源: 01-no-code-ingestion 目录
- 无法直连数据源时的 Agent 安装方法:安装 taosX-Agent
【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考