
TDengine 零代码接入 Apache Pulsar通过 taosExplorer 配置 Pulsar 数据写入任务【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengineApache Pulsar 是云原生的开源分布式消息与流处理平台在 IIoT 场景中常被用作设备数据的缓冲与分发层。本文基于 TDengine 的零代码数据写入体系完整讲解如何在 taosExplorer 界面中创建一个从 Pulsar 消费数据并写入 TDengine 集群的数据同步任务涵盖连接信息、认证机制、采集参数、Payload 解析、字段拆分、数据过滤、表映射、高级选项与异常处理策略并补充了底层数据转换规则与任务管理机制读者按文操作即可完成 Pulsar 历史数据迁移或实时数据流入库。功能概述Pulsar 数据接入的价值Apache Pulsar 是一个云原生的开源分布式消息与流处理平台支持多租户与持久化消息常用于解耦数据生产端与消费端。TDengine 可以高效地从 Pulsar 读取数据并将其写入 TDengine实现两类典型目标历史数据迁移将 Pulsar 中已堆积的消息批量回放并写入 TDengine实时数据流入库持续消费 Pulsar 中的实时消息低延迟地落入 TDengine 时序表。整个接入过程在 TDengine 的零代码数据写入体系taosExplorer taosX中完成。taosX 是负责实际数据采集与写入的组件taosExplorer 提供浏览器端的配置界面用户无需编写任何代码即可完成Pulsar → TDengine的数据管道搭建并且支持在写入过程中对数据进行提取、过滤和转换从而统一命名空间并提升数据质量减少额外的 ETL 组件。完整的体系说明与数据源列表参见零代码数据写入总览。创建任务创建 Pulsar 数据写入任务的前置条件是已部署并登录 taosExplorer且已在运行 taosX 采集器Agent的机器上完成连接配置参见安装 Agent。新增数据源在数据写入页面中点击新增数据源按钮进入新增数据源页面。配置基本信息在新增数据源页面中依次配置以下基本信息名称输入任务名称例如test_pulsar类型在下拉列表中选择Pulsar代理非必填项。如有需要可以在下拉框中选择指定的代理也可以先点击右侧的创建新的代理新建代理目标数据库在下拉列表中选择一个目标数据库也可以先点击右侧的创建数据库按钮创建新库。配置连接信息在Broker Server中填写 Pulsar Broker 的地址格式为host:port例如192.168.2.131:66506650 为 Pulsar 的默认 TCP 服务端口。只需要填写一个有效的 broker server 地址即可taosX 会通过该地址完成 Pulsar 集群的元数据发现与连接。认证机制如果 Pulsar 服务端开启了相关认证机制此处需要填写认证信息。目前支持Basic Auth / JWT / mTLS / Custom Authentication四种认证机制请按服务端实际配置进行选择。如果服务端没有配置任何认证可跳过此步骤不填写。Basic Auth 认证选择Basic-Auth认证机制输入 Pulsar 服务端配置的用户名和密码。JWT 认证选择JWT认证机制输入 Pulsar 服务端签发的 JWT token 信息。JWT 是 Pulsar 常用的认证方式适用于使用pulsar-admin或 client 配置了 token 认证的集群。配置 mTLS 证书认证如果服务端开启了 mTLS 加密认证此处需要启用 mTLS 并配置客户端证书、私钥及 CA 证书等相关内容实现双向 TLS 认证。Custom Authentication 认证选择Custom Authentication输入服务器自定义的认证信息即可适用于使用了自定义认证插件Authentication Provider的 Pulsar 集群。配置采集信息在采集配置区域填写采集任务相关的配置参数这些参数直接决定 taosX 以何种方式消费 Pulsar 消息。超时时间当从 Pulsar 消费不到任何数据超过 timeout 后数据采集任务会退出。默认值是0 ms。当 timeout 设置为0时会一直等待直到有数据可用或者发生错误。主题Topic填写要消费的 Topic 名称。可以配置多个 TopicTopic 之间用逗号分隔例如persistent://public/default/tp1,persistent://public/default/tp2。Pulsar Topic 使用完整的persistent://tenant/namespace/topic形式。消费者名称填写消费者标识填写后会生成带有taosx前缀的消费者 ID。如果打开末尾处的开关则会把当前任务的任务 ID 拼接到taosx之后、输入的标识之前从而在多任务并发消费同一 Topic 时避免消费者 ID 冲突。订阅名称填写订阅名标识填写后会生成带有taosx前缀的订阅 ID。同样地打开末尾处的开关后任务 ID 会拼接到taosx之后、输入的标识之前。Initial Position选择从哪个位置开始消费数据有两个选项默认值为EarliestEarliest请求最早的位置即从 Topic 中最老的消息开始消费适合历史数据迁移场景Latest请求最晚的位置即仅消费新到达的消息适合纯实时入库场景。字符编码配置消息体编码格式taosX 在接收到消息后使用对应的编码格式对消息体进行解码以获取原始数据。可选项为UTF_8、GBK、GB18030、BIG5默认为UTF_8。当 Pulsar 中消息体为中文等非 UTF-8 编码时需按生产端实际编码调整此项否则解码后会出现乱码。配置完成后点击连通性检查按钮检查数据源是否可用。若检查失败请根据页面返回的错误提示调整 Broker 地址、认证信息或 Topic 配置。配置 Payload 解析在Payload 解析区域填写 Payload 解析相关的配置参数。这一步的作用是把 Pulsar 消息体非结构化字符串解析为结构化字段为后续的拆分、过滤与映射提供基础。关于解析、转换规则的通用定义可参考零代码数据写入总览中的数据提取、过滤和转换章节。解析Pulsar 消息体为 JSON 格式数据时有三种获取示例数据的方法点击从服务器检索按钮从配置的 Pulsar 服务器实时获取示例数据点击文件上传按钮上传 CSV 文件将文件内容作为示例数据在消息体中直接手动填写 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}]解析结果会以结构化表格的形式呈现。点击放大镜图标可查看预览解析结果确认字段名与字段类型是否符合预期。在解析阶段除自动解析简单的 key/value 结构外还可使用JSON Path表达式从嵌套 JSON 中提取感兴趣的字段例如$[data][voltage]voltage对于非 JSON 的文本消息体可使用正则表达式的命名捕获组提取多个字段对于一次上报需要拆分为多行的数据还可使用 UDT 自定义 rhai 脚本解析脚本仅支持 json 格式原始数据输出必须是数组。注意JSON 属性名称中不能含有.如果含有则必须使用别名alias将名称转义。字段拆分在从列中提取或拆分中填写从消息体中提取或拆分的字段。例如将message字段拆分成message_0和message_1这 2 个字段选择split提取器separator分隔符填写-number拆分数量填写2。点击新增可以添加更多提取规则点击删除可以删除当前提取规则。点击放大镜图标可查看预览提取/拆分结果。拆分规则除split外还支持regex正则表达式同样使用命名捕获组命名提取字段与convert转换填写 JSON map 对象做值映射拆分后的字段命名规则为{原字段名}_{顺序号}。数据过滤在过滤中填写过滤条件满足条件的数据行才会被写入 TDengine。例如填写id ! 1则只有id不为 1 的数据才会被写入 TDengine。点击新增可以添加更多过滤规则点击删除可以删除当前过滤规则。点击放大镜图标可查看预览过滤结果。过滤条件表达式的结果必须是 boolean 类型支持的语法要点如下数值类型int/float支持比较操作符、!、、、、字符串类型除比较操作符外还支持is_empty()、contains()、starts_with()、ends_with()、len等字符串函数BOOL 类型可直接使用变量或!取反例如inuse、!inuse多个判断表达式可使用逻辑操作符、||、!组合例如location.starts_with(beijing) voltage 200时间戳字段可使用between_time_range(ts, t1, t2)做时间范围过滤例如只允许最近 7 天内的数据入库between_time_range(ts, -604800, 0)通过 regex、split 规则解析出的字段均为 string 类型如需按数值过滤可先使用parse_int(56)、parse_float(12.3)做类型转换。表映射在目标超级表的下拉列表中选择一个目标超级表也可以先点击右侧的创建超级表按钮新建。选择后页面会加载出超级表的所有 tags 和 columns源字段根据名称自动使用 mapping 规则映射到目标超级表的 tag 和 column。在映射中填写目标超级表中的子表名称例如t_{id}。其中${}作为占位符占位符中可以是一个源字段也可以是 string 类型字段的函数处理mapping 支持设置缺省值。点击预览可以查看映射的结果确认子表名与字段映射是否符合预期。除直接映射mapping外支持的映射规则还包括value常量输入的常量值直接入库generator生成器目前仅支持时间戳生成器now入库时会将当前时间入库join字符串连接器可指定连接字符拼接多个源字段format字符串格式化工具如${year}-${month}-${day}占位符内还支持pad、trim、sub_string、replace等字符串处理函数sum多个数值型字段做加法计算expr数值运算表达式支持、-、*、/四则运算及sin、cos、sqrt、exp、ln、floor、ceiling、round等数学函数例如将摄氏温度转为华氏温度temperature * 1.8 32。配置高级选项高级选项区域默认折叠点击右侧可展开。Pulsar 属于消息队列类数据源常见配置项如下最大读取并发数限制数据源连接数或读取线程数默认0表示由采集器自动配置批次大小单次发送的最大消息数或行数默认常见为1000写入并发数量同时写入 TDengine 的并发任务数量。此外从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 - keep1到now 100y可选归档、丢弃、报错主键时间戳空可选归档、丢弃、报错、使用当前时间将当前时间填充到空的时间戳字段复合主键空可选归档、丢弃、报错表名长度溢出子表表名最大 192 字符可选归档、丢弃、报错、截断、截断且归档表名非法字符检查子表表名中是否包含.等特殊字符可选归档、丢弃、报错、非法字符替换为指定字符串例如a.b替换为a_b表名模板变量空值可选丢弃、留空如a_{x}转换为a_、变量替换为指定字符串如a_{x}转换为a_b列名不存在可选归档、丢弃、报错、自动增加缺失列自动修改表结构增加列修改成功后重试列名长度溢出列名最大 64 字符可选归档、丢弃、报错列自动扩容开关选项打开时列数据长度超长将自动修改表结构并重试列长度溢出可选归档、丢弃、报错、截断、截断且归档数据异常其他未在上方列出的异常可选归档、丢弃、报错连接超时目标库连接超时时间单位秒取值范围 1~600临时存储文件位置缓存文件的位置实际生效位置为${data_dir}/tasks/:id/{location}归档数据保留天数非负整数0表示无限制归档数据可用空间0~65535其中0表示无限制归档数据文件位置归档文件的位置实际生效位置为${data_dir}/tasks/:id/{location}归档数据失败处理策略当写入归档文件报错时的处理策略可选删除旧文件删除后仍无法写入则报错并停止任务、丢弃丢弃即将归档的数据、报错并停止任务。完整说明见异常处理策略资源。建议在生产任务中为表不存在开启自动建表、为主键时间戳空开启使用当前时间并合理设置归档与缓存策略以提升数据落库率。创建完成与任务管理点击提交按钮完成创建 Pulsar 到 TDengine 的数据同步任务回到数据源列表页面可查看任务执行情况。提交成功后任务状态会切换至运行中若提交失败可通过查看该任务的活动日志查找错误原因。在任务列表页面还可以对任务进行启动、停止、查看、删除、复制等操作并查看各个任务的运行情况包括写入的记录条数、流量等指标。任务运行中页面会展示健康状态Ready、Idle、Active、Pending、Busy、Bounce、SourceError、SinkError、Fatal 等其中 Busy 表示写入队列已满可能存在性能瓶颈需要调整并发与批次参数SourceError/SinkError 表示数据源或写入端出错taosX 会自动尝试重连。关于断点恢复taosX 大部分数据源支持从上次写入的断点恢复。Pulsar 与 Kafka 同属消息队列类数据源其消费进度由 Pulsar 服务端的订阅游标subscription cursor机制管理可以推断任务重启后可基于订阅的持久化游标位置继续消费避免重复或遗漏数据。对于实时消费链路建议同时配合缓存实时数据等高级选项将消费数据优先写入持久化磁盘再落库增强任务重启后的数据完整性。小结本文完整梳理了在 taosExplorer 中创建 Pulsar → TDengine 数据写入任务的全流程从新增数据源、配置 Broker 连接与四种认证机制到采集参数超时时间、多 Topic、消费者/订阅命名、Initial Position、字符编码的设置再到 Payload 的解析、字段拆分、数据过滤与表映射最后通过高级选项与异常处理策略保障任务的吞吐与数据质量。掌握这些配置后即可将 Pulsar 中的设备消息以零代码方式持续汇入 TDengine为工业物联网场景提供统一、可靠的时序数据底座。【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考