使用 TDengine taosExplorer 将 Apache Pulsar 数据接入 TDengine:无代码数据接入指南
【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine
本文基于 TDengine 开源仓库中的 Pulsar 接入文档 编写,系统讲解如何通过 taosExplorer(TDengine 图形化管理界面)以零代码方式创建从 Apache Pulsar 到 TDengine 集群的数据接入任务。读者将掌握 Pulsar 数据源的连接配置、四种认证机制、采集参数设置、Payload 解析、字段拆分、数据过滤、表映射、异常处理策略等完整实战流程,并了解底层 taosX 数据接入组件的行为约定。
功能概述
Apache Pulsar 是一个云原生的开源分布式消息与流处理平台,提供多租户、持久化存储、跨地域复制等能力,广泛用于消息队列和流式数据处理场景。TDengine 作为面向工业物联网(IIoT)等场景设计的高性能时序数据库,可以通过数据接入组件taosX高效地从 Pulsar 读取消息并写入 TDengine,从而实现历史数据迁移或实时数据流入库两种典型场景。
整个过程完全通过 taosExplorer 的图形界面完成,无需编写任何代码,属于 TDengine 数据接入体系中的"无代码接入(No-Code Ingestion)"能力。taosExplorer 是 TDengine 集群自带的可视化运维管理界面,其中"数据写入"页面聚合了多种数据源的接入向导,Pulsar 是其中之一。
创建任务
在 taosExplorer 的数据写入页面中,点击+新增数据源按钮,进入新增数据源页面,开始创建 Pulsar 数据接入任务。
新增数据源
进入新增数据源页面后,需要先填写任务的基本信息:
- 名称:输入任务名称,例如
test_pulsar; - 类型:在下拉列表中选择Pulsar;
- 代理(Agent):非必填项,如有需要可以在下拉框中选择指定的代理,也可以先点击右侧的+创建新的代理。代理即 taosX Agent,负责执行具体的数据采集与写入任务;
- 目标数据库:在下拉列表中选择一个目标数据库,也可以先点击右侧的+创建数据库按钮在线创建。
配置连接信息
在连接信息区域填写 Pulsar Broker 的地址:
- Broker Server,例如:
192.168.2.131:6650。
只需要填写一个有效的 broker server 地址即可,Pulsar 客户端会通过该地址发现集群中的其他 Broker 并进行路由。地址格式为host:port,其中 6650 是 Pulsar Broker 的默认服务端口。
认证机制
如果服务端开启了相关认证机制,此处需要填写认证信息。当前支持Basic Auth / JWT / mTLS / Custom Authentication四种认证机制,请按实际情况进行选择;如果服务端没有配置任何认证,可以跳过此步骤,认证字段留空即可。
Basic Auth 认证
选择Basic-Auth认证机制,输入用户名和密码。适用于 Pulsar 服务端启用了基于用户名密码的简单认证场景。
JWT 认证
选择JWT认证机制,输入 JWT token 信息。JWT(JSON Web Token)是 Pulsar 常用的无状态认证方式,Token 由 Pulsar 管理员通过密钥或密钥对签发,填写时需要保证 Token 与 Pulsar 服务端配置的认证密钥一致。
配置 mTLS 证书认证
如果服务端开启了 mTLS(双向 TLS)加密认证,此处需要启用 mTLS 并配置相关内容。mTLS 要求客户端与服务端互相验证证书,通常需要上传或配置客户端证书、私钥以及 CA 证书,用于建立可信的双向加密通道。
Custom Authentication 认证
选择Custom Authentication,输入服务器自定义的认证信息即可。该选项面向启用了自定义认证插件(Custom Authentication Provider)的 Pulsar 服务端,认证信息的格式由服务端定义。
配置采集信息
在采集配置(Collection Configuration)区域填写采集任务相关的配置参数,这是决定数据消费行为的关键部分。
超时时间(Timeout)
在超时时间中填写超时时间。当从 Pulsar 消费不到任何数据,且持续时间超过该超时值后,数据采集任务会退出。默认值是0 ms:
- 当 timeout 设置为
0时,会一直等待,直到有数据可用,或者发生错误; - 当 timeout 大于
0时,若在指定时间内没有可消费的数据,任务将自动退出,避免任务长期空转占用资源。
主题(Topic)
在主题中填写要消费的 Topic 名称。可以配置多个 Topic,Topic 之间用逗号分隔,例如:
persistent://public/default/tp1,persistent://public/default/tp2Pulsar 的 Topic 完整名遵循persistent://<tenant>/<namespace>/<topic>的格式,示例中public为租户,default为命名空间,tp1、tp2为具体的 Topic 名称。多个 Topic 会被同一个消费者统一消费。
消费者名称(Consumer Name)
在消费者名称中填写消费者标识,填写后会生成带有taosx前缀的消费者 ID。例如输入标识foo,生成的消费者 ID 为taosxfoo。如果打开末尾处的开关,则会把当前任务的任务 ID拼接到taosx之后、输入的标识之前,例如任务 ID 为100时生成taosx100foo。
这一设计保证不同任务的消费者 ID 互不冲突,方便在 Pulsar 服务端按消费者维度进行监控和排查。
订阅名称(Subscription Name)
在订阅名称中填写订阅名标识,填写后会生成带有taosx前缀的订阅 ID。生成规则与消费者名称一致:输入标识前自动添加taosx前缀;如果打开末尾处的开关,则把当前任务的任务 ID 拼接到taosx之后、输入的标识之前。
订阅(Subscription)是 Pulsar 消费模型中的核心概念,订阅名决定了消费的进度游标归属,同一订阅名下的多个消费者可以共享消费进度(Exclusive / Shared / Failover 等订阅模式由 Pulsar 服务端配置决定)。
Initial Position
在Initial Position的下拉列表中选择从哪个位置开始消费数据,有两个选项,默认值为Earliest:
- Earliest:用于请求最早的位置,即从 Topic 中可用的最早消息开始消费。适合历史数据迁移、全量回放场景;
- Latest:用于请求最晚的位置,即从当前最新消息开始消费,仅消费新建任务之后产生的新数据。适合仅关注增量实时数据的场景。
该参数决定了新建消费任务首次连接时的消费起点;在已有订阅进度的情况下,实际消费位置以 Pulsar 服务端保存的订阅游标为准。
字符编码
在字符编码中配置消息体编码格式。taosX 在接收到消息后,使用对应的编码格式对消息体进行解码,从而获取原始数据。可选项为:
UTF_8(默认)GBKGB18030BIG5
如果消息体是中文等多字节文本,请根据消息生产端的实际编码选择合适的字符集,否则可能出现乱码或解析失败。
完成以上配置后,点击连通性检查按钮,可立即检查数据源是否可用(包括 Broker 地址可达性、认证信息正确性等)。
配置 Payload 解析
在Payload 解析区域填写 Payload 解析相关的配置参数,将 Pulsar 消息体解析为结构化字段,这是把消息转换为 TDengine 表数据前的关键步骤。
解析(Parsing)
有三种获取示例数据的方法:
- 点击从服务器检索(Retrieve from Server)按钮,从 Pulsar 实时获取示例数据;
- 点击文件上传(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这 2 个字段:
- 选择split提取器;
- separator填写
-(分隔符); - number填写
2(拆分后的字段数量)。
对示例数据"hello-world"按-拆分后,将得到message_0 = hello、message_1 = world两个新字段。
点击新增(Add)可以添加更多提取规则,点击删除(Delete)可以删除当前提取规则。点击放大镜图标可查看预览提取/拆分结果。
数据过滤(Data Filtering)
在过滤(Filter)中填写过滤条件。例如填写id != 1,则只有id不为 1 的数据才会被写入 TDengine,实现消息的按需筛选,减少无效数据入库。
点击新增(Add)可以添加更多过滤规则,多条规则之间为叠加生效关系;点击删除(Delete)可以删除当前过滤规则。点击放大镜图标可查看预览过滤结果,确认过滤条件命中情况。
表映射(Table Mapping)
在目标超级表(Target Supertable)的下拉列表中选择一个目标超级表,也可以先点击右侧的创建超级表(Create Supertable)按钮新建。
在映射(Mapping)中填写目标超级表中的子表名称与字段映射规则:
- 子表名称支持模板变量,例如
t_{id},其中{id}会被实际数据中id字段的值替换; - 根据需求填写映射规则,将消息字段映射到超级表的列;
- 映射支持设置缺省值(默认值),当某字段在消息中缺失时使用默认值填充。
点击预览(Preview)可以查看映射的结果,确认子表名称与列映射生成是否符合预期。
配置高级选项
高级选项(Advanced Options)区域默认折叠,点击右侧的>可展开。针对 Pulsar 这类消息队列数据源,常见的高级选项包括:
- 最大读取并发(Maximum Read Concurrency):限制数据源的连接数或读取线程数。默认
0表示由连接器自动配置;当源端响应较慢且允许更高并发时可适当调大; - 批量大小(Batch Size):单次发送的消息或行数的最大值,常见默认值为
1000; - 写入并发(Write Concurrency):指定可同时写入 TDengine 的任务数。
配置异常处理
异常处理策略(Exception Handling Strategy)区域默认折叠,点击>可展开。常见的处理策略包括:
- 归档(Archive):将无效数据写入归档文件,默认位置为
${data_dir}/tasks/<id>/<datetime>,不写入目标数据库; - 丢弃(Discard):忽略无效数据;
- 报错(Error):报告错误;
- 缓存(Cache):当目标连接失败或资源不足时,将数据写入缓存文件,待目标恢复后继续入库。
可以针对以下异常场景分别配置策略:
- 目标连接超时:归档、丢弃、报错或缓存;
- 目标数据库不存在:归档、丢弃或报错;
- 表不存在:归档、丢弃、报错,或自动建表后重试;
- 主时间戳超出范围(
now - keep1至now + 100y):归档、丢弃或报错; - 主时间戳为空:归档、丢弃、报错,或使用当前时间;
- 复合主键为空:归档、丢弃或报错;
- 表名超过 192 字符:归档、丢弃、报错、截断,或截断并归档;
- 表名包含非法字符(如
.):归档、丢弃、报错,或用配置的字符串替换非法字符; - 表名模板变量为空:丢弃、变量留空,或用配置的字符串替换;
- 列不存在:归档、丢弃、报错,或自动补列后重试;
- 列名超过 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 到 TDengine 的数据同步任务。回到数据源列表页面可查看任务的执行状态,包括任务运行是否正常、消费进度以及写入情况等。
关联资料与深入阅读
- 本文对应文档:docs/en/08-data-ingest-and-delivery/01-no-code-ingestion/20-pulsar.md
- 同系列 Kafka 数据接入指南(结构与 Pulsar 高度一致,可作为对照参考):docs/en/08-data-ingest-and-delivery/01-no-code-ingestion/08-kafka.md
- 异常处理策略完整定义:resources/_03-exception-handling-strategy.mdx
- 消息队列类数据源高级选项定义:resources/_02-advanced-options-mq.mdx
- 数据接入功能总览与健康状态说明:01-no-code-ingestion/index.md
- 安装数据接入代理(taosX Agent):01-install-agent.md
【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考