FEATURED · 精选文章

使用 TDengine taosExplorer 将 Apache Pulsar 数据接入 TDengine:无代码数据接入指南

发布时间 / 2026/9/12 19:46:36
来源 / 创域科博编辑部
栏目 / 资讯中心
使用 TDengine taosExplorer 将 Apache Pulsar 数据接入 TDengine:无代码数据接入指南 使用 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 接入文档 编写系统讲解如何通过 taosExplorerTDengine 图形化管理界面以零代码方式创建从 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 信息。JWTJSON 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 名称。可以配置多个 TopicTopic 之间用逗号分隔例如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的下拉列表中选择从哪个位置开始消费数据有两个选项默认值为EarliestEarliest用于请求最早的位置即从 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到655350表示不限归档位置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 Agent01-install-agent.md【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED — 相关阅读

相关资讯

LATEST — 最新资讯

最新发布

TODAY — 本日精选

新闻

WEEKLY — 本周精选

新闻

MONTHLY — 本月精选

新闻